Skip to content
Charles Shao
Go back

Kafka Consumer 深挖:poll 循环、offset 提交与 rebalance,怎么做到不丢不重不卡顿

views

前两篇立好了地基:Kafka 是一份分布式可重放日志Producer 用分区策略把消息高效攒批塞进去。写入侧讲透了,这一篇补齐对称的另一半——消息怎么被可靠、高效、不抖动地读出来

Consumer 看似只是个 while(true) { poll(); process(); } 循环,但生产里几乎所有”消息丢了""消息重复了""消费突然卡住""一发布就整组停摆”的事故,根子都在这个循环的细节里。这篇把它拆开讲透。

TL;DR

Table of contents

Open Table of contents

1. 消费者的核心循环:poll 到底做了什么

消费者代码简单得有点欺骗性:

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
    for (ConsumerRecord<String, String> r : records) process(r);
    consumer.commitSync();
}

poll() 远不止”拉一批消息”。它是消费者与集群交互的唯一心跳节拍器,一次 poll 背后至少做了:拉取消息、推进 rebalance(加入/退出组、领取分区分配)、拉取元数据、(旧版本里)触发自动提交。几乎所有协调动作都挂在 poll——这解释了后面很多”反直觉”的行为。

Kafka 消费者 poll 循环与心跳线程模型图。上半部分是『应用线程(唯一拉取线程)』的 while(true) 循环:poll(timeout) 拉一批记录 → process(records) 处理这一批 → commit 提交 offset → 再回到 poll,形成闭环。下半部分是『心跳线程(0.10.1+ 后台独立线程)』:按 heartbeat.interval.ms 定时发送 heartbeat。两者之间标注 poll 还负责推进 rebalance、拉元数据。图下方黄色注解说明两条存活判定各管一段:①心跳线程按 heartbeat.interval.ms 发心跳,超过 session.timeout.ms 没收到心跳就判定消费者掉线;②应用线程两次 poll 间隔超过 max.poll.interval.ms 就判定假死、主动离组触发 rebalance;因此单批处理别太久、别在 poll 之外阻塞太久。右下角粉色注解区分 position(下次要读的位置,随 poll 前进)与 committed offset(已确认处理到,随 commit 前进),两者之差就是崩溃时的重复/丢失风险窗口。

1.1 一个消费者实例只能一个线程用

KafkaConsumer 不是线程安全的。多线程共享一个实例会直接抛 ConcurrentModificationException。这是刻意的设计:把”拉取 + 位点管理”约束在单线程里,语义最简单。想并行,见 §9。

1.2 position vs committed offset:风险窗口在哪

回顾核心原理篇 §3.3:要严格区分两个位置。

两者之间的差,就是崩溃时消息重复或丢失的风险窗口——你在哪一步提交、提交了什么,直接决定语义。这就是下一节的主题。

2. offset 提交:四种姿势与不丢不重的时机

“消息会不会丢/会不会重”不取决于 Kafka,取决于你在处理的哪一步提交 offset

offset 提交时机与交付语义对比图,分上下两块。上块『先提交,后处理 → at-most-once(可能丢,绝不重)』:poll 拉到消息 → commit offset → process 处理,注解指出若在 commit 之后、process 之前崩溃,位点已推进,这批消息再也不会重发,于是丢失。下块『先处理,后提交 → at-least-once(绝不丢,可能重)』:poll 拉到消息 → process 处理(写库/计费)→ commit offset,注解指出若在 process 之后、commit 之前崩溃,位点没推进,这批消息会被重投造成重复,需靠消费端幂等(唯一键/去重表)兜住。底部提示:生产默认采用 at-least-once + 消费端幂等;自动提交(enable.auto.commit=true)是按 auto.commit.interval.ms 定时后台提交,介于两者之间且时机不可控,对不丢要求高的链路应关掉改手动提交。

2.1 自动提交:省心,但时机不可控

默认 enable.auto.commit=true:消费者会在 poll 时,按 auto.commit.interval.ms(默认 5s)把上一次 poll 返回的位点后台提交掉

自动提交最大的坑:它提交的是”已经 poll 出来的”位点,而不是”已经处理完的”位点。如果你 poll 出一批、处理到一半崩了,但这批的位点已被自动提交——没处理完的那部分就丢了(滑向 at-most-once)。反之处理慢、还没到 5s 就崩,又会重复(at-least-once)。时机完全不由你控制,所以对可靠性敏感的链路一律关掉它。

