Agent Communication · Research Brief

让 Agent 协作
消息送得达,
任务接得住、过程可追溯

围绕跨进程、异步长任务与多会话并发,比较 6 个开源候选,选择可恢复、可维护的通信投递底座。

Scope & architecture

协议约定任务,投递层保障消息流转

本次选型聚焦 Agent 之间的传输与投递基础设施,服务跨进程、异步长任务、多会话并发。大模型、Agent 框架与 MCP / A2A 协议本身的选型不在本次范围内。

治理与审计
权限判定 · 授权上下文 · 追踪与审计由治理体系与业务服务负责
协议语义
A2A 任务描述 · MCP 工具调用约定任务、能力与交互语义
本次选型
消息队列 / 事件总线 · 传输与投递持久化、投递确认、会话隔离、保序与重放
运行时编排
持久化工作流 · Actor · 进程内消息负责执行编排与进程内协作
使用边界

长任务跨服务运行时,需要独立投递层

网络中断需要恢复,下游变慢需要缓冲,会话并发需要隔离。任务协议与消息队列分层协作,才能支撑这些要求。

同一进程内的协作

优先使用 Agent 框架自带消息机制,无需额外引入中间件。

仅使用直连 HTTP / SSE 时,持久化、失败重放、会话保序与背压需要额外机制承接。

Requirements to evaluation

六项业务需求,对应六个评价维度

从项目希望达成的结果出发描述需求,再落到可比较的评价口径。需求 R01–R06 与评分 D1–D6 一一对应,采用连续编号便于查阅。

关键需求业务侧希望达成的结果对应评价口径权重
R01接入既有 Agent 协作链路既有任务描述、工具调用和能力发现能够接入投递底座,尽量减少专用接口与语义转换。D1 · A2A / MCP 协议兼容性是否原生支持;需要多少协议映射与适配工作。20%
R02长任务执行结果有保障断链或重启后任务不静默丢失,执行结果能够确认送达;重复投递不造成重复业务副作用。D2 · 可靠性与投递语义投递保证、事务能力、终态确认与幂等支撑。25%
R03多会话独立推进,故障后可续接并发会话互不干扰,同一会话或授权链路保持因果顺序,并能从指定位置恢复消息历史。D3 · 会话隔离与可重放会话粒度通道、保序与断点重放。20%
R04参与方身份可信,敏感内容受控核验发送方身份与授权范围,防止冒充;按机构内或跨机构边界保护消息内容。D4 · 安全与加密加密层级、身份校验、防冒充与细粒度权限。15%
R05核心依赖稳定,有持续维护基础关键能力具有稳定版本、持续发行与实际采用基础,避免把演示能力或频繁变动接口作为长期依赖。D5 · 成熟度与生态版本稳定性、发行节奏与采用规模。10%
R06在一年周期内完成生产验证与交接结合适配开发、专项验证与运维准备,判断 12 个月内能否达到生产可用并完成交接。D6 · 一年内落地风险达到生产可用的把握;高分表示风险较低。10%

从已有能力到项目可交付,分开判断

D5 成熟度:看方案当前的版本、维护和采用基础。D6 落地风险:看本项目能否在一年内完成集成、验证与上线准备。成熟方案未必接入成本低,新能力也不能仅凭设计契合就进入生产。

评分方法:六维整数权重,统一五分制

六维权重依次为 20%、25%、20%、15%、10%、10%,合计 100%。总分 = Σ(维度得分 × 权重百分比)/ 100;展示值四舍五入至两位小数。技术维度采用能力档位;成熟度与落地风险分别按依赖稳定和项目交付判断。半分是定性判断,适用前提见评分依据。

5 原生支持,开箱即用
4 原生支持,少量配置
3 需要定制开发
2 能力显著不足
1 不支持或引入风险
20%25%20%15%10%10%
依次为 D1 协议 / D2 可靠性 / D3 会话 / D4 安全 / D5 成熟 / D6 落地
Weighted comparison

六个候选,一个有边界的选型结论

按六维整数权重重新计算,满分 5 分。排名用于比较本场景的适配程度;协议、持久化与生产准入条件需结合逐项依据理解。

加权总分 · 五分制
012345
生产主线快速验证跨机构预研对照 / 观察

六维评分矩阵

