AI Infra 下的消息中枢:Apache Pulsar 如何支撑多模态与 Agent 编排
2026/9/16 1:52:58 网站建设 项目流程

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 下的消息中枢,需求清单已经发生了变化:

  1. 海量 Topic 支持。AI 场景下每个数据源、每个模型任务甚至每个用户会话都可能需要独立的主题,Topic 数量轻松上千上万,传统的主题较少时优势明显的方案容易吃力。
  2. 流量削峰和积压隔离。训练任务是周期性的,推理流量是突发的,消息系统必须能扛住长时间积压,并且某个 Topic 的积压不能拖垮其他 Topic 的消费。
  3. 存储与计算独立扩缩容。数据量增长但算力有限,或者算力充足但存储紧张,这两种情况都要能独立调整。
  4. 跨地域复制和多租户隔离。AI 平台通常服务多个团队,租户之间的数据隔离和权限控制必须原生支持。
  5. 稳定的顺序语义和消息回溯。训练样本的顺序在某些场景下会影响模型收敛,消息系统需要支持按序消费和从任意位置重放。

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 特征流与样本流的消息建模

在多模态处理中,我倾向于把消息分成三类流:

  1. 样本流(Sample Stream):用于训练和评估的完整样本,包含文本、图像路径、音频路径、标注信息等。这类消息体积较大,但频率相对低,适合用批量生产的方式写入。
  2. 特征流(Feature Stream):模型中间层的输出向量,用于向量检索、多模态对齐或者实时推荐。这类消息体积较小,但吞吐极高,是 Pulsar 最擅长的场景。
  3. 事件流(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 Topic,Agent 启动时从指定 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.conf)

  • managedLedgerDefaultEnsembleSize:Segment 副本数,默认 2。
  • managedLedgerDefaultWriteQuorum:写入需要的确认副本数,默认 2。
  • managedLedgerDefaultAckQuorum:写入完成的最小确认数,默认 2。
  • maxUnackedMessagesPerConsumer:单个消费者未确认消息数上限,默认 50000,要根据消费逻辑调整。

BookKeeper 配置(bookkeeper.conf)

  • journalMaxWriteRequests:写入队列长度,影响写入吞吐。
  • 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/v1
  • ai-platform/feature-extract/image/embedding/v1
  • ai-platform/agent/task/event/v1

这么设计的好处是:租户隔离清晰、数据类型清晰、业务域清晰、版本号清晰。后续做权限控制、数据保留策略、监控告警时,都可以按前缀匹配,不需要逐条配置。版本号放在末尾,是为了兼容管道升级时的双跑场景——新旧版本同时运行,互相不干扰。

5.4 生产端与消费端的工程实现要点

生产端我用 Java 客户端写了一个比较典型的异步发送逻辑:

Producer<String> 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(); CompletableFuture<MessageId> 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 支持积压,但磁盘总有上限。排查思路是:

  1. 先看积压大小:bin/pulsar-admin topics stats-internal {topic},关注backlogSize字段。
  2. 判断积压原因到底是消费者挂了还是消费速度跟不上。消费速度跟不上就需要增加消费者并行度。
  3. 如果积压时间较长,且旧消息已经不需要处理,可以直接用pulsar-admin topics skip跳过积压,或调低保留策略快速清理旧数据。

记住一个原则:积压隔离是 Pulsar 的优势,但积压本身还是需要监控和告警,不要让“能积压”变成“不处理”的借口。

6.2 消费顺序错乱如何排查

用 Shared 模式消费时,消息处理顺序天然不保证。如果你发现某个多模态对齐任务的结果错乱,排查步骤:

  1. 确认生产端是否用同一个分区键(Message Key)。
  2. 确认消费端是否使用 Key_Shared 订阅。
  3. 检查 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 + BookKeeper)Broker 本地磁盘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:

  1. 你正在做公司级的 AI 平台底座,消息中枢要服务多个团队、多个业务域,租户隔离是硬需求。
  2. 你的数据管道有多模态数据流入,不同数据源之间需要异步汇聚、对齐、重试。
  3. 你在做 Agent 编排,需要灵活的事件总线和多订阅模型。
  4. 你预见到未来 1~2 年 Train 和 Inference 的流量会指数级增长,但峰值不可预测,要求系统能独立扩缩容。

我自己在实践中的一个感受是:Kafka 是消息队列,Pulsar 在 AI Infra 场景下更像一个数据基础设施组件。它那个“发布订阅 + 持久化存储 + 灵活回溯”的组合,跟 AI 管道的思维模式天然匹配。

8. 最后再分享一点体会

从分布式演进的角度看,AI Infra 的消息中枢不是一个可以事后补的组件。等到多模态管道和 Agent 编排已经跑起来,再引入一个消息层,重构成本会成倍增加。我在实际项目中踩过不少坑,最深的体会是:消息模型的设计要高瞻远瞩一点,把 Topic 规划、订阅模式、保留策略这些基础决策在第一时间做对,后面会省掉很多返工的麻烦。

Pulsar 在 AI Infra 场景里的价值,本质上来自于它把“流的存储”和“流的计算”分开思考。多模态数据是流的集合,Agent 的事件循环是流的交互,分布式系统的节点协作也是流的传递。当你的底层组件具备了对“流”的良好抽象,上层的 AI 应用才能更自然地生长出来。

如果你们团队正在规划 AI 平台的底座架构,我的建议是:不要只盯着模型和算力,也把消息中枢当作一个一等公民来设计。选型阶段花一周时间跑通 Pulsar 的原型,验证一下多模态汇聚和 Agent 事件编排这两个核心场景,远比上线后再迁移要划算。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询