- Published on
Kafka Share Groups: From Streams to Queues
- Authors

- Name
- Charles Chen
Kafka 长期以来最鲜明的特征,是把数据组织成有序 partition,再由 Consumer Group 按 partition 分配给消费者。这种模型非常适合日志、CDC 和流计算,但并不适合所有“任务队列”场景。
Kafka 4.2 的发布公告宣布 Kafka Queues(Share Groups)达到 production-ready,并带来了 RENEW acknowledgement、adaptive batching、消费数量限制和更完整的 lag metrics。Apache Kafka 4.2 发布公告
这不是给 Consumer Group 增加一个小配置,而是 Kafka 在消费语义上的一次扩展:Kafka 现在可以同时承担 Stream 和 Queue 两类工作。
1. 为什么传统 Consumer Group 不总是适合任务队列
传统 Consumer Group 的基本关系可以简化为:
Topic
├── Partition 0 ── Consumer A
├── Partition 1 ── Consumer B
└── Partition 2 ── Consumer C
同一个 Consumer Group 中,一个 partition 在同一时刻只由一个 consumer 处理。因此,常见的并行度上限近似为:
consumer parallelism <= partition count
这不是缺陷,而是 Kafka 用于保持 partition 内顺序、进行可预测扩展的设计。Kafka 官方设计文档也明确说明,consumer group 通过 partition assignment 换取顺序和可扩展性。Kafka Consumer Group 设计
问题出现在任务耗时高度不均匀时。例如一个 partition 中依次出现:
Task A: 100 ms
Task B: 2 s
Task C: 60 s
Task D: 200 ms
如果 Task C 阻塞了当前 consumer,Task D 即使只需要 200 ms,也要等在它后面。这是典型的 head-of-line blocking。
对于日志解释、OCR、embedding、模型推理、异步 API 调用以及图片和视频处理,任务处理时间经常由外部服务决定。此时,固定的 partition ownership 往往不是最自然的调度模型。
2. Share Groups 带来了什么
Share Group 是与 Consumer Group 并列的一种 group 类型。多个 share consumer 共同处理同一批 topic 的记录,partition 可以同时分配给多个 consumer:
Kafka topic
│
▼
Share Group
├── Worker A
├── Worker B
├── Worker C
└── Worker D
Kafka 4.2 的 KafkaShareConsumer 文档说明,share group 中的 consumer 数量可以超过 topic 的 partition 数,并且多个 consumer 可以共享同一个 partition;代价是不能再把传统的 partition 顺序保证直接套用到记录处理上。KafkaShareConsumer API
Share Group 的关键特征包括:
| 维度 | Consumer Group | Share Group |
|---|---|---|
| 主要模型 | Partition-oriented stream | Record-oriented queue |
| 一个 partition 的成员 | 通常一个 consumer | 可以有多个 share consumer |
| 并行度上限 | 受 partition 数约束 | consumer 数量可超过 partition 数 |
| 处理完成方式 | 提交消费位置 | 对记录进行 acknowledgement |
| 顺序语义 | partition 内顺序是核心能力 | 不应依赖全局或严格 partition 顺序 |
| 典型场景 | 日志、CDC、流计算 | 独立任务、工作队列、异步处理 |
它更接近 RabbitMQ、Amazon SQS 或 Celery 的工作队列语义,但并不意味着 Kafka 完全变成了这些系统。Share Group 仍然运行在 Kafka 的 topic、partition 和 broker 架构之上。
3. Kafka 4.2 的生产化能力
3.1 RENEW 适合长耗时任务
Share consumer 获取记录时,记录会处于一个有时限的 acquisition lock 中。Kafka 4.2 的 AcknowledgeType.RENEW 表示记录仍在处理中,消费者可以持续续租这个锁。AcknowledgeType API
这对 60 秒甚至更长的 Agent task 很重要:worker 可以在处理期间定期调用 RENEW,避免记录在任务尚未完成时被其他 worker 再次投递。KafkaShareConsumer 的 RENEW 示例
但 RENEW 只是延长处理租约,不是“执行成功”的确认,也不等于 exactly-once。worker 崩溃、网络分区、租约过期或 acknowledgement 丢失时,记录仍然可能再次投递。
3.2 Adaptive batching 降低协调开销
Kafka 4.2 为 group coordinator 和 share coordinator 引入 adaptive batching。系统根据工作负载自动调整 batch linger time,减少固定批量等待带来的延迟和吞吐折中。Kafka 4.2 发布公告
这解决的是 broker 协调与批处理效率问题,不是业务任务本身的调度策略。应用仍然需要根据任务时延、worker 数量、外部 API 限流和重试成本设置合理的批大小。
3.3 消费数量限制与 lag metrics
4.2 增加了对 fetched records 数量的软限制和严格限制,并补充了 share partition lag 的持久化和查询能力。管理员可以通过 kafka-share-groups.sh 查看 share group 的起始 offset 和 lag,也可以查看成员与 group state。Share Group 运维命令
监控侧还提供按 acknowledgement 类型统计的指标,例如 Accept、Release、Reject 和 Renew 的速率。Share Group 监控指标
因此,Share Group 的生产化不只是“能消费记录”,还包括了长任务续租、协调批处理、消费边界和运行观测这几块必要的运维能力。
4. 为什么 Agent 系统更需要 Queue 模型
假设 Kafka 收到 flink-error-events,每条事件触发一个独立任务:
Kafka
↓
Share Group
↓
Agent workers
├── 调用 LLM 解释日志
├── 查询 Milvus
├── 写入 Prometheus
└── 调用 Flink REST API
不同任务可能分别耗时 100 ms、2 s、15 s 或 60 s。Share Group 允许 worker pool 更细粒度地领取记录,从而把扩展单位从“partition”变成“可独立完成的任务”。
这个模型特别适合以下条件:
- 每条记录可以独立处理,或者依赖关系已经被显式建模。
- 任务处理时间差异很大。
- 允许记录在失败后被重新投递。
- 外部副作用可以通过幂等键或去重状态保护。
- 业务更关心任务完成率、重试次数和队列 lag,而不是严格的 partition 顺序。
反过来,如果业务依赖事件顺序、按 key 聚合,或者需要从一致的状态快照恢复,就不应该因为 Share Group 更像队列而强行迁移。
5. Share Groups 不替代 Flink
Kafka Share Groups 和 Flink SlotSharingGroup 也不是一回事。前者属于 Kafka 的消费模型,后者属于 Flink 的任务资源调度。
| 对比项 | Kafka Share Groups | Flink SlotSharingGroup |
|---|---|---|
| 所属组件 | Kafka | Flink |
| 解决问题 | 哪个 worker 消费哪条记录 | 哪些 operator subtask 可以共享 slot |
| 关注点 | 消费与 acknowledgement 语义 | 任务部署与资源分配 |
| 是否改变数据消费语义 | 是 | 否 |
| 是否决定状态计算 | 否 | 否,主要影响部署布局 |
| 典型用途 | 独立任务、异步 worker pool | 控制 operator 的 slot sharing |
Flink 官方文档说明,SlotSharingGroup 用于定义哪些任务可以共享 slot;slot sharing 能提高资源利用率,但不是完整的 CPU 隔离机制。Flink 作业调度与 SlotSharingGroupFlink 架构中的 slot 说明
边界可以这样记:
Kafka Share Groups = Message Scheduling
Flink SlotSharingGroup = Task Placement
Flink Stateful Processing = State + Time + Recovery
因此,调用 LLM 做一次日志解释可以放进 Share Group;但 10 分钟窗口聚合、event-time join、CEP、CDC stateful transform 和依赖 checkpoint 恢复的 exactly-once 流处理,仍然应该由 Flink 负责。Flink 的 checkpoint 会保存输入位置和 operator state,并在恢复时从一致位置重放数据。Flink Stateful Stream Processing
6. 最大的生产风险:重复副作用
队列系统最容易被误解的地方,是“消息被确认”不等于“外部动作只执行了一次”。
下面这个流程仍然有重复执行窗口:
Kafka
↓
Worker
↓
POST /restart-job 第一次请求成功
↓
ACK 丢失或 worker 崩溃
↓
Kafka 再次投递
↓
POST /restart-job 第二次请求再次执行
Kafka 官方设计文档在讨论 acknowledgement 时也指出,消费者完成处理后、发送 acknowledgement 前发生故障,会造成重复处理。Share Group 的 acknowledgement 机制管理的是记录交付状态,不会自动替外部 HTTP API、Shell 命令或数据库副作用建立跨系统事务。
生产 worker 至少应围绕下面的状态设计:
Kafka record
↓
Worker
↓
Idempotency Store
├── idempotency_key
├── task_id
├── attempt
├── status
└── result_hash
↓
External Side Effect
↓
ACCEPT / RELEASE / REJECT
推荐做法是:
- 用稳定的
idempotency_key绑定业务动作,而不是只使用 Kafka offset。 - 将
task_id、attempt、状态和结果摘要持久化,能够判断重试是新任务还是重复投递。 - 对外部 API 使用服务端幂等键;对无法幂等的 Shell 或管理操作增加分布式锁、状态检查和人工保护阈值。
- 把
RENEW当成租约续期机制,不把它当成成功确认。 - 分开监控
Renew、Release、Reject、lag、处理时延和外部副作用成功率。
不要采用下面这种没有边界的链路:
Kafka → LLM → Shell
它没有幂等控制、没有副作用审计,也无法清楚区分“任务处理中”“任务失败”和“动作已经成功但 acknowledgement 丢失”。这类系统迟早会把一次网络抖动变成凌晨三点的重复操作事故。
7. 选择模型的决策表
| 问题 | 更适合 Consumer Group / Flink | 更适合 Share Group |
|---|---|---|
| 是否依赖 partition 或 key 顺序 | 是 | 否或顺序可放弃 |
| 是否需要窗口、join、CEP 或状态 | 是 | 否 |
| 任务耗时是否高度不均匀 | 一般 | 是 |
| worker 数是否需要超过 partition 数 | 不适合 | 是 |
| 是否能承受重复投递 | 视 sink 契约而定 | 必须设计幂等 |
| 主要关注点 | 状态一致性与事件时间 | 任务完成、重试与队列 lag |
实际系统可以同时使用两种模型,而不是二选一:Flink 负责从业务事件中识别需要执行的任务,Kafka Share Group 负责把独立任务分发给 Agent worker,结果再回到 Kafka 或其他服务层。
8. 总结
Kafka 4.2 的 Share Groups 让 Kafka 在保留 Stream 能力的同时,正式拥有了更成熟的 Queue 消费模型:
- Consumer Group 仍然适合有序 partition 消费、CDC 和流计算。
- Share Group 适合耗时不均、可独立执行、需要动态 worker pool 的任务。
RENEW解决长耗时记录的 acquisition lock 续期问题,但不提供外部副作用的 exactly-once。- adaptive batching、消费数量限制和 lag metrics 让 Share Group 更接近可运维的生产能力。
- Flink
SlotSharingGroup是资源布局概念,不是 Kafka 消费模型,也不能替代 Flink 的状态、时间和恢复语义。 - Agent worker 必须配合幂等键、任务状态、重试审计和外部副作用保护。
最重要的判断不是“Kafka 现在能不能做队列”,而是每条数据到底需要哪一种语义:需要状态和时间,就交给 Flink;需要把独立任务可靠地分发给一组耗时不同的 worker,就考虑 Share Group。