Kafka Share Groups 与队列语义:KIP-932 共享消费模型

深入剖析 Kafka KIP-932 引入的 Share Groups 共享消费模型:从经典消费组的分区独占困境出发,讲解共享分区、记录级获取与确认、获取锁与锁超时、投递计数与重投上限、poison message 与死信队列模式、并发度突破分区数上限的并行消费、与消费组在位移模型和顺序性上的差异对比,以及适用场景、版本限制与生产迁移实践

Kafka 从诞生起就是「日志」而非「队列」:分区内的消息有严格顺序,一个分区同一时刻只能被消费组内的一个消费者持有。这套模型把流式处理做到了极致,却在「任务分发」场景里处处掣肘——分区数成了并行度的硬上限,一个慢任务能拖住整个分区的进度。KIP-932 带来的 Share Groups(共享组)正是为这类场景设计的第二套消费模型:它在同一份日志之上叠加了队列语义,让多个消费者可以共享同一个分区,按记录粒度获取与确认。本文从经典消费组的痛点讲到共享组的内部机制与生产落地。

1. 从消费组到队列语义

1.1 经典消费组的分区独占模型

消费组(consumer group)的核心不变式只有一条:同一个分区在同一个组内,任何时刻只被一个消费者持有。协调者负责把分区尽量均匀地分给组内成员,成员数量超过分区数时,多出来的消费者只能空转。

这个设计服务于「有序流处理」:既然一个分区只被一个消费者读,那么位移(offset)就是一份简单的单调递增游标,提交位移、重置位移、按位点回溯都变得直白。代价则是并行度与分区数被死死绑定。

1.2 分区数决定并行度上限的痛点

考虑一个典型的任务分发场景:用 Kafka 做图片转码的任务队列。假设主题有 12 个分区,转码任务耗时从 200ms 到 30s 不等。问题立刻显现:

  • 想扩到 48 个并发消费者,必须先把分区数扩到 48;而扩容分区会改变键到分区的映射,破坏既有顺序。
  • 12 个分区里若有 3 个分区被「大图任务」堵住,对应消费者线程空转等待,另外 9 个消费者即使空闲也无法接手——它们不能碰别人持有的分区。
  • 单个消费者处理慢触发 max.poll.interval.ms,分区被强制 revoke 触发再平衡,整个组跟着抖一下。

根因是分区是分配的最小单位。只要分配粒度停留在分区,负载不均是结构性的,靠调参无法根治。

1.3 队列语义的对照

RabbitMQ、AWS SQS 这类消息队列的模型完全不同:

  • 消息(记录)是分配的最小单位,消费者逐条拉取,谁有空谁拿,不存在「分区独占」。
  • 消息被取走后进入「不可见」状态,处理成功才删除;处理失败或超时则重新可见,被其他消费者取走。
  • 并发度只受消费者数量与下游能力限制,与「分区数」无关。
  • 顺序性基本不保证(除非用 FIFO 队列并接受吞吐折损)。

Kafka 的 Share Groups 就是把 SQS 这套「记录级获取 + 确认 + 重投」语义搬到了日志之上,同时保留日志的高吞吐与持久化。它不取代消费组,而是与消费组并存,让你按场景选模型。

2. Share Group 核心模型与 KIP-932

2.1 Share Group 与共享分区

一个 share group 是一组共享消费的客户端集合,它订阅若干主题。与消费组的关键区别在于:组内的分区是共享的,而不是独占的。多个 share consumer 可以同时从同一个分区读取记录,各自处理各自拿到的那些记录。

协调者仍然存在(由某个 broker 担任),但它的职责从「分配分区」变成了「分配记录」——更准确地说,是从「派发」变成了「授权」:消费者主动来要记录,协调者决定给它哪些、以及它当前是否还持有对这些记录的锁。

2.2 记录级获取