2.2 commitSync:可靠、阻塞、会重试

关掉自动提交后,最直接的是先处理、后同步提交

props.put(ENABLE_AUTO_COMMIT_CONFIG, false);
// ...
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
    for (ConsumerRecord<String, String> r : records) process(r);
    consumer.commitSync();   // 阻塞直到 broker 确认;失败会自动重试
}

commitSync() 提交本批最后一条的下一个 offset,阻塞等待 broker ack,遇可重试错误会自动重试。缺点是每批都要等一次网络往返,高吞吐下会拖慢消费。

2.3 commitAsync:快,但失败不重试

consumer.commitAsync((offsets, ex) -> {
    if (ex != null) log.warn("commit failed for {}", offsets, ex);
});

异步提交不阻塞,吞吐更高。代价:失败不自动重试(重试可能把更小的 offset 覆盖掉更大的、造成回退)。所以生产里的经典组合是:

try {
    while (true) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
        for (ConsumerRecord<String, String> r : records) process(r);
        consumer.commitAsync();     // 常态:异步,不阻塞主循环
    }
} finally {
    try { consumer.commitSync(); }  // 关闭/异常前:同步兜底,确保最后一次落定
    finally { consumer.close(); }
}

2.4 提交指定位点:更细粒度的控制

commitSync 默认提交整批。想处理完一部分就提交一部分(缩小风险窗口),可显式指定:

Map<TopicPartition, OffsetAndMetadata> toCommit = new HashMap<>();
for (ConsumerRecord<String, String> r : records) {
    process(r);
    toCommit.put(new TopicPartition(r.topic(), r.partition()),
                 new OffsetAndMetadata(r.offset() + 1));   // 注意:提交的是"下一条"= 当前 offset + 1
    if (someBatchBoundary) consumer.commitSync(toCommit);
}

一个高频 off-by-one 坑:提交的 offset 语义是”下次从这里开始读”,所以要提交 已处理的最大 offset + 1,而不是 offset 本身。提交成 offset 会导致重启后重复消费最后一条

3. poll 循环的三个”生死参数”

消费者被判”死”的方式有两种,对应三个参数。搞混它们是”无限 rebalance”事故的根源。

参数默认谁在看超了会怎样
heartbeat.interval.ms3s心跳线程发心跳的频率(应 < session 的 1/3)
session.timeout.ms45s(新版)Group Coordinator这段时间没收到心跳 → 判消费者掉线,触发 rebalance
max.poll.interval.ms5minGroup Coordinator两次 poll 间隔超它 → 判”假死”,主动离组

关键在于 0.10.1 之后心跳被拆到独立后台线程:即使你的处理逻辑正卡着,心跳照样在发,session.timeout.ms 这条线不会误杀你。真正卡死你的是另一条线——

max.poll.interval.ms:如果单批消息处理太久(比如一批 500 条,每条要写一次慢下游),迟迟不回来调下一次 poll,Coordinator 就认为你”假死”、把你踢出组 → rebalance → 分区给别人 → 你处理完回来发现位点没了、又要重新加入 → 恶性循环,整组越来越慢

对策三选一(常组合):

  1. 减小 max.poll.records(默认 500):一次少拉点,缩短单批处理时间。
  2. 调大 max.poll.interval.ms:给慢处理留足预算。
  3. 把慢处理异步化(§9),让主循环快速回到 poll

4. Rebalance 深挖:分区分配、生命周期与 Listener

核心原理篇 §5 讲了 rebalance 的 Eager vs Cooperative 两种协议。这里从消费者视角补齐三件事:分区按什么策略分给消费者(assignor)、rebalance 在你代码里触发哪些回调、以及怎么把停摆降到最小。

4.1 分区分配策略(assignor):谁拿到哪些分区

消息进哪个分区由 Producer 的分区策略决定(见 Producer 篇);写进去之后,这些分区再分给消费者组里的哪个消费者——这由消费侧的 partition assignment strategypartition.assignment.strategy)决定。它不改吞吐/延迟的批处理逻辑,却直接决定消费端的负载均衡与 rebalance 代价