数字保留原报告评分;点击任一评分,即可在下方依据区域查看对应说明。颜色按得分区间展示,半分保留原值。高分越接近绿色,能力缺口越接近橙红色。

5≥4.543.5–432.5–321.5–21≤1
六个候选方案的六维定性评分与加权总分
候选方案 / 所处层级D1
协议
20%
D2
可靠
25%
D3
会话 / 重放
20%
D4
安全
15%
D5
成熟
10%
D6
落地
10%
总分
/ 5

安全维度:RocketMQ 的 3.5 分反映其链路加密边界;SLIM 的 5 分来自报告所强调的消息层端到端加密能力。评分不等同于安全认证或生产验收结果。部分高分带有额外集成前提,详见逐项评分依据。

能力轮廓 · 同一尺度比较优势与短板
RocketMQ
逐维对比

各维度满分均为 5。雷达图展示原始得分,权重仅用于加权总分。

Scores & evidence

选一个方案,逐项看分数与依据

下拉选择候选方案,再点击左侧评分项,查看同一维度的评分理由、资料依据与适用前提。条形长度按五分制展示,权重只用于加权总分。

4.5–5 较高3.5–4 良好2.5–3 需补齐1.5–2 缺口较大
六项评分 点击查看对应依据

条形长度 = 得分 / 5
D6 分数越高,表示落地风险越低。

以上依据由现有报告与官方资料整理,属于定性判断;部分评分依赖额外集成。阅读完整 36 项说明 ↗
Core features · one case · distinctive code

六个方案,各有明确的特色与取舍

先了解核心特点、优势与局限,再看它们如何完成同一个材料核验案例。

Apache RocketMQ · LiteTopic

RocketMQ 的核心特点是用 LiteTopic 为每个会话建立轻量子通道:大量会话可以共用父主题,各自组织消息、订阅与消费位置。优势是会话进度、顺序与断线续读围绕同一标识建立,适合异步长任务;RocketMQ-AI 还提供协议集成入口。局限是 LiteTopic 与集成模块较新,需要固定兼容版本;跨机构正文保密和外部业务幂等仍需应用补充。

NATS JetStream + Agent Protocol

NATS + Agent Protocol 的核心特点是用 subject 寻址,把请求、分块响应与流中反问串成一次 Agent 交互。优势是核验过程可以先向客户追问,再沿原响应流继续执行,交互组织清楚。局限是 Agent Protocol 自身采用 Core NATS,不能自动提供持久任务与历史恢复;可靠长任务需接入 JetStream 并保存等待状态,既有 A2A / MCP 语义也需要映射。

AGNTCY SLIM

SLIM 的核心特点是围绕加密会话与群组成员组织协作,借助 MLS 保护消息正文,中转数据面按名称路由。优势是跨机构参与者可以在受控群组中交换内容,邀请或移除成员成为明确操作。局限是安全会话不能直接替代持久任务账本、离线结果留存和审计历史;密钥、身份互信、绑定版本及配套存储仍需管理。

Apache Kafka

Kafka 的核心特点是把业务过程写成可保留的追加日志,让不同消费组各自推进位置。优势是编排、审计、分析可以独立读取同一批事件,单个读者暂停不会改变其他读者的进度,历史也可重新消费。局限是分区与位点不是客户会话专属通道,Agent 协议、会话筛选与状态恢复需要适配;新建平台还涉及分区、容量和运维建设。

Redis Streams

Redis Streams 的核心特点是通过消费组记录待确认任务,再把超时未完成项转交给其他 Worker。优势是追加、领取、PEL、接管与确认都有直接命令,已有 Redis 的团队容易搭建轻量任务流。局限是持久性取决于 AOF / RDB、复制和故障切换配置;恢复循环、历史裁剪、协议接口与业务幂等需要自行组织,缓存经验不能直接替代任务治理。

RobustMQ · mq9

mq9 的核心特点是把 Agent 的能力发现、专属邮箱与消息优先级放入同一组接口。优势是编排者可以先找会做核验的 Agent,再投递到其邮箱,接收者暂时离线仍可在留存窗口内取件。局限是一个 Agent 邮箱可能混合多个会话,任务状态与审计还需应用组织;RobustMQ 官方仍标记早期开发、尚未生产就绪,当前更适合隔离预研。

一个具体案例:100 份材料的协同核验

“请核验这 100 份客户材料。缺页时先问我,加急件优先处理;需要跨机构复核时保护材料正文。我断线后还要继续看进度,核验过程要留给审计回看。”

