尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

AI Infra 下的消息中枢:Apache Pulsar 如何支撑多模态与 Agent 编排

AI Infra 下的消息中枢:Apache Pulsar 如何支撑多模态与 Agent 编排 AI Infra 演进到后半场消息中间件这件事被很多人低估了。模型训练、多模态数据流转、Agent 编排这些场景听起来跟“发消息”八竿子打不着但真正把系统拆开看每一层都在依赖一个可靠的消息中枢。最近我在梳理分布式架构选型时重新审视了 Apache Pulsar发现它在 AI Infra 这个语境下的价值比过去单纯做业务消息队列时要突出得多。这篇文章就聊聊 Pulsar 在分布式演进、多模态处理和 Agent 时代的定位以及我实际搭建和踩坑过程中的一些体会。先说清楚 Pulsar 是什么。Apache Pulsar 是一个分布式消息流平台核心由 Broker 和 BookKeeper 两层组成Broker 负责接入和路由BookKeeper 负责持久化存储。它最大的特点是把计算和存储分离这跟 Kafka 那种 Broker 本地磁盘存储的架构思路完全不同。对 AI Infra 来说这个差异直接决定了系统能不能撑住多模态数据的高吞吐、Agent 事件流的突发性以及训练管道和在线推理之间的数据缓冲。这篇文章适合正在做 AI 平台底座、多模态数据处理管道、Agent 编排系统的工程师和技术决策者。如果你只是用过 Kafka想了解 Pulsar 在 AI 场景下到底强在哪里也能从中获得不少可以直接落地的参考。我会从架构原理、场景拆解、实际配置和问题排查几个维度展开尽量把“为什么需要它”和“它到底怎么用”讲透。1. 分布式演进下的 AI Infra为什么消息中枢成了底座1.1 AI Infra 的三个阶段与消息层的角色变化早期的 AI 基础设施其实很简单训练脚本跑在单机或者几台 GPU 服务器上数据通过文件系统直接读取模型产出结果后写到数据库里整个链路几乎没有“消息”的概念。后来进入分布式训练和推理阶段数据需要多机并行读特征工程需要实时流转模型服务需要动态伸缩这时候中间层必须有一个能承上启下的组件。到了现在这个阶段AI Infra 已经明显呈现出“数据管道 模型服务 业务编排”三层分离的形态数据管道层负责采集、清洗、增强、特征提取多模态数据在这里汇聚。模型服务层负责推理、微调、评估甚至多模型协同。业务编排层负责 Agent 的任务调度、上下文管理、工具调用和结果回传。这三层之间的通信不能靠硬编码的 HTTP 调用因为层与层之间是异步的、突发的、高吞吐的。比如一个视频理解任务从数据管道拿到视频帧序列经过抽帧、OCR、音频转写、图像描述多个子任务每个子任务的产出都可能触发后续步骤这种场景天然就是一个事件流。消息中枢在这时候扮演的角色类似人体的血液循环系统——不直接产生业务价值但所有器官都靠它输送养分和带走代谢物。1.2 分布式系统对消息中间件的“新需求清单”很多团队在选型时还是拿五六年前的标准来衡量消息队列这是不够的。AI Infra 下的消息中枢需求清单已经发生了变化海量 Topic 支持。AI 场景下每个数据源、每个模型任务甚至每个用户会话都可能需要独立的主题Topic 数量轻松上千上万传统的主题较少时优势明显的方案容易吃力。流量削峰和积压隔离。训练任务是周期性的推理流量是突发的消息系统必须能扛住长时间积压并且某个 Topic 的积压不能拖垮其他 Topic 的消费。存储与计算独立扩缩容。数据量增长但算力有限或者算力充足但存储紧张这两种情况都要能独立调整。跨地域复制和多租户隔离。AI 平台通常服务多个团队租户之间的数据隔离和权限控制必须原生支持。稳定的顺序语义和消息回溯。训练样本的顺序在某些场景下会影响模型收敛消息系统需要支持按序消费和从任意位置重放。Pulsar 在架构上就针对这些问题做了设计。Broker 无状态化意味着接入层可以随意伸缩BookKeeper 作为存储层独立扩展Segment 分片机制支持海量数据持久化。对于 AI Infra 这种读多写多、流量波动大、数据生命周期长的场景它的适配度比传统队列高很多。2. Pulsar 架构深度拆解它能成为 AI 消息中枢的根本原因2.1 计算与存储分离到底意味着什么要理解 Pulsar 的价值就要理解它跟 Kafka 的本质区别。Kafka 的 Broker 既负责计算接收、路由、消费协调也负责存储消息写在本地磁盘这带来一个经典问题分区数增加时每个 Broker 的存储和连接开销同步增长Broker 扩容往往需要同时迁移数据。对 AI Infra 里那种“保留海量数据做重放”的场景存储成本和组织复杂度都会上升。Pulsar 把 Broker 和 BookKeeper 拆开Broker 层是无状态的只做消息的收发、路由、订阅管理和元数据协调用 ZooKeeper 或 etcd 管理元数据。BookKeeper 层是有状态的负责持久化消息数据以 Segment 为存储单元分布在多个 Bookie 节点上。这个分离带来的直接好处是你可以在计算密集时多部署 Broker在存储吃紧时单独扩 Bookie两者互不干扰。对 AI 平台来说这意味着训练数据增长不会直接消耗在线的 Broker CPU在线推理的突发流量也不会影响底层存储的稳定性。另外Pulsar 的消息存储在 Segment 层面做多副本默认 2~3 副本故障自动切换数据的可靠性比单机磁盘高得多。训练数据的完整性对于 AI 项目来说特别重要少一个帧、丢一条标注都可能影响整个数据集的有效性。2.2 多租户、Topic 层级和订阅模型的价值Pulsar 的逻辑模型采用三层结构Tenant租户→ Namespace命名空间→ Topic。租户隔离是 AI 平台刚需不同算法团队、不同项目组天然适合映射到不同的租户权限、存储配额、消息保留策略都可以在租户级别统一设置。在 Topic 之下Pulsar 订阅模型有四种模式Exclusive一个 Topic 只能被一个消费者消费适合严格的单消费者场景。Shared多个消费者共享消费一个 Topic 的消息适合任务并行分发。Failover多个消费者中只有一个活跃消费其他作为备用适合高可用场景。Key_Shared按照消息 key 将消息路由到固定的消费者既保证多个消费者的并行度又能保证同一个 key 的消息顺序。这套订阅模型对于 Agent 场景特别有价值。一个 Agent 的任务流可以被看作一个 Topic多个 Worker 可以用 Shared 模式并行处理不同会话但同一个会话的上下文同一个 session id通过 Key_Shared 模式路由到同一个 Worker避免上下文分裂。这是 Kafka 需要靠额外设计才能实现的语义。2.3 消息保留与回溯训练数据的“时间旅行”几乎所有的消息队列都支持消费者消费后删除消息但 AI 场景里这个过程是反过来的——我们希望保留消息以便后续重新处理。比如一个多模态数据管道运行时某个模型的版本更新了需要重新生成所有视频片段的描述信息或者特征提取代码改动后要基于原始数据重新计算特征。如果原始消息已经被消费并删除就不得不重新从源头拉数据成本非常高。Pulsar 支持基于时间和基于存储大小的消息保留策略。你可以设置一个 Topic 保留过去 7 天或者 100GB 的数据消费者即使已经消费完也可以从任意位置重新订阅和回溯。这个能力对“数据回放”和“模型迭代”来说是隐藏的加速器。3. 多模态时代Pulsar 如何处理多种异构数据流3.1 多模态数据管道的真实形态多模态不是“把文本、图片、视频放在一起”这么简单真正的多模态数据管道是高度并发且异构的。我见过一个典型的视频理解项目数据处理流程是这样的视频文件上传后触发拆分任务切成 5~10 秒的片段。每个片段同时走三条分支抽帧生成图像序列、抽取音频轨道、提取字幕文本。图像序列进入图像描述模型音频进入 ASR 模型字幕直接进文本清洗模块。三个分支的产出最终汇聚到一个对齐模块按时间戳对齐生成多模态样本。这个流程中分支任务是典型的扇出Fan-out模型汇聚阶段是扇入Fan-in模型。如果用传统 HTTP 同步调用每个阶段的失败都要做复杂的重试和补偿如果直接用数据库轮询则IO开销和时延都不可控。用消息中间件每个分支独立生产消息、独立消费谁慢了谁快了消息队列自动做缓冲和削峰。3.2 特征流与样本流的消息建模在多模态处理中我倾向于把消息分成三类流样本流Sample Stream用于训练和评估的完整样本包含文本、图像路径、音频路径、标注信息等。这类消息体积较大但频率相对低适合用批量生产的方式写入。特征流Feature Stream模型中间层的输出向量用于向量检索、多模态对齐或者实时推荐。这类消息体积较小但吞吐极高是 Pulsar 最擅长的场景。事件流Event Stream数据管道中的状态变化比如任务开始、任务失败、模型版本切换、数据质量告警。这类消息用于监控和编排延迟要求最高。三类流建议放在不同的 Namespace 下配置不同的存储策略和消费模式。样本流用持久化策略保存较长时间特征流用基于大小的淘汰策略保留最近 N GB事件流设置较短的 TTL 避免积压。3.3 多模态统一处理中的顺序性与一致性多模态融合里有一个很隐蔽的坑数据对齐的顺序一致性。比如视频的音频流和图像流分别经过不同的预处理服务两个服务的处理速度不同可能导致最终汇聚时顺序错乱。Pulsar 的 Key_Shared 订阅模式可以解决这个问题——将视频 ID 作为消息 key确保同一个视频的所有消息都路由到同一个消费者在消费者内部做对齐缓冲避免跨节点的一致性协调。我在实际项目中用了一个相对简洁的方案每个视频片段生成一个 UUID作为消息 key 发往同一个 Topic消费端用 ConcurrentHashMap 做窗口对齐等待某个 UUID 的全部子任务产出到达后再组装成完整样本发送到下游。这个方案在 100 并发流量下表现很稳定Pulsar 的有序路由是保证对齐逻辑简单的前提。4. Agent 时代Pulsar 如何支撑 Agent 编排与事件驱动4.1 Agent 的运行时架构需要什么Agent 的本质是一个事件驱动的循环系统接收任务、拆解子任务、调用工具、整合结果、再决策。这个过程不是线性的而是多轮、多分支、可中断、可重试的。这就对底层基础设施提出了几个要求任务队列必须支持多消费者并行也要支持同一会话的顺序。子任务的状态变化要能触发后续动作需要事件总线的能力。每个 Agent 会话的上下文要能被保存和恢复不能因为某个 Worker 挂了就丢失会话。工具调用的结果回传要做到可靠投递不能漏消息。Pulsar 的 Topic 模型天然适合做 Agent 的任务总线。一个 Agent 的完整生命周期可以映射为一组 Topic任务输入、任务事件、工具调用请求、工具调用结果、最终回复。每个环节之间用消息解耦既方便扩展 Worker又能支持复杂的编排逻辑。4.2 任务编排中的消息流设计我设计过一个比较通用的 Agent 消息模型分享出来供参考task-inputTopic接收用户的初始请求消费者是 Agent 编排器。task-eventTopic记录 Agent 的状态变化开始、调用工具、暂停、恢复、完成、失败所有监控和日志系统订阅这个 Topic。tool-requestTopic编排器把工具调用请求发到这里各个工具执行器以 Shared 模式消费。tool-responseTopic工具执行器把执行结果回传编排器用 Key_Shared 模式按会话 ID 消费保证同一个 Agent 的多个工具结果顺序处理。这个设计的核心是让 Agent 编排器本身保持无状态。编排器的多个实例都消费 task-input但同一个会话的所有后续消息通过 session_id 作为 key始终路由到同一个实例这样即使编排器扩缩容会话状态也不会错乱。4.3 会话上下文持久化与长时记忆Agent 场景还有一个 Kafka 处理起来比较费劲的问题长会话的状态冗余。用户在 Agent 里的每一轮交互都需要带上历史上下文如果每次都从数据库加载完整历史延迟和 IO 都难以接受。Pulsar 的消息回溯能力可以当做一个轻量级的“会话日志”使用。每次交互的消息都发给 session-log TopicAgent 启动时从指定 offset 回溯读取最近 N 条消息快速恢复上下文。基于时间的保留策略可以自动清理陈旧会话不需要额外的过期任务。这个方案不一定适合所有场景但对中低并发的内部业务系统来说实现成本低、效果直观。5. 实操基于 Pulsar 搭建 AI 多模态 Agent 消息中枢5.1 部署选型与资源规划Pulsar 的部署方式有几种裸机集群、容器化部署、云服务。对于中小团队我建议先容器化部署用 Helm Chart 在 Kubernetes 上跑。如果是为了试验和学习也可以用 Docker Compose 起一个单机实例快速验证功能。以下是单机试验模式的 docker-compose 关键片段供参考services: pulsar: image: apachepulsar/pulsar:3.2.0 container_name: pulsar command: bin/pulsar standalone ports: - 8080:8080 - 6650:6650 volumes: - ./pulsar-data:/pulsar/data生产环境需要注意资源分配Broker 是 CPU 和内存密集型建议至少 4C8G 起步Bookie 是磁盘密集型务必使用 SSD 并开启多个 Journal 目录。元数据服务 ZooKeeper 至少 3 节点。消息副本数建议默认 2但训练数据等关键 Topic 可以设为 3。5.2 核心配置项解析Pulsar 的配置项很多但做 AI 场景的消息中枢重点盯这几个Broker 配置broker.confmanagedLedgerDefaultEnsembleSizeSegment 副本数默认 2。managedLedgerDefaultWriteQuorum写入需要的确认副本数默认 2。managedLedgerDefaultAckQuorum写入完成的最小确认数默认 2。maxUnackedMessagesPerConsumer单个消费者未确认消息数上限默认 50000要根据消费逻辑调整。BookKeeper 配置bookkeeper.confjournalMaxWriteRequests写入队列长度影响写入吞吐。dbStorage_writeCacheMaxSizeMb写入缓存大小建议根据内存情况调大。dbStorage_readAheadCacheMaxSizeMb预读缓存大小对回溯消费影响明显。我在实际项目中调过最有效的一个参数是managedLedgerDefaultMarkDeleteRateLimit默认情况下限速很保守导致积压消息的删除速度跟不上生产速度长期运行磁盘占用持续走高。调高这个限制后情况改善明显。消息保留策略配置用 pulsar-admin 设置 Topic 或 Namespace 的保留策略bin/pulsar-admin namespaces set-retention \ --size 100G \ --time 7d \ my-tenant/ai-training这条命令的意思是在my-tenant/ai-training命名空间下的所有 Topic保留最近 100GB 或 7 天的消息哪个先达到就按哪个策略清理。5.3 Topic 规划与命名规范Topic 命名规范这块我总结了一套在 AI Infra 下比较实用的约定{tenant}/{namespace}/{data-type}/{domain}/{entity}/{action}举个例子ai-platform/raw-data/video/uploaded/v1ai-platform/feature-extract/image/embedding/v1ai-platform/agent/task/event/v1这么设计的好处是租户隔离清晰、数据类型清晰、业务域清晰、版本号清晰。后续做权限控制、数据保留策略、监控告警时都可以按前缀匹配不需要逐条配置。版本号放在末尾是为了兼容管道升级时的双跑场景——新旧版本同时运行互相不干扰。5.4 生产端与消费端的工程实现要点生产端我用 Java 客户端写了一个比较典型的异步发送逻辑ProducerString producer client.newProducer(Schema.STRING) .topic(ai-platform/raw-data/video/uploaded/v1) .enableBatching(true) .batchingMaxMessages(1000) .batchingMaxPublishDelay(5, TimeUnit.MILLISECONDS) .compressionType(CompressionType.LZ4) .create(); CompletableFutureMessageId future producer.sendAsync(message); future.whenComplete((msgId, ex) - { if (ex ! null) { // 发送失败需要重试或进入死信队列 } });几个要点批量发送开启后吞吐提升明显尤其是特征流这种千字节级别的小消息压缩选 LZ4 或 ZSTD对文本和 JSON 数据收益很大异步发送一定要处理回调发送失败的消息不能静默丢弃建议重试几次后写入死信 Topic。消费端更要注意的是消费确认模型。AI 管道里消息处理通常是“先持久化再确认”也就是先把处理结果写到数据库或对象存储确认成功后手动 ack避免消息处理成功但确认失败的重复消费问题。consumer.receiveAsync().thenAccept(message - { // 1. 处理消息 // 2. 写结果到存储 // 3. 手动确认 consumer.acknowledge(message.getMessageId()); });6. 常见问题与排查技巧实录6.1 积压消息导致磁盘飙高怎么办一个典型场景特征提取服务挂了一个小时消息在 Topic 里不断积压。虽然 Pulsar 支持积压但磁盘总有上限。排查思路是先看积压大小bin/pulsar-admin topics stats-internal {topic}关注backlogSize字段。判断积压原因到底是消费者挂了还是消费速度跟不上。消费速度跟不上就需要增加消费者并行度。如果积压时间较长且旧消息已经不需要处理可以直接用pulsar-admin topics skip跳过积压或调低保留策略快速清理旧数据。记住一个原则积压隔离是 Pulsar 的优势但积压本身还是需要监控和告警不要让“能积压”变成“不处理”的借口。6.2 消费顺序错乱如何排查用 Shared 模式消费时消息处理顺序天然不保证。如果你发现某个多模态对齐任务的结果错乱排查步骤确认生产端是否用同一个分区键Message Key。确认消费端是否使用 Key_Shared 订阅。检查 Key_Shared 模式下消费者的数量变化——如果消费者动态增减Pulsar 可能需要短暂的重新哈希期间部分消息的顺序会有中断。这个问题的本质是顺序保证是有成本的。Key_Shared 在 Pulsar 中并非全链路强一致如果对顺序有极高要求建议在消费者内部再做一层按 key 的分组缓冲。6.3 多线程消费的正确姿势一个消费者实例内部可以用多线程处理消息但 ack 的顺序要小心。Pulsar 的 ack 支持累计确认也就是 ack 一个消息 id 会把它之前的消息一起确认。如果你用了多线程乱序处理不能随便用累计 ack否则会误确认未处理的消息。我的做法是每个线程处理完消息后用consumer.acknowledgeAsync单独确认对应的 message id不依赖累计语义。虽然性能略低于累计 ack但避免了重复消费带来的数据错乱——在 AI 训练管道里一条重复的样本可能导致整个 epoch 的评估指标失真这种风险不值得冒。6.4 消息体过大的问题与方案多模态数据里经常出现几十 MB 的二进制内容。虽然 Pulsar 默认的单条消息上限是 5MB可以通过maxMessageSize调整但我不建议把视频帧或者音频直接塞进消息。正确做法是消息体只保存对象存储的路径和元数据真正的二进制内容放到 MinIO 或者云对象存储里。这样消息体积控制在 KB 级Pulsar 的吞吐优势才能发挥出来也避免 Broker 和 Bookie 承担不必要的存储压力。之前有个项目没遵循这个原则直接把 Base64 编码的图像塞进消息结果 Topic 里积压了上百 GB 数据消费端反序列化也很慢。后来改成“路径引用”模式整个管道的吞吐提升了 3 倍以上。这个教训值得分享给所有做多模态管道的人。7. Pulsar 与 Kafka 的选型对照AI Infra 场景怎么选7.1 关键维度对比很多团队在 Pulsar 和 Kafka 之间纠结我的建议是把决策维度限定在 AI Infra 的实际需求内而不是看基准测试的数字。下表是我在选型时常用对照维度PulsarKafka对 AI Infra 的影响存储模型计算存储分离Broker BookKeeperBroker 本地磁盘Pulsar 扩缩容更灵活Kafka 分区规模久了会重Topic 数量上限十万级理论值万级以内较优AI 场景下 Topic 数量多Pulsar 更抗压消息保留与回溯原生支持按时间和大小保留任意回溯通过 log retention 和 offset 重置功能偏弱训练数据重放和模型迭代需要后者订阅模式Exclusive、Shared、Failover、Key_Sharedconsumer group 为主顺序与并行冲突Agent 的场景需要灵活的消费语义运维复杂度组件多Broker、Bookie、ZooKeeper组件相对少小团队初期 Kafka 上手快Pulsar 后期收益大流量隔离存储层共享但 Topic 间积压隔离性好分区之间共享 Broker 资源隔离弱平台化后多个团队共用隔离性重要7.2 什么时候还是选 Kafka我也不是无脑推 Pulsar。如果你的场景是标准的日志收集、链路追踪、离线数仓同步对 Topic 数量要求不高团队维护经验主要靠 Kafka那 Kafka 依然是一个低成本、高效率的选择。只要流量规模没有突破单集群数万分区、数据重放需求不频繁Kafka 的简单性就是优势。7.3 什么时候建议用 Pulsar出现这几个信号我建议认真考虑 Pulsar你正在做公司级的 AI 平台底座消息中枢要服务多个团队、多个业务域租户隔离是硬需求。你的数据管道有多模态数据流入不同数据源之间需要异步汇聚、对齐、重试。你在做 Agent 编排需要灵活的事件总线和多订阅模型。你预见到未来 1~2 年 Train 和 Inference 的流量会指数级增长但峰值不可预测要求系统能独立扩缩容。我自己在实践中的一个感受是Kafka 是消息队列Pulsar 在 AI Infra 场景下更像一个数据基础设施组件。它那个“发布订阅 持久化存储 灵活回溯”的组合跟 AI 管道的思维模式天然匹配。8. 最后再分享一点体会从分布式演进的角度看AI Infra 的消息中枢不是一个可以事后补的组件。等到多模态管道和 Agent 编排已经跑起来再引入一个消息层重构成本会成倍增加。我在实际项目中踩过不少坑最深的体会是消息模型的设计要高瞻远瞩一点把 Topic 规划、订阅模式、保留策略这些基础决策在第一时间做对后面会省掉很多返工的麻烦。Pulsar 在 AI Infra 场景里的价值本质上来自于它把“流的存储”和“流的计算”分开思考。多模态数据是流的集合Agent 的事件循环是流的交互分布式系统的节点协作也是流的传递。当你的底层组件具备了对“流”的良好抽象上层的 AI 应用才能更自然地生长出来。如果你们团队正在规划 AI 平台的底座架构我的建议是不要只盯着模型和算力也把消息中枢当作一个一等公民来设计。选型阶段花一周时间跑通 Pulsar 的原型验证一下多模态汇聚和 Agent 事件编排这两个核心场景远比上线后再迁移要划算。
返回列表