Kafka 消费者三种分区分配策略对比图,横向三块。第一块『RangeAssignor(默认)』:一个 Topic T 有 4 个分区 P0 P1 P2 P3,按段切分——C1 拿到前段 P0、P1,C2 拿到后段 P2、P3;旁注:多 topic 时各 topic 的同号分区容易堆到同一个消费者上,导致负载倾斜。第二块『RoundRobinAssignor』:同样 4 个分区被逐个轮流分配——C1 拿 P0、P2,C2 拿 P1、P3;旁注:跨 topic 拉平、最大化消费者利用率,但 rebalance 时分配可能大幅变动。第三块『(Cooperative)StickyAssignor』:分配结果同样均衡,C1 拿 P0、P1,C2 拿 P2、P3;旁注:分配均衡的同时,在 rebalance 时尽量保留原有归属,减少分区迁移与状态重建,是推荐选择。底部横贯一条共同铁律:组内一个分区只归一个消费者;消费者数量超过分区数时,多出来的消费者空转待命作为故障备份。

Assignor分配逻辑特点 / 取舍
Range(默认之一)每个 topic 单独按 分区数 / 消费者数 连续分段,除不尽时靠前的消费者多拿一个简单;但多 topic 时同号分区易堆到同一消费者 → 倾斜
RoundRobin把所有分区逐个轮流发给消费者跨 topic 拉平、最大化消费者利用率;rebalance 时改动可能较大
Sticky / CooperativeSticky像 RoundRobin 一样求均衡,但尽量保留已有分配rebalance 时分区迁移最少、状态重建最少(推荐,展开见 §4.3)
Custom继承 AbstractPartitionAssignor、重写 assign()按优先级 / 机房等自定义分配逻辑

默认值:3.0+ 的默认是 [RangeAssignor, CooperativeStickyAssignor](先用 Range,具备条件时可平滑切到 Cooperative)。无论哪种 assignor,一条铁律不变(核心原理篇 §3.6):组内一个分区只归一个消费者;消费者数量多于分区数时,多出来的消费者会空转待命(充当故障切换的热备)。

4.2 Rebalance 生命周期与 ConsumerRebalanceListener

选定 assignor 决定”怎么分”;而真正触发一次 rebalance 时,你的代码里会依次触发哪些回调、每个回调该做什么,才是不丢不重的关键

Kafka 消费者 rebalance 生命周期时序图,参与方为 Consumer、ConsumerRebalanceListener、Group Coordinator(broker)。流程:触发条件是成员增减、订阅变化或分区数变化;① Group Coordinator 在消费者下次 poll 时通知需要 rebalance;② Consumer 回调 onPartitionsRevoked(即将失去的分区),注解说明应在此提交 offset、flush 本地状态,Cooperative 模式只回调被迁走的少数分区;③ Consumer 向 Coordinator 发 JoinGroup 重新入组、上报订阅;④ Coordinator 选出 leader consumer 执行 assignor 分配;⑤ Consumer 发 SyncGroup 领取新分配方案;⑥ Consumer 回调 onPartitionsAssigned(新分到的分区),注解说明应在此 seek 到起始位点、重建本地缓存;⑦ Consumer 恢复消费。最后注解:异常掉线未能优雅收尾时会回调 onPartitionsLost。

订阅时挂一个 listener,就能在分区被收走 / 分到时插入自己的逻辑:

consumer.subscribe(List.of("ad-events"), new ConsumerRebalanceListener() {
    @Override
    public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
        // 即将失去这些分区:把已处理进度提交掉,避免换主后重复
        consumer.commitSync(currentOffsets);
        flushLocalState(partitions);   // 有状态消费者:把本地聚合/缓存落盘
    }
    @Override
    public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
        // 刚分到这些分区:需要的话 seek 到自定义位点、预热本地缓存
        rebuildLocalState(partitions);
    }
    @Override
    public void onPartitionsLost(Collection<TopicPartition> partitions) {
        // 非正常失去(如会话超时):分区可能已被别人接管,这里做兜底清理即可
        discardLocalState(partitions);
    }
});

onPartitionsRevoked 是防重复的关键点:在分区被收走之前把位点提交掉,接管者就能从正确的地方接着读,而不是重放一大段。有状态消费者(本地维护聚合/缓存)尤其要在这里 flush,在 onPartitionsAssigned 里重建。

4.3 少停摆的两把钥匙:CooperativeSticky + static membership

把 rebalance 抖动降到最低,靠两个正交的旋钮配合:

partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
group.instance.id=billing-consumer-3   # 每个实例一个稳定且唯一的 id