客户会话记作 S-001…S-100,单个核验任务记作 T-001。以下分别说明整体实现路径,代码着重展示各方案最有特色的调用。

下载完整特色与代码文档 ↗

Apache RocketMQ · LiteTopic生产主线 · 加权分 4.68 / 5

把并行核验结果分到各自会话,重连后继续读。

01 · 这个案例怎样实现

编排服务先通过受控目录选择核验 Agent,将 T-001 等请求投递到任务主题;材料正文以受控引用或应用加密负载传递。核验端把“已开始、缺签章页、已完成”写入 materials.progress 父主题下的 S-001 LiteTopic,S-002 等会话使用其他子通道。缺页时,应用记录等待状态并投递补件问题,收到答复后续做;加急件可使用独立任务队列或应用调度。网页重连后订阅原会话并从保存位置续读,审计服务在保留期内读取事件。特色集中在会话消息的归属与恢复入口;能力发现、优先调度与跨机构正文保护由配套服务完成。

STEP 01共用父主题materials.progress
STEP 02独立子通道S-001 与 S-002
STEP 03各写各的进度缺页 / 完成分别归档
STEP 04定向订阅客户 1 只看 S-001
02 · 核心代码

下面同样发送三条消息,但 setLiteTopic 直接决定它们属于哪一个会话通道。

适用前提 · Java gRPC SDK;provider 与 producer 是已初始化的 ClientServiceProvider 和 Producer,producer 可发送到预建的 Lite 父主题 materials.progress。官方 LiteTopic 文档要求 broker ≥ 5.5.0、gRPC SDK ≥ 5.1.0。下方是方法片段,需置于相应类内并导入 SDK 类型;消费端另用 LitePushConsumer 订阅所需子通道。

黄色边线标出直接体现方案特色的调用。

关键机制:父主题相同,LiteTopic 随会话改变JAVA · 核心片段
void publishProgress(ClientServiceProvider provider, Producer producer,
                     String sessionId, String payload) throws ClientException {
    Message event = provider.newMessageBuilder()
        .setTopic("materials.progress")
        .setLiteTopic(sessionId)
        .setBody(payload.getBytes(java.nio.charset.StandardCharsets.UTF_8))
        .build();
    producer.send(event);
}

// 以下三次调用置于已初始化的业务方法中
publishProgress(provider, producer, "S-001", "缺签章页;event_seq=1");
publishProgress(provider, producer, "S-002", "核验完成;event_seq=1");
publishProgress(provider, producer, "S-001", "已收到补件;event_seq=2");
03 · 如何理解这些调用
01 · setTopic:共用基础资源

三个事件都发到 materials.progress;父主题由部署侧预建。此处的规模设计不需要 100 个普通 Topic。

02 · setLiteTopic:真正区分案卷的一行

第 1、3 条进入 S-001,第 2 条进入 S-002。消费端需要使用对应子通道订阅;仅在负载里写 session_id 不会自动获得这种通道划分。

03 · event_seq:业务顺序仍由应用定义

同一子通道默认单队列有序,但多个生产者、业务重试和状态变更仍需稳定事件 ID / 序号。示例文字用于看清归属,正式系统应使用结构化事件。

04 · 恢复读取:还需要消费状态

网页重连后由应用重新订阅 / 查询自己的通道,并根据保留期与消费位置恢复。LiteTopic 命名本身不完成前端重连、任务状态查询或工具幂等。

中断恢复与结果一致性

客户网页断线后,可由应用重新关联 S-001 的子通道并恢复读取,实际范围受消息保留期、子通道 TTL 和消费位置影响。执行 Agent 退出是另一种故障,需要任务重新投递、业务状态恢复与副作用去重;重新订阅进度不能替代这些处理。

  • LiteTopic 的隔离是通道和订阅粒度;不代表每个任务独享 CPU、磁盘或完整租户权限。单通道默认单队列,也有吞吐限制。
  • 子通道 TTL 与消息保留期决定可恢复窗口,过期后不能无限重放。任务重试时仍需业务幂等;把结果写入与任务确认衔接好,避免丢终态或重复执行。
NATS JetStream + Agent Protocol快速验证 · 加权分 4.05 / 5

发现缺页时,在当前响应流中反问并继续。

01 · 这个案例怎样实现