共享消费的获取是记录级的。消费者向协调者发出获取请求(acquire),协调者从订阅的分区中挑出「当前未被锁定」的记录,连同它们的锁一起返回。返回的记录进入消费者的「已获取但未确认」集合。

# 记录级获取的语义要点
# 1) 协调者按分区轮流挑记录,天然做了负载均衡
# 2) 已被别的消费者锁定的记录不会被重复派发
# 3) 消费者拿到的是「记录 + 获取锁」,锁有租期
# 4) 处理完成后必须显式确认,否则锁到期会被回收

因为派发的粒度是记录,某个消费者被大任务堵住时,其他消费者会继续从同一个分区取走后面的小任务——负载均衡从「静态分分区」变成了「动态分记录」。

值得强调的是「锁定」与「位移」的根本差别。消费组的位移是一条单调游标:它假设「游标之前的所有记录都处理完了」。共享组没有这个假设——记录之间没有隐含的先后完成关系,第 100 条可能已 ACCEPT,第 50 条还在重投。协调者维护的不是「读到哪」,而是一张记录状态表:每条记录处于「可派发 / 已锁定 / 已完成 / 已死信」四态之一。这张表是共享组一切行为(去重、重投、死信)的单一事实来源,也是它无法提供精确一次的原因——状态表与下游写入无法进入同一个事务。

2.3 ShareConsumer API

客户端侧新增了 ShareConsumer 接口(Java 客户端在 Kafka 4.0 引入,早期以 KafkaShareConsumer 形式在 3.7+ 提供预览)。它的用法与 KafkaConsumer 相似但语义不同:

Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("group.id", "image-transcode-share");          // 共享组 id
props.put("key.deserializer",   StringDeserializer.class.getName());
props.put("value.deserializer", StringDeserializer.class.getName());
props.put("share.acknowledgement.mode", "explicit");     // 显式确认

try (ShareConsumer<String, String> consumer = new KafkaShareConsumer<>(props)) {
    consumer.subscribe(List.of("image-tasks"));
    while (running) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
        for (ConsumerRecord<String, String> rec : records) {
            try {
                transcode(rec.value());                       // 业务处理
                consumer.acknowledge(rec, AcknowledgeType.ACCEPT);
            } catch (RetryableException e) {
                consumer.acknowledge(rec, AcknowledgeType.RELEASE); // 释放,让别处重试
            } catch (FatalException e) {
                consumer.acknowledge(rec, AcknowledgeType.REJECT);  // 拒绝,计入投递计数
            }
        }
    }
}

注意 acknowledge 是逐记录的,而不是提交一个位移。这决定了共享组没有「位移」这个概念——进度是按记录状态维护的。

2.4 与消费组的定位差异

一句话概括定位:消费组面向「流」,共享组面向「任务」。消费组关心的是「我读到哪了」,共享组关心的是「这条消息处理完了没」。前者天然有序、可回溯、可精确一次;后者天然并发、可重投、最终一致。两者不是替代关系,同一个主题可以被消费组和共享组同时订阅,互不干扰。

如果你还在理解消费组的基础模型,可以先回看 https://plumephp.com/kafka-consumer/ 与 https://plumephp.com/kafka-consumer-group-rebalance/,那里讲了分区分配与再平衡的完整流程。

3. 记录级确认与获取锁

3.1 acknowledge 的三种语义

共享组的确认(acknowledgement)有三种类型,对应三种处理结果:

确认类型含义后续行为
ACCEPT处理成功记录标记完成,不再投递
RELEASE暂时处理不了释放锁,立即可被其他消费者重取
REJECT处理失败释放锁并累加投递计数,达上限后进死信

ACCEPT 是正常路径;RELEASE 用于「依赖不可用、稍后重试」,它不计入投递计数,因此不能用来无限重试同一批坏消息——那会变成忙等;REJECT 才是「这条消息我处理不了」的信号,它会推进投递计数,最终触发死信逻辑。

3.2 获取锁与锁超时