升级注意:从 Eager 切到 Cooperative 需要滚动升级(两代协议不能直接混跑),要按官方指引分两步走,别一把梭。

5. 消费端的交付语义:at-least-once + 幂等是主流

把 §2 的结论和核心原理篇 §9 的交付语义表接起来,落到消费端代码:

void process(ConsumerRecord<String, String> r) {
    String eventId = extractEventId(r.value());        // 业务唯一键,不是 offset
    // 幂等落库:重复投递会命中唯一键冲突而被安全忽略
    jdbc.update("INSERT INTO billing(event_id, amount) VALUES(?, ?) ON CONFLICT (event_id) DO NOTHING",
                eventId, extractAmount(r.value()));
}

6. 错误处理:别让一条”毒丸”卡死整条链路

生产里最容易被忽略的是处理失败怎么办。默认写法 for (r : records) process(r) 有个致命隐患:

Poison pill(毒丸消息):某条消息格式错误/反序列化失败/业务异常,如果不处理直接抛,整个 poll 循环就崩了;重启后从同一位点又读到它、又崩——消费永久卡在这一条上,LAG 一路飙升

三种常见对策:

策略做法适用
跳过 + 记录try/catch 住单条异常,记日志/指标后跳过允许丢个别坏消息的场景
死信队列(DLQ)处理失败的消息转发到 xxx-dlq topic,主链路继续大多数生产场景(可事后重放/人工介入)
重试 topic失败消息进 xxx-retry,延迟后重试,多次失败再进 DLQ依赖下游、需要重试的场景
for (ConsumerRecord<String, String> r : records) {
    try {
        process(r);
    } catch (Exception e) {
        dlqProducer.send(new ProducerRecord<>("ad-events-dlq", r.key(), r.value()));
        log.error("moved to DLQ: partition={} offset={}", r.partition(), r.offset(), e);
    }
}
consumer.commitSync();   // 坏消息已进 DLQ,主链路位点可安全推进

6.1 用 pause/resume 做背压

当下游(数据库、外部 API)扛不住时,别硬 poll——用 pause 暂停某些分区、消化完再 resume

if (downstreamOverloaded()) consumer.pause(consumer.assignment());
// ... 等下游恢复 ...
consumer.resume(consumer.paused());

pause 期间 poll 仍会调用(保持心跳、不触发 max.poll.interval 超时),但不返回被暂停分区的数据——这是既做背压、又不掉线的正解。

7. Consumer Lag:消费健康的头号信号

核心原理篇 §6.5 给了 LAG 的定义,这里把它作为运维第一指标展开。

Consumer Lag 拆解图。画一个 Partition 的日志,offset 从左到右递增:offset 0、1 已消费(灰色),offset 2 标为 CURRENT-OFFSET(已提交位点,绿色),offset 3、4、5 是尚未消费的滞后消息(粉色),offset 6 标为 LOG-END-OFFSET(可消费末端,约等于高水位 HW,蓝色)。下方注解:LAG = LOG-END-OFFSET − CURRENT-OFFSET = 6 − 2 = 4 条没消费;它是判断消费者跟不跟得上生产的头号信号;持续增长说明处理能力不足,应加分区/加消费者/并行处理;突然跳高多半是消费卡住、rebalance 或下游阻塞。

# 查每个分区的 LAG(CURRENT-OFFSET / LOG-END-OFFSET / LAG 三列一目了然)
docker exec kafka /opt/kafka/bin/kafka-consumer-groups.sh \
  --bootstrap-server localhost:9092 --describe --group billing

怎么读 LAG:

生产里一般用 Burrow / Kafka Exporter + Prometheus 把 LAG 做成告警。对延迟敏感的广告链路,LAG 告警往往比 CPU/内存告警更早、更准地反映”要出事了”。

8. subscribe vs assign:自动分配还是手动接管

两种拿到分区的方式,语义完全不同:

方式分区分配Consumer Group适用
subscribe(topics)由组协调 + assignor 自动分配,会 rebalance参与 group绝大多数场景
assign(partitions)手动指定分区,不 rebalance、不参与组管理不参与需要精确控制/自己管分配(如按分区并行、CDC 全量读)

assign 配合 seek 能实现精确重放:

consumer.assign(List.of(new TopicPartition("ad-events", 0)));
consumer.seek(new TopicPartition("ad-events", 0), 12345);   // 从指定 offset 开始
// 也可 seekToBeginning / seekToEnd,或按时间戳定位 offsetsForTimes(...)