编排方通过 Agent Protocol 的发现与请求契约联系核验端。处理 S-001 时缺少签章页,核验端发送 query,附带 reply_subject,询问“补充后继续,还是先核验其余材料”;编排方将客户选择送回这个地址,Agent 沿当前交互继续响应。若客户明天才回复,等待状态必须持久化;任务和进度接入 JetStream,发布获得确认,消费在结果保存后 ACK,网页与审计按留存策略续读。加急调度和跨机构正文保护仍由应用实现。下面分别展示流中反问与持久消费,二者需要应用适配衔接。

STEP 01开始核验既有响应流已建立
STEP 02发 query缺页 → 给出两种选择
STEP 03回复到指定地址reply_subject 收到选择
STEP 04继续响应流执行剩余核验步骤
02 · 核心代码

第一段直接展示 Agent Protocol 的 query;第二段解释持久任务通道如何另行适配。

适用前提 · 第一段使用已连接的 nats-py nc;response_subject 是请求建立的响应地址,处理器已发送协议要求的首个 status ack。调用方需解析 query 并回复指定 subject。该片段只封装中途提问,后续 response 分块与结束帧由请求处理器发送。第二段需要开启 JetStream,流配置仅为首次部署示例。

黄色边线标出直接体现方案特色的调用。

关键机制:在响应流中反问,而不是重新提交任务PYTHON · 核心片段
import json
import uuid

async def ask_about_missing_page(nc, response_subject):
    reply_subject = nc.new_inbox()
    answers = await nc.subscribe(reply_subject)
    await nc.flush()  # 先确认已订阅回复地址
    query = {"type": "query", "data": {
        "id": str(uuid.uuid4()),
        "reply_subject": reply_subject,
        "prompt": "缺签章页:补充材料,还是先核验其余部分?",
    }}
    await nc.publish(response_subject, json.dumps(query).encode())
    try:
        answer = await answers.next_msg(timeout=30)
        return answer.data.decode()  # 本例调用方回复纯文本选择
    finally:
        await answers.unsubscribe()
先保存结果,再确认任务PYTHON · 核心片段
import json
import nats
from nats.js.api import RetentionPolicy

async def inspect_once():
    nc = await nats.connect("nats://localhost:4222")
    js = nc.jetstream()
    await js.add_stream(name="TASKS", subjects=["tasks.>"])
    await js.add_stream(name="EVENTS", subjects=["events.>"],
                        retention=RetentionPolicy.LIMITS)
    task = {"task_id": "T-001", "session_id": "S-001"}
    await js.publish("tasks.inspect", json.dumps(task).encode(),
                     headers={"Nats-Msg-Id": task["task_id"]})
    worker = await js.pull_subscribe("tasks.inspect", durable="inspect-workers")
    for msg in await worker.fetch(1):
        task = json.loads(msg.data)
        result = {"task_id": task["task_id"], "status": "completed"}
        await js.publish("events." + task["session_id"],
                         json.dumps(result).encode(),
                         headers={"Nats-Msg-Id": task["task_id"] + ":done"})
        await msg.ack()
    await nc.drain()
03 · 如何理解这些调用
01 · type=query:双方能识别的反问

缺页问题进入当前响应流。query 不是终态,也不关闭响应流;Agent 收到答案后可继续产生 response。

02 · reply_subject:把问题与答案对上

每次提问建立独立 inbox,先订阅再发布问题。调用方回复到这个地址;本例使用规范允许的纯文本回复。

03 · 等待与状态:协议之外的应用职责

30 秒是本例的等待时间。长期人工补件需存下问题、任务状态、截止时间和恢复入口;超时并不自动意味着业务任务失败。

04 · JetStream:把“可对话”与“可恢复”连接起来

第二段的 durable、发布确认与 ACK 属于持久投递层,需与协议处理器适配。Nats-Msg-Id 只在配置窗口内帮助发布去重,外部副作用仍需 task_id 幂等。

中断恢复与结果一致性

Core NATS 的这次临时反问不会因进程重启自动恢复。应用保存等待状态,或把任务与事件写入 JetStream,再实现重新关联请求与回复。第二段先保存结果再 ACK,可减少终态遗漏,但仍有重复执行窗口;长任务还需配置 ACK 等待 / 进度确认。

  • Agent Protocol v0.3 文档将 JetStream 支持的至少一次投递与端到端加密列在其范围之外;持久流代码需要作为独立适配层接入。
  • S-001 subject 只是逻辑分类,隔离还需要账户和 subject 权限。协议中的 query 可表达补充问题,但人工审批状态、期限与审计仍由应用实现。