每条被获取的记录都带着一把获取锁(acquisition lock)。锁的租期由协调者控制,对应参数 share.record.lock.duration.ms(默认约 30 秒,具体以 broker 端 group.share.record.lock.duration.ms 为准)。锁有两个作用:

  • 去重:锁没到期前,这条记录不会被派发给别的消费者,避免重复处理。
  • 容错:消费者崩溃后锁自然到期,协调者把记录回收,重新派发给其他消费者。

锁超时是共享组的核心健康指标。如果消费者处理一条记录的时间超过了锁租期,锁会在处理途中到期,记录被回收并可能被另一个消费者同时处理——出现重复处理。所以长耗时任务必须配合续租。

3.3 隐式确认与续租

为了避免「处理到一半锁过期」,共享组提供两种机制:

  • 续租(renew):客户端在后台线程周期性延长它持有记录的锁租期,只要消费者还活着并且还在处理,锁就不会丢。Java 客户端的 KafkaShareConsumer 默认开启自动续租。
  • 隐式确认(implicit acknowledgement):把 share.acknowledgement.mode 设为 implicit 后,记录在下一次获取请求发出时被自动确认,无需显式 acknowledge。它适合「处理逻辑简单、poll 循环即处理完成」的场景,能省掉显式调用的样板代码,但把确认时机交给了框架,语义上更接近 at-most-once 的滑动窗口。
// 隐式确认:poll 返回下一批时,上一批自动 ACCEPT
props.put("share.acknowledgement.mode", "implicit");
// 显式确认(推荐用于长耗时任务):自己决定 ACCEPT / RELEASE / REJECT
props.put("share.acknowledgement.mode", "explicit");

3.4 锁释放与消费者退出

消费者正常关闭时,客户端会主动释放它持有的所有未确认记录(相当于对每条发 RELEASE),让它们立刻回到可派发池,而不是等锁自然过期。若进程被 kill -9,则只能等锁到期。这就是为什么优雅停机在共享组里同样重要:它能把「重投延迟」从锁租期(几十秒)缩短到毫秒级。

4. 超时、重投与投递计数

4.1 投递计数 delivery count

协调者为每条记录维护一个投递计数:这条记录被派发出去多少次。计数随每次派发递增,并在被 ACCEPT 时归零(记录完成)。它有两个用途:

  • 识别 poison message:一条记录被反复投递却始终得不到 ACCEPT,说明它「有毒」——要么内容非法,要么处理逻辑对它必然失败。
  • 触发死信:计数达到上限后,协调者不再重投,转而把记录投递到配置的死信主题。

4.2 重投上限与死信主题

# broker 端相关配置(server.properties 或动态配置)
group.share.delivery.count.limit=5
# 超过该投递次数后,记录被认定为 poison message
# 共享组会把它投递到主题名加 .dlq 后缀的死信主题

# 锁租期(broker 侧默认值,客户端可按需协商)
group.share.record.lock.duration.ms=30000
# 单次获取请求最多返回多少条记录
group.share.max.records=500
# 记录在完成前的最大锁定时间上限(防止无限续租)
group.share.max.lock.duration.ms=60000

# 预建死信主题(分区数与保留策略按排查需要设置)
kafka-topics.sh --bootstrap-server broker:9092 \
  --create --topic image-tasks.dlq \
  --partitions 12 --replication-factor 3 \
  --config retention.ms=604800000

默认投递上限通常较小(如 5 次)。一旦某条记录被投递满 5 次仍未 ACCEPT,共享组会把它写入 <原主题名>.dlq 之类的死信主题,并从正常派发池中移除。这样既保证了「坏消息不会无限循环」,又保留了「事后人工排查」的可能。

需要强调的是,死信主题需要提前存在且分区数等配置合理,否则投递失败。生产上建议为死信主题单独配置保留策略与告警。

