| 失效链接处理 |
|
Kafka 宣布正式接入AI
相关截图:
![]() 主要内容:
先认识 Kafka
相信大家都用过或者听过Kafka。比如消息队列、日志收集、实时数仓这些应用场景,Kafka都可以派上用场。我们用过一次就会发现,它干的其实是同一件事:把已经发生的事件,按顺序、可回溯地送出去。
订单付了款,传感器抖了一下,用户在直播间丢下一句弹幕。这些都是事件。Kafka 把它们写成日志,切成多个分区,让不同的下游各自按自己的速度去读。读慢了,数据还在磁盘上;读挂了,换一台机器从上次的位置接着读。生产者和消费者互不耽误,中间靠 Topic 把两边隔开。
它擅长的是高吞吐和持久化,不擅长的是“想一想”。一条弹幕该不该判成负面情绪,一张工单要不要标成高风险,这些判断要算力、要模型、要几毫秒甚至几秒。Kafka 把原材料送到工位,但是具体处理由工位上的人来决定。
这次接入,接的是哪一层
2026 年,这件事从零散实验走到了可以对照的接口上,分两头看。
社区这边,KIP-1318 提议为 Apache Kafka 提供第一方 MCP Server。它是一个独立进程,放在 tools/mcp-server,不进 Broker。对外走 MCP 这套工具协议,对内只包一层已有的 Admin、Producer、Consumer。AI 助手可以因此去建 Topic、看消费滞后、管 ACL、看 Connect,本地用标准输入输出,远程用 HTTP。目标集群是 Kafka 4.0 及以上的 KRaft 模式。实现跟踪在 KAFKA-20436,目前仍是提案,还没进正式发行版。
平台这边已经能用。Confluent 把开源 MCP Server、托管的只读 MCP Server,以及给编程助手用的 Agent Skills 做成了可上线的工具。助手可以按你现有的权限查看环境、Topic、连接器和指标。Flink SQL 里也可以调用 AI 函数:AI_COMPLETE 做文本补全,AI_TOOL_INVOKE 把 MCP 工具交给模型。数据还在流里,标签和摘要顺手写上。
两头合在一起,就是这次“正式接入”的实际含义。Kafka 没有改行去当模型运行时,它把自己的操作和数据,接到了 AI 已经在用的标准口上。
Kafka 是事件骨干,推理放在外面
开发环境里有一种很顺手的写法:消费者拉到一条消息,当场调用大模型,拿到结果再写进另一个 Topic。笔记本上跑几条,一切正常。放到生产上,这是会把自己拖死的写法。
原因不在“调用了 AI”,而在消费者什么时候才回到 poll()。Kafka 用 max.poll.interval.ms 判断这个消费者是不是还活着,默认 5 分钟。poll() 一次可以拿走一批消息,默认上限 500 条。外部大模型一次来回常常要 1 到 10 秒,还带着速率限制。500 条依次同步等待,这一批还没写完,距离上一次 poll() 早就超过 5 分钟。
协调器会认为该消费者已离开,触发消费组重平衡,分区被分给别人。别人同样同步去调模型,同样超时,于是再平衡一次。分区来回倒手,消费进度看起来就像停住了。
|


苏公网安备 32061202001004号