AGNTCY SLIM跨机构预研 · 加权分 3.88 / 5

机构协作围绕“谁是加密群组成员”展开。

01 · 这个案例怎样实现

材料核验由应用编排并记录 T-001 / S-001,机构 A 与 B 的 Agent 建立启用 MLS 的 GROUP 会话。疑点需要机构 C 复核时,先验证 C 的身份与授权,再邀请它加入,等待成员操作完成后交换材料摘要;复核结束后移除 C,再继续后续消息。中转负责路由,正文由端点处理。缺页追问、加急排序、网页断线续读和审计通过配套任务状态与事件存储实现;安全会话的可靠投递不能替代这些持久机制。代码突出成员加入 / 移除如何改变后续加密协作范围。

STEP 01建立 GROUP + MLSA / B 安全会话
STEP 02邀请授权复核者等待 C 加入完成
STEP 03群组内交换正文中转负责路由
STEP 04移除复核者等待完成后继续交流
02 · 核心代码

用群组成员的加入与移除,直接展示 SLIM 的会话设计中心。

适用前提 · local_app、conn_id、group_name 由初始化过程提供;group_name 为 slim.Name 类型。端点身份、凭据、密钥与路由已配置,受邀 Agent 已监听兼容会话。调用形式参照官方 group.py;业务侧已确认受邀者的授权资格。

黄色边线标出直接体现方案特色的调用。

关键机制:MLS 群组 + 邀请 / 移除成员PYTHON · 核心片段
from datetime import timedelta
import slim_bindings as slim

async def consult_reviewer(local_app, conn_id, group_name):
    cfg = slim.SessionConfig(
        session_type=slim.SessionType.GROUP,
        max_retries=5, interval=timedelta(seconds=5), metadata={},
        mls_settings=slim.MlsSettings(
            header_integrity_validation_percent=100,
            max_seen_control_message_ids_size=None),
    )
    ctx = local_app.create_session(cfg, group_name)
    await ctx.completion.wait_async()
    session = ctx.session
    # 示例邀请两位已获业务授权的跨机构成员
    for name in ["bank-b/verify/agent", "bank-c/review/agent"]:
        peer = slim.Name.from_string(name)
        await local_app.set_route_async(peer, conn_id)
        invited = await session.invite_async(peer)
        await invited.wait_async()
    await session.publish_async("S-001:请复核示例摘要".encode(), None, None)
    # 实际流程应等待复核结果,再执行下面的移除
    reviewer = slim.Name.from_string("bank-c/review/agent")
    removed = await session.remove_async(reviewer)
    await removed.wait_async()
03 · 如何理解这些调用
01 · GROUP + mls_settings:安全会话单位

创建群组并显式配置 MLS;正文保护还取决于成员端点的相容配置、身份与密钥。这里只给会话核心,不替代完整安全初始化。

02 · invite_async:参与方变化成为会话操作

先配置路由,再邀请成员,等待完成后发送正文。set_route_async 解决“怎么送到”,业务授权解决“是否可以邀请”。

03 · publish_async:向已建立会话交流

正文进入群组会话;中转节点负责路由。路由名称与其他元数据不因正文加密而自动隐藏。

04 · remove_async:影响后续交流的成员集合

移除操作完成后继续会话。移除不能收回成员此前已经接收的内容,也不能自动撤销该 Agent 在业务数据库中的权限。

中断恢复与结果一致性

会话重试参数处理传输中的重试。进程重启后要找回昨天的任务与审计事件,仍需应用建立持久状态和历史账本;MLS 的正文保护与长期消息保存分别实现。

  • 群组操作需要与业务授权、会话生命周期和密钥管理衔接;代码中的受邀资格是假定已完成的前置判断。
  • 示例突出成员变更,省略复核结果的接收与会话关闭。生产代码应等待业务结果、处理超时并释放会话资源。
Apache Kafka日志基线 · 加权分 3.48 / 5

一次记录核验过程,编排、审计和分析各自读取。

01 · 这个案例怎样实现