一个实用的排查闭环是:给死信主题配一条「问题记录」消费流水线,自动把死信内容、投递次数、首次与末次投递时间、关联的业务键解析出来,推送告警并归档到可检索的存储。这样毒消息不会悄无声息地躺在死信主题里无人问津——它会被转换成一张待处理的工单。投递计数的意义也正在于此:它不是用来「重试到成功」的,而是用来「判定失败并转交」的。

4.3 与 SQS 可见性超时的类比

熟悉 SQS 的读者可以把这套机制直接对号入座:

共享组概念SQS 对应概念
获取锁租期可见性超时 Visibility Timeout
ACCEPTDeleteMessage
RELEASEChangeMessageVisibility 置零
REJECT投递计数递增,达上限进 DLQ
delivery countApproximateReceiveCount
死信主题Dead Letter Queue

差别在于:Kafka 的锁由协调者集中管理,续租是显式协议;SQS 的可见性由队列服务控制,续租靠 ChangeMessageVisibility。语义高度同构,迁移成本主要落在客户端 API 上。

4.4 poison message 的处理策略

对毒消息,工程上有三条路可选,按推荐度排序:

  1. REJECT 交给死信:把「我处理不了」明确表达出来,让框架按投递上限自动隔离。最简单也最不容易出错。
  2. REJECT 后附带结构化信息:在业务侧捕获异常,把原始记录与错误堆栈一起写到一个「问题记录」主题,便于排查。
  3. RELEASE 兜底:仅用于「临时故障、稍后一定能成功」的情况,且必须有退避,否则会变成 CPU 空转的忙等循环。

千万不要用 RELEASE 处理永久性错误——投递计数不增,死信永远不触发,记录会在集群里无限打转。这是共享组最常见的生产事故。

5. 并发度与并行消费

5.1 并行度不再受分区数限制

共享组最重要的收益:消费者数量可以超过分区数。12 个分区的主题,可以起 48 个 share consumer 并发消费,协调者按记录派发,谁空闲谁多拿。并行度从「分区数上限」解放为「消费者数上限」,而消费者数只受下游资源约束。

这直接解决了 1.2 节里的三个痛点:扩容不必动分区、慢任务不会堵住整个分区、消费者数量与主题分区数解耦。

5.2 批处理与 max.poll.records

共享组的 poll 同样支持批量返回,参数与消费组类似:

  • max.poll.records:单次 poll 最多返回多少条记录。共享组里它决定了一次获取请求带走多少条锁——设得越大,单次获取的锁越多,锁管理开销越低,但单条记录的锁被闲置的时间也越长。
  • max.poll.interval.ms:两次 poll 的最大间隔。共享组里它的含义从「分区会被 revoke」变成「消费者被判定失联,其未确认记录被回收重投」。
// 面向吞吐:一次多拿、批量处理
props.put("max.poll.records", 500);
// 面向低延迟:一次少拿,快速确认
props.put("max.poll.records", 10);

5.3 多个消费者共享同一分区

共享组的派发是「按分区轮流」的:协调者会尽量均匀地从各个订阅分区中挑记录,避免某个分区被整体拖后。这意味着即使某个分区里堆着一条处理 30 秒的慢任务,同分区的其他记录仍会被派给别的空闲消费者。

但要注意:同一分区的记录仍然可能被并行处理,这恰恰是顺序性丧失的根源。如果业务对分区内顺序有硬要求,共享组不合适。

5.4 吞吐特性

共享组的吞吐模型与消费组不同:

  • 上限更高:并行度不受分区数约束,理论上可以吃满所有消费者与下游的算力。
  • 单位开销更大:记录级锁、确认、续租都带来额外的协调者往返,单条记录的固定开销高于消费组的「批量读 + 位移提交」。
  • 适合处理耗时长的场景:记录越大、处理越慢,锁与确认的开销占比越低,共享组的相对优势越明显;反之,处理极快的高吞吐流场景,消费组的批量位移模型反而更划算。

