Published on

Kafka Share Groups: From Streams to Queues

Authors
  • avatar
    Name
    Charles Chen
    Twitter

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 GroupShare Group
主要模型Partition-oriented streamRecord-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 类型统计的指标,例如 AcceptReleaseRejectRenew 的速率。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”变成“可独立完成的任务”。

这个模型特别适合以下条件:

  1. 每条记录可以独立处理,或者依赖关系已经被显式建模。
  2. 任务处理时间差异很大。
  3. 允许记录在失败后被重新投递。
  4. 外部副作用可以通过幂等键或去重状态保护。
  5. 业务更关心任务完成率、重试次数和队列 lag,而不是严格的 partition 顺序。

反过来,如果业务依赖事件顺序、按 key 聚合,或者需要从一致的状态快照恢复,就不应该因为 Share Group 更像队列而强行迁移。

Kafka Share Groups 和 Flink SlotSharingGroup 也不是一回事。前者属于 Kafka 的消费模型,后者属于 Flink 的任务资源调度。

对比项Kafka Share GroupsFlink SlotSharingGroup
所属组件KafkaFlink
解决问题哪个 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_idattempt、状态和结果摘要持久化,能够判断重试是新任务还是重复投递。
  • 对外部 API 使用服务端幂等键;对无法幂等的 Shell 或管理操作增加分布式锁、状态检查和人工保护阈值。
  • RENEW 当成租约续期机制,不把它当成成功确认。
  • 分开监控 RenewReleaseReject、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 消费模型:

  1. Consumer Group 仍然适合有序 partition 消费、CDC 和流计算。
  2. Share Group 适合耗时不均、可独立执行、需要动态 worker pool 的任务。
  3. RENEW 解决长耗时记录的 acquisition lock 续期问题,但不提供外部副作用的 exactly-once。
  4. adaptive batching、消费数量限制和 lag metrics 让 Share Group 更接近可运维的生产能力。
  5. Flink SlotSharingGroup 是资源布局概念,不是 Kafka 消费模型,也不能替代 Flink 的状态、时间和恢复语义。
  6. Agent worker 必须配合幂等键、任务状态、重试审计和外部副作用保护。

最重要的判断不是“Kafka 现在能不能做队列”,而是每条数据到底需要哪一种语义:需要状态和时间,就交给 Flink;需要把独立任务可靠地分发给一组耗时不同的 worker,就考虑 Share Group。