应用网关将核验请求与状态映射为事件,按 session_id 写入 materials.events。编排、审计、分析分别使用自己的消费组:编排读到缺页事件后向客户发问并保存等待状态,收到补件事件再推进任务;审计保留过程,分析计算缺页率。审计维护时暂停不会推进其他组的 offset,恢复后可接着读,新的回放组可重读保留历史。网页按会话筛选进度并恢复业务位置;加急件由独立主题或调度处理,客户权限与跨机构正文保护由网关负责。key 帮助会话内顺序,但不会自动建立客户级权限或独立消费通道。

共享事件源写入一份日志key = session_id
独立读者 1编排组推进下一项任务
独立读者 2审计组独立保存 / 回看过程
独立读者 3分析组独立计算缺页率
02 · 核心代码

下面用三个不同的 group.id 读取同一主题;最后建立新的回放组,展示“自己的书签”。流程图后三格是并行读者。

适用前提 · 已部署 Kafka,预建 materials.events,安装 confluent-kafka;本例假设各消费组是首次使用,消息仍在保留期内。auto.offset.reset=earliest 仅在没有有效提交位点时生效;已有消费组要显式重置或定位 offset。需另补 poll 循环、成功后的 commit、错误处理与 close。

黄色边线标出直接体现方案特色的调用。

关键机制:同一份事件,四个独立的消费位置PYTHON · 核心片段
import json
from confluent_kafka import Producer, Consumer

p = Producer({"bootstrap.servers": "localhost:9092"})
p.produce("materials.events", key=b"S-001",
          value=json.dumps({"status": "missing_page", "event_seq": 1}).encode())
p.flush()

def open_reader(group_id):
    c = Consumer({
        "bootstrap.servers": "localhost:9092",
        "group.id": group_id,
        "enable.auto.commit": False,
        "auto.offset.reset": "earliest",
    })
    c.subscribe(["materials.events"])
    return c

orchestrator = open_reader("orchestrator")
audit = open_reader("audit")
analytics = open_reader("analytics")
# 新组无旧位点,从当前保留的最早事件开始核对;不移动其他组位点
replay = open_reader("audit-replay-B-001-new")
03 · 如何理解这些调用
01 · key:局部有序的组织依据

S-001 作为事件键,在既定分区策略下使同会话事件进入同一分区。key 不是独立会话通道,也不单独提供会话级权限。

02 · 不同 group.id:分别读,而不是抢一条消息

三个服务的消费组不同,所以每组可读取同一份事件。同组的多个实例则分担分区,不能用同一个 group.id 来实现三个系统各读一份。

03 · offset:一个组自己的书签

业务处理成功后提交该组位点;audit 停止推进不会拖动 orchestrator 的位点。代码省略 poll / commit,实际处理器必须实现。

04 · 新的回放组:回看不干扰生产读取

audit-replay-B-001-new 没有旧提交位点,earliest 从保留窗口的起点读。应用还需按批次 / 会话筛选;数据已过期时不能再从日志找回。

中断恢复与结果一致性

审计服务重启可从本组已提交位点续读,未提交部分可能再读。对 Kafka 内“读事件、写结果、提交输入位点”的链路可用事务,并配合 read_committed;模型调用、外部 HTTP 与数据库副作用仍需应用幂等。

  • 顺序保证属于分区,不是全局顺序。增加分区或变更分区器会影响同键路由,历史与新事件需按迁移方案处理。
  • 按 key 分类不是独立的会话 ACL。若每会话都建主题 / 分区会增加资源负担;主题权限、会话访问控制和协议映射需要自行设计。
Redis Streams轻量基线 · 加权分 2.88 / 5

核验 Worker 退出后,从待确认列表接管任务。

01 · 这个案例怎样实现

编排服务把 T-001 与 S-001 写入 tasks:inspect,核验 Worker 通过 XREADGROUP 领取。缺页时,业务保存等待状态并通知客户;有补件后生成续做任务。Worker A 如果领取后退出而没有 XACK,任务留在 PEL;Worker B 查看待确认状态,按合理空闲阈值 XAUTOCLAIM,重新核验,先将结果写入 events:S-001 再确认输入。网页和审计从会话流的保存 ID 续读,保留窗口由裁剪策略控制。加急队列、协议入口、跨机构正文保护与副作用幂等由应用补充;命令恢复逻辑和 Redis 持久配置共同决定故障行为。