选型经验:单条处理时间在几十毫秒以上、任务间可乱序、需要弹性扩缩消费者——用共享组;要求顺序、要精确一次、单条处理极快——用消费组。

做一个粗略的容量估算来直观感受差异。假设 12 个分区、单条处理平均 2 秒,消费组最多 12 个消费者并发,理论吞吐约 6 条/秒;若把分区扩到 60(代价是键映射变化与元数据膨胀),也才 30 条/秒。而共享组起 60 个消费者,吞吐同样约 30 条/秒,但不需要动分区;再起 120 个消费者可到约 60 条/秒,并行度线性增长,瓶颈只在下游。反过来,若单条处理只有 2 毫秒,消费组批量拉取 + 位移提交的单条开销远低于共享组的记录级锁与确认往返,此时共享组反而更慢。这就是「处理越慢、共享组越占优」的定量解释:记录级协调开销是常数,处理耗时越大,它在总耗时中的占比越低。

另一个实践要点是消费者数与分区的解耦收益。消费组里消费者数超过分区数是纯粹浪费(多出的成员空转还占内存);共享组里消费者数可以远超分区数,调度器会自然把记录摊到所有消费者上。这意味着「一个主题配固定分区数、按流量动态调消费者」成为一种可行且经济的运维模式,不再需要为了扩容而做一次高风险的改分区操作。

一个容易被忽略的细节是「消费者组的伸缩惯性」。消费组扩容时,新增消费者要等一轮再平衡才能真正拿到分区,且分配粒度是整块分区,扩容收益是阶梯式的;共享组扩容时,新消费者发一次获取请求就能立刻拿到记录,扩容收益几乎是线性的。反过来,消费组缩容时若用 CooperativeSticky,未受影响的消费者不停顿;共享组缩容则只是少了一个「要记录的人」,在途记录靠锁回收,同样平滑。对于「白天扩容、夜间缩容」的弹性负载,共享组的伸缩曲线明显更贴合。

此外,共享组的「背压」表达方式也不同。消费组靠「不 poll」来背压,但一旦长时间不 poll 会被判定失联触发再平衡;共享组靠「不 acquire」来背压,协调者只是暂时不派发记录给这个消费者,不会影响其他消费者,也不会引发任何组级别的状态变化。这让共享组更适合「下游时快时慢」的场景——慢的时候自动降速,快的时候自动提速,而不必担心触发一次全组抖动。

6. 与消费组的差异对比

6.1 逐维度对比

维度消费组 Consumer Group共享组 Share Group
分配粒度分区记录
位移模型提交 offset,按位点回溯无 offset,按记录状态确认
并行度上限分区数消费者数
顺序性分区内严格有序不保证(同分区可并行)
确认粒度批次位移单条记录 ACCEPT/RELEASE/REJECT
再平衡成员变化触发分区重分配无再平衡,记录动态派发
重投靠位移回退或手动 seek锁超时或 RELEASE 自动重投
投递计数无有,用于毒消息识别与死信
死信需自行实现内置死信主题机制
精确一次支持(事务/幂等)不支持,至多近似一次

6.2 顺序性与再平衡的取舍

消费组的两个「负担」在共享组里消失了:没有再平衡——成员增减不再引起全组停顿,因为本来就没有分区要重分配;没有位移提交——进度是记录状态的副产品,不需要消费者维护游标。代价是顺序性:同分区记录可能被不同消费者并行处理,处理完成顺序与写入顺序无关。

这本质上是一次交换:用顺序性和精确一次,换来了弹性并发和队列语义。选型时先问自己「业务能否容忍乱序」——能,就考虑共享组;不能,就留在消费组。

一个快速的决策清单:

  • 记录之间是否要求按写入顺序处理?是则消费组。
  • 是否要求 exactly-once(下游写入与消费进度原子)?是则消费组。
  • 单条处理是否可能耗时超过再平衡容忍度?是则共享组。
  • 消费者数量是否可能超过分区数?是则共享组。
  • 是否存在大量「可丢弃」或「可重试」的独立任务?是则共享组。
  • 是否需要把处理失败的消息自动隔离到死信?是则共享组。