重放(replay)是 Kafka 的杀手级能力:对账、重算、修数据时,把位点重置到过去重新消费一遍即可(CLI 用 kafka-consumer-groups.sh --reset-offsets)。这也是”消费不删数据”(核心原理篇 §1)带来的直接红利。

9. 消费端并行模型:怎么把吞吐提上去

单消费者单线程处理,吞吐上限就是”单线程处理速度 × 分到的分区数”。想更快,两条路:

9.1 加消费者(横向)—— 简单,但受分区数限制

同一 group 内加消费者,Kafka 自动把分区摊开。上限 = 分区数,多出来的消费者空转(核心原理篇 §3.6)。要更多并行度,得先加分区。

9.2 消费者内线程池(纵向)—— 更灵活,但要自己管 offset

一个消费者拉取、把消息丢进线程池并行处理:

poll 拉一批 → 按 key/分区分发到 worker 线程池并行处理 → 处理完再统一提交

优点是不受分区数限制、CPU 利用率高。代价(都是坑)

经验:能靠”加分区 + 加消费者”解决就别上消费者内线程池——后者把 offset 管理的复杂度全揽到了自己身上。只有当”分区数已经很多、单条处理又重(如调模型/查外部)“时,纵向并行才值得。

10. 落到 AdTech:消费链路的真实取舍

广告系统里,同一份 ad-events 事件流通常被多条消费链路各成一组并行消费(核心原理篇 §13):

消费链路关注点消费端配置取向
计费 / 对账绝不丢、可重放、幂等关自动提交、手动 commitSync、业务唯一键幂等、必要时重放
实时特征 / 样本高吞吐、低延迟CooperativeSticky + static membership、并行处理、盯 LAG
实时预算 / 反作弊低抖动、快感知max.poll.recordspause/resume 背压、DLQ 隔离坏消息

三条铁律贯穿始终:

  1. 发布/扩缩容频繁 → 上 CooperativeStickyAssignor + group.instance.id,把 rebalance 停摆降到最低,保住高峰期端到端延迟 SLA。
  2. 计费链路 → at-least-once + 幂等,绝不用自动提交赌”大概率不丢”。
  3. LAG 是第一告警信号,比资源指标更早暴露”消费跟不上事件洪峰”。

11. 生产反模式与踩坑

这些坑几乎都指向同一根源:没把”poll 循环 + 位点由自己管 + 处理与提交的先后”这套机制当回事。回到 §1 的循环模型,多数坑都能提前预判。

12. 常见误解 ↔ 正解

常见误解正解
消费者处理完消息 Kafka 就自动记住了进度靠你提交 offset;不提交,重启会重读
自动提交 = 处理完才提交自动提交的是”已 poll 出”的位点,与是否处理完无关,可能丢也可能重
心跳正常就不会被踢心跳只管 session.timeout;处理太慢超 max.poll.interval.ms 照样被踢
提交 offset 就是提交当前这条的编号提交的是”下次从哪读”= 已处理最大 offset + 1
Cooperative rebalance 会全组停摆只撤回需迁移的少数分区,其余照常消费,无全组停摆
加消费者就能无限提速上限 = 分区数;再多只能空转
一条坏消息顶多丢它自己不处理会卡死整个分区消费,LAG 持续增长
commitAsynccommitSync 随便用async 快但失败不重试;生产用”async 常态 + sync 兜底”
一个 KafkaConsumer 可以多线程共享非线程安全,多线程共享直接抛异常
LAG 高就是要加机器先看是”整体增长”还是”个别分区倾斜”——后者是热点 key,加机器没用

13. 速查表

到这里,“消息怎么写进去、怎么读出来”就凑齐了。但还有一个反复出现却没讲透的问题:Kafka 到底怎么保证”绝不丢”、以及”精确一次”是怎么做到的? 下一篇Kafka 可靠性与 Exactly-Once 深挖会把副本、ISR、acks × min.insync.replicas、幂等生产者与事务一次讲透。


延伸阅读

本系列内部串读:

一手资料与优质教程:


views
Share this post on:

Previous Post
Kafka 性能内核:磁盘系统凭什么跑出内存级吞吐
Next Post
Kafka Producer 深挖:分区策略与 Sticky Partitioner 是怎么把延迟砍半的