• Kafka 成了 Agent 的「共享内存」,Flink 成了它的「大脑」


    Confluent 官方博客里举了个例子:一位客户写信来问订单到哪了,客服 Agent 查了账户,看到状态是「已发货」,就回了句安抚的话。

    问题是,那笔订单四十分钟前就取消了。

    Agent 没犯任何它看得见的错——它只是基于一份六小时才刷新一次的快照在推理。这不是模型不够聪明,是它手里的上下文是 stale 的。

    大多数 AI 项目卡住,根子都在这:Agent 想「聪明」,但给它喂的是一堆过期的、割裂的记忆。

    那如果把记忆这件事,交给数据流本身来做呢?

    一、先把旧问题摆清楚:Agent 的「记性」是断的

    传统 Agent 的两块「记忆」都存在问题:

    • 短期记忆(对话上下文、刚查到的数据)——存在进程内,崩一次就没了,多实例之间还不同步。
    • 长期记忆(业务知识、历史决策)——存在某个库里,Agent 每次要「去问」,数据更新了它也未必知道。

    所以才会出现「订单取消了还回已发货」这种尴尬:记忆里那份快照早就过期了,Agent 却把它当真相。

    Confluent 这季度在 Confluent Cloud 的 Apache Flink 上原生提供 Streaming Agents,做的就是把这两块记忆都接到数据流上——让 Agent 的记性,跟着数据一起流动。

    注意这句话的差别:这不是「把你的 Agent 连到 Kafka」,而是 Agent 直接嵌进 Flink 流里,变成常驻算子。

    以前 Agent 挂在管线外面,定时去拉数、调一次模型、返回结果,管线对它是个黑盒。现在是事件一来,它在流里就地推理,结果直接回流进流。

    具体技术面其实挺厚:

    • 在 Flink SQL 里直接对远程 LLM 做原生推理;
    • 持续生成 RAG 向量,让实时事件流本身就是检索语料;
    • 通过 MCP 调工具(查库、调 API、改配置);
    • 内置异常检测和 Auto-ARIMA 预测,直接在时序流上跑。

    一句话:Flink 不再只是搬数据,它成了 Agent 的「大脑」——一边接收最新鲜的事件,一边当场出决策。

    pipeline

    三、Kafka 成了「共享内存」:多 Agent 靠事件流协作

    单个 Agent 嵌进流还不够,真正有意思的是多 Agent 怎么协作。

    Confluent 把 A2A 协议(Agent-to-Agent,已经是开放协议)直接接进了 Flink。于是多 Agent 编排变成一套很「微服务」的结构:

    • Kafka 当短期共享内存(short-term shared memory);
    • Flink 当实时路由(real-time routing);
    • 每个 Agent 是一个带脑子的有状态微服务,彼此之间没有硬编码依赖。

    这就是标题那半句的意思:Agent 之间不再靠代码里写死的「A 调 B」来通信,而是把状态和消息丢进 Kafka 这个共享内存,谁上线、谁下线、谁扩容,互不影响。

    memory

    四、可重放:这套「记忆」最值钱的地方

    Streaming Agents 里,每一次输入都不可变地落日志,每一次决策都可重放。

    这一个特性,同时解决三件让 Agent 难以上生产的事:

    1. 失败恢复——崩了不用从头跑,从断点重放;
    2. 逻辑测试——同样的输入喂进去,跑一百遍验证 Agent 逻辑;
    3. 决策审计——某次自动审批怎么做的,顺着事件流完整回溯。

    对传统 Agent 框架而言,决策是个黑盒,出了问题只能从头再跑,还不一定复现。Streaming Agents 把「黑盒」拆成了「一串可回放的不可变事件」。

    对金融、医疗、政务这类强审计、强合规的场景,这条是上生产的硬门槛,不是锦上添花。

    replay

    五、工程细节:为什么 KIP-932 是刚需

    Agent 的特性是它的负载是突发、并行的,一个触发器能瞬间炸出成百上千个并行任务。

    而 Kafka 老的那套 1:1 分区-消费者约束,是给「可预测的人规模吞吐」设计的:一个分区只能被一个消费者独占。Agent 那种突发 burst 一来,直接卡死。

    Confluent 顺手把 KIP-932 Share Groups 带进了 Streaming Agents——弹性 many-to-many 消费,打破 1:1 硬限制。

    对 Agent 工作负载来说,这不是优化,是能跑起来的前提。


    回到开头:客服 Agent 在过期数据上回了一句「已发货」。把记忆交给数据流,听起来很美,但你现在的系统落在哪一关?评论区扣个字母:

    • A Agent 还完全在管线外,定时拉数、黑盒决策,stale 数据是常态。
    • B 已经让 Agent 碰实时流了,但决策不可重放,出问题只能从头跑。
    • C 多 Agent 想协作,但还在代码里硬编码谁调谁,扩一个就改一片。

    顺便说说:你最想让 Agent 接管的那条实时业务流,是什么?

  • 相关阅读:
    IDEA:commit 提交插件
    android——自定义加载按钮LoadingButton
    N1中openwrt实现不插网线就能上网,通过wifi连接路由器
    Qt6.3学习笔记 QFrame::HLine与QFrame::VLine 改变线条颜色
    javascript复习之旅 13.1 模块化(上)
    3-3、python中内置数据类型(集合和字典)
    LeetCode75——Day16
    requests库中r.content 与 r.read() 的使用方式
    使用 L293D 电机驱动器 IC 和 Arduino 控制直流电机
    python字典按照 值进行排序 sorted
  • 原文地址:https://www.cnblogs.com/Jackeyzhe/p/22888221