这份清单不能覆盖所有情况,但它把选型从「凭感觉」变成了「按约束逐条核对」。

6.3 exactly-once 的现状

共享组目前不提供精确一次语义。原因是记录级的获取与确认跨越多个消费者与协调者,无法像消费组那样把「位移提交」和「下游写入」放进同一个 Kafka 事务里。共享组保证的是至少一次(at-least-once):锁超时、RELEASE、消费者崩溃都可能导致重复处理。

因此使用共享组必须做到消费幂等:以记录的业务唯一键去重,或用下游的天然幂等写(如 upsert、带唯一约束的 insert)。这一点与消费组的重复消费场景同源,可以延伸阅读 https://plumephp.com/kafka-delivery-semantics/ 理解 Kafka 三种投递语义的边界。

7. 适用场景、限制与生产实践

7.1 适合共享组的场景

  • 任务队列 / 作业分发:图片转码、报表生成、邮件发送等「一任务一处理」的负载,任务之间无顺序依赖。
  • 长耗时处理:单条处理几百毫秒到几十秒,分区独占模型下容易触发再平衡,共享组的锁续租能稳定承接。
  • 突发流量削峰:需要快速扩缩消费者应对流量峰谷,而不想动分区数。
  • 从 SQS / RabbitMQ 迁移:已有队列语义的代码,迁移到共享组的心智成本最低。
  • 消费组做不了的「慢消费者」问题:某个消费者因外部依赖变慢时,共享组会自动把记录分给别人,不会拖住整个分区。

7.2 不适合与限制

  • 强顺序需求:同分区顺序处理、按 key 有序,共享组无法满足。
  • 精确一次:需要 EOS 的场景必须用消费组加事务。
  • 版本门槛:共享组依赖 broker 端支持(Kafka 3.7 起提供预览,4.0 起正式可用),旧集群无法启用;客户端 API 也在演进,早期版本接口可能不兼容。
  • 死信主题需预建:否则毒消息无法隔离。
  • 协调者压力:记录级派发让协调者承担更多协调流量,超大规模共享组需要评估 broker 负载。

7.3 关键配置参数表

参数作用建议
group.id共享组标识与消费组命名区分,避免混淆
share.acknowledgement.mode确认模式长耗时用 explicit,简单处理用 implicit
max.poll.records单次获取记录数吞吐优先调大,延迟优先调小
max.poll.interval.ms消费者失联判定大于单批最长处理时间
group.share.record.lock.duration.ms锁租期大于单条处理时间,或依赖自动续租
group.share.delivery.count.limit投递上限3~10,超过进死信
share.auto.offset.reset无位点时的起始位置一般用 earliest,避免漏读历史任务

7.4 迁移与上线注意事项

从消费组迁移到共享组,或从 SQS 迁移过来,需要注意几点:

  • 幂等先行:迁移前先把消费逻辑改成幂等,共享组至少一次语义会把重复暴露出来。这与 https://plumephp.com/kafka-vs-rabbitmq-vs-redis-streams/ 里讨论的队列选型是同一类工程问题。
  • 监控换指标:消费组的 consumer lag 在共享组里换成「未确认记录数」「平均投递次数」「死信速率」。原有的 lag 监控在共享组上语义不同,需要重建看板。

重建监控看板时,建议至少覆盖以下指标:

  • 未确认记录数(in-flight records):每个消费者当前持有但未确认的记录数量。持续接近 max.poll.records 说明处理跟不上获取速度。
  • 平均投递次数:所有记录的投递次数均值。大于 1 说明存在重投,需结合锁超时与 RELEASE 次数定位来源。
  • 锁超时次数:因租期到期被回收的记录数。这是重复处理最直接的信号,应尽量压到接近零。
  • RELEASE / REJECT 速率:RELEASE 高说明依赖频繁不可用;REJECT 高说明数据质量或逻辑有问题。
  • 死信速率:单位时间进入死信主题的记录数。任何非零值都应触发告警与人工排查。
  • 协调者派发延迟:从记录可派发到实际被取走的时延,反映消费者是否足够。