STEP 01A 领取后退出任务留在 PEL
STEP 02B 查看待确认XPENDING
STEP 03B 接管超时项XAUTOCLAIM
STEP 04处理后确认保存结果 → XACK
02 · 核心代码

直接展示故障接管路径;inspect 是业务提供的核验函数,需要按 task_id 幂等。

适用前提 · r 为 Redis(decode_responses=True) 连接。部署时预建 inspect-workers 消费组,按需要选择起始 ID;例如 XGROUP CREATE tasks:inspect inspect-workers 0-0 MKSTREAM。XAUTOCLAIM 要求 Redis ≥ 6.2。示例省略认证、持久化配置与错误处理。

黄色边线标出直接体现方案特色的调用。

关键机制:待确认任务转交给 worker-bPYTHON · 核心片段
def recover_pending(r, inspect, idle_ms):
    # r 为 Redis(decode_responses=True),消费组已预建
    pending = r.xpending("tasks:inspect", "inspect-workers")
    cursor = "0-0"
    while True:
        claim = r.xautoclaim(
            "tasks:inspect", "inspect-workers", "worker-b",
            min_idle_time=idle_ms, start_id=cursor, count=10)
        cursor, entries = claim[0], claim[1]
        for msg_id, task in entries:
            result = inspect(task)  # 业务按 task_id 防止重复副作用
            r.xadd("events:" + task["session_id"], result)
            r.xack("tasks:inspect", "inspect-workers", msg_id)
        if cursor == "0-0":
            break
    return pending  # 本轮接管前的待确认概况
03 · 如何理解这些调用
01 · PEL / XPENDING:看见“领走但未确认”

XREADGROUP 读取后进入待确认列表,XPENDING 提供待确认概况。Worker 离线不意味着条目自动被另一个 Worker 完成。

02 · XAUTOCLAIM:转移待确认归属

符合 min_idle_time 的任务转交给 worker-b,返回的 entries 才是本次接管的消息。idle_ms 应结合正常任务耗时与进度维护配置,过短可能抢走仍在运行的任务。

03 · cursor:接管扫描也需要遍历

返回游标供本轮继续扫描,0-0 表示到达末尾。应用要周期性启动下一轮恢复,并监控没有活跃执行者的任务。

04 · 先结果再 XACK:接管不等于完成

结果写成功后确认任务。两条操作之间崩溃可能导致重复处理;XACK 只清该组的待确认状态,不删除 Stream 正文,业务副作用仍需去重。

中断恢复与结果一致性

结果 XADD 成功、输入 XACK 前退出,会导致再次处理并可能产生重复结果。这里的两条命令没有跨业务操作事务;需应用去重。能否在 Redis 故障后恢复,还取决于 AOF / RDB、复制与部署配置。

  • Streams 提供数据结构与消费机制,持久性要结合部署配置判断。不要仅凭使用 XADD 就宣称任务绝不会丢失。
  • XTRIM、XDEL 和容量策略会影响历史回放。用 events:<session> 分键组织数据,还需考虑 key 数量、保留策略与访问权限。
RobustMQ · mq9观察项 · 加权分 2.85 / 5

先找能核验的 Agent,再给其邮箱投递加急任务。

01 · 这个案例怎样实现

核验 Agent 注册“材料完整性核验”能力与邮箱地址;编排方 DISCOVER 后,从返回且获授权的候选中选择执行者。S-001 的补件任务加急,发送到其邮箱时标记 urgent。接收者暂时离线,消息在邮箱留存窗口内等待;上线后 FETCH,处理并保存结果,再 ACK。应用以 task_id / session_id 关联反问、补件、终态与网页进度,并另建客户会话历史供审计;跨机构正文保护也由应用实现。优先级类别不构成实时处理承诺。代码演示原生命令,标准 A2A 接入走 mq9.a2a facade,生产资格仍受上游开发状态限制。

STEP 01DISCOVER 能力找能核验材料的 Agent
STEP 02选择实际邮箱agent.inspector
STEP 03SEND + urgent离线时待办留邮箱
STEP 04上线后 FETCH处理并保存 → ACK
02 · 核心代码

原生协议命令:用 Bash + NATS CLI 展示注册、发现和投递,需要支持 mq9 的 broker。

适用前提 · 需要自行部署支持 mq9 的 RobustMQ 与 NATS CLI,并把 CLI 指向该 broker。下面的地址为本机教学示例。CLI 请求是 mq9 原生协议,不是 A2A 的标准 HTTP 请求。

黄色边线标出直接体现方案特色的调用。

找得到 Agent,也能把任务送入邮箱BASH · 核心片段
nats --server nats://localhost:4222 request '$mq9.AI.AGENT.REGISTER' \
  '{"name":"agent.inspector","mailbox":"agent.inspector","payload":"材料完整性核验"}'
nats --server nats://localhost:4222 request '$mq9.AI.AGENT.DISCOVER' \
  '{"semantic":"材料完整性核验","limit":3}'
nats --server nats://localhost:4222 request '$mq9.AI.MAILBOX.CREATE' \
  '{"name":"agent.inspector","ttl":3600}'
nats --server nats://localhost:4222 request '$mq9.AI.MSG.SEND.agent.inspector' \
  --header 'mq9-priority:urgent' \
  '{"task_id":"T-001","session_id":"S-001","task":"inspect"}'
nats --server nats://localhost:4222 request '$mq9.AI.MSG.FETCH.agent.inspector' \
  '{"group_name":"inspect-workers","deliver":"earliest"}'
03 · 如何理解这些调用
01 · REGISTER / DISCOVER

注册时公布能力描述与邮箱地址。示例先看发现结果,再演示发送到已知邮箱;生产代码应从实际返回结果中选择,不能始终硬编码同一个地址。

02 · MAILBOX.CREATE / SEND

创建有 TTL 的邮箱,按 Agent 地址投递。邮箱存储待处理消息,业务 session_id 随负载传递,用于把多个任务的结果关联起来。

03 · FETCH / ACK

FETCH 读取后,执行者应先处理并保存结果,再请求 $mq9.AI.MSG.ACK.agent.inspector;请求体包含 group_name、mail_address 和本次 FETCH 返回的 msg_id,不能填固定示例 ID。

04 · 优先级与协议

urgent 表示优先级类别,不构成硬实时 SLA。若需要标准 A2A,还需采用并核对 mq9.a2a SDK 接入路径;原生命令自身并不表达完整 A2A 任务状态机。

中断恢复与结果一致性

执行者下线后再上线读取邮箱,是该设计试图解决的核心场景。邮箱 TTL、确认位点和消息保留策略决定可恢复范围;正式采用前还需以固定版本验证重启、重复投递和集群故障行为。

  • RobustMQ 官方仓库明确标注项目仍处早期开发、尚未达到生产就绪。这里只用于解释设计特色,不能把演示可用扩大为已通过生产验收。
  • 一个 Agent 邮箱可能混有多个会话。多会话独立访问、长期事件账本、任务终态以及外部副作用幂等仍需应用设计。
Decision

以 RocketMQ 为主线,按场景补充能力

选型优先考虑协议接入、可靠投递、会话恢复与一年内的生产落地把握,结合验证结果决定进入生产的范围。

生产投递底座

RocketMQ · LiteTopic

承接 Agent 之间的任务与结果投递,重点验证会话隔离、终态确认和断点恢复。

早期快速验证

NATS · JetStream

用较低部署成本跑通端到端;保留协议适配边界,避免早期接口变化进入生产关键路径。

跨机构预研

AGNTCY SLIM

针对跨组织消息保密验证端到端加密、身份与成员撤销,再决定专用通道方案。

其余候选的采用边界

mq9 继续观察。Kafka 已有集群可评估复用,新建系统不优先引入;Redis Streams 不作为本报告建议的生产主线。

Sources & glossary

调研依据与关键术语

技术能力、协议与版本状态的主要资料入口如下。评分属于调研判断,版本、许可与能力以实际采用版本的官方文档及验证结果为准。

常用术语,一句话说明
至少一次投递

允许重复送达,消费方必须处理幂等。

会话级隔离

每个会话独立流转,限制单会话故障影响。

传输层加密

保护通信链路,中转终结链路后可能读取内容。

消息层加密

保护消息本身,内容只有授权端点能够读取。

重放

从历史指定位置重新投递消息流,用于恢复与复现。

终态确认

执行结果持久化成功后,才确认任务消费完成。

优先把协议接入、可靠投递与故障恢复做实

生产主线采用 RocketMQ,NATS 支撑快速验证,SLIM 为跨机构保密预留技术路径。最终进入生产的范围以专项验证和验收结果确认。