这些指标共同刻画了共享组的健康度,与消费组的 lag + rebalance 指标一一对应,是迁移时不可省略的一步。

  • 灰度共存:同一个主题可以同时被消费组和共享组订阅,可以先用共享组承接一部分任务类型,验证稳定后再全量切。
  • 优雅停机:务必实现 shutdown hook,主动释放未确认记录,把重投延迟压到最低。
  • 锁租期与处理时长的关系:上线前压测出 P99 处理耗时,确保锁租期或续租机制覆盖它,否则会周期性出现重复处理。

优雅停机的代码骨架值得固化到每个共享组消费者里:

Runtime.getRuntime().addShutdownHook(new Thread(() -> {
    running = false;              // 让 poll 循环退出
    consumer.wakeup();            // 唤醒阻塞中的 poll
}));

// poll 循环退出后,close() 会主动释放本消费者持有的所有未确认记录
try (ShareConsumer<String, String> consumer = new KafkaShareConsumer<>(props)) {
    consumer.subscribe(List.of("image-tasks"));
    while (running) {
        ConsumerRecords<String, String> recs = consumer.poll(Duration.ofMillis(500));
        for (ConsumerRecord<String, String> r : recs) {
            handle(r);            // 处理成功后按需 acknowledge
        }
    }
}   // close() 触发未确认记录的 RELEASE

如果没有 shutdown hook,进程被 SIGTERM 直接终止后,它持有的记录只能等锁租期(默认几十秒)到期才被重投,会造成一段可观测的「消费空窗」。加了 hook 之后,重投延迟从「锁租期」降到「毫秒级」,这对延迟敏感的线上任务队列是必要的。

7.5 常见坑清单

  • 用 RELEASE 处理永久错误:投递计数不涨,死信不触发,消息无限打转。永久错误必须 REJECT。
  • 锁租期短于处理时长且未续租:处理途中锁过期,记录被并发处理,重复率飙升。
  • 忘记预建死信主题:毒消息无法隔离,REJECT 后无处可去。
  • 把共享组当消费组用:需要顺序或精确一次却选了共享组,事后补幂等代价极高。
  • max.poll.records 过大:一次拿走几百条锁,其中大半在排队,锁被长时间闲置,吞吐反而下降。
  • 忽略协调者负载:共享组规模上去后,记录级协调流量可能成为 broker 瓶颈,需提前压测。

8. 总结

Share Groups 给 Kafka 补上了缺失已久的第二块拼图:在日志的有序流模型之外,叠加了一套记录级的队列语义。它的核心是把分配粒度从分区下沉到记录——协调者动态派发记录并授予获取锁,消费者按条确认(ACCEPT / RELEASE / REJECT),锁超时触发重投,投递计数驱动死信。这套机制换来的收益是并行度脱离分区数、弹性扩缩、慢任务不堵分区;代价是顺序性丧失、精确一次不可用、单条开销上升。

选型判断可以浓缩成一句话:要顺序和 EOS,用消费组;要弹性并发和任务分发,用共享组。两者并存而非互斥,同一个主题可以让不同业务按各自语义消费。落地时记住三条铁律——消费逻辑必须幂等、永久错误必须 REJECT、锁租期必须覆盖处理时长。把这三条守住,共享组就能在任务队列场景里稳定地替你扛住分区独占模型扛不动的负载。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. Kubernetes 上的 Kafka:Strimzi Operator 生产实践
  2. ksqlDB 流式 SQL:流表模型、窗口聚合与生产运维
  3. Kafka 消费延迟诊断:Lag 定位、分区倾斜与治理