前两篇立好了地基:Kafka 是一份分布式可重放日志,Producer 用分区策略把消息高效攒批塞进去。写入侧讲透了,这一篇补齐对称的另一半——消息怎么被可靠、高效、不抖动地读出来。
Consumer 看似只是个 while(true) { poll(); process(); } 循环,但生产里几乎所有”消息丢了""消息重复了""消费突然卡住""一发布就整组停摆”的事故,根子都在这个循环的细节里。这篇把它拆开讲透。
TL;DR
- 消费的核心是一个单线程循环:
poll → 处理 → 提交 → 再 poll。一个KafkaConsumer实例不是线程安全的,拉取只能在一个线程里做;心跳自 0.10.1 起交给一个后台独立线程,但存活判定分两条线(见下)。 - 两条存活判定各管一段:①心跳线程按
heartbeat.interval.ms发心跳,session.timeout.ms内没心跳 → 判定掉线;②应用线程两次poll间隔超max.poll.interval.ms→ 判定”假死”、主动离组。单批处理太慢 = 无限 rebalance 的头号原因。 - offset 提交决定交付语义:先处理后提交 = at-least-once(不丢可能重);先提交后处理 = at-most-once(可能丢)。自动提交是”定时后台提交”,时机不可控,要”不丢”就关掉它、改手动。
- 提交有四种姿势:自动提交、
commitSync(可靠、阻塞)、commitAsync(快、可能提交失败不重试)、commitSync(offsets)(精确提交指定位点)。生产常用”async 常态 + sync 兜底 + rebalance 前强制提交”。 - Rebalance = 分区分配 + 生命周期:分区靠 assignor(Range / RoundRobin / Sticky) 分给消费者;触发时依次回调
onPartitionsRevoked(提交位点、flush 状态)→ JoinGroup/SyncGroup →onPartitionsAssigned(seek、重建缓存)。用CooperativeStickyAssignor+ static membership 把抖动降到最小。 - Consumer Lag = LOG-END-OFFSET − CURRENT-OFFSET,是消费健康的头号信号;持续增长要加并行度,突然跳高多半是卡住/rebalance/下游阻塞。
- 错误处理要有预案:poison pill(毒丸消息)会把消费卡死,用 try/catch + 死信队列(DLQ)/重试 topic 隔离;用
pause/resume做背压。 - 想提消费吞吐:先加分区再加消费者(上限 = 分区数),或在消费者内部用线程池并行处理(代价是要自己管好 offset 提交与保序)。
Table of contents
Open Table of contents
- 1. 消费者的核心循环:poll 到底做了什么
- 2. offset 提交:四种姿势与不丢不重的时机
- 3. poll 循环的三个”生死参数”
- 4. Rebalance 深挖:分区分配、生命周期与 Listener
- 5. 消费端的交付语义:at-least-once + 幂等是主流
- 6. 错误处理:别让一条”毒丸”卡死整条链路
- 7. Consumer Lag:消费健康的头号信号
- 8. subscribe vs assign:自动分配还是手动接管
- 9. 消费端并行模型:怎么把吞吐提上去
- 10. 落到 AdTech:消费链路的真实取舍
- 11. 生产反模式与踩坑
- 12. 常见误解 ↔ 正解
- 13. 速查表
- 延伸阅读
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 上——这解释了后面很多”反直觉”的行为。
1.1 一个消费者实例只能一个线程用
KafkaConsumer 不是线程安全的。多线程共享一个实例会直接抛 ConcurrentModificationException。这是刻意的设计:把”拉取 + 位点管理”约束在单线程里,语义最简单。想并行,见 §9。
1.2 position vs committed offset:风险窗口在哪
回顾核心原理篇 §3.3:要严格区分两个位置。
- position(当前位置):消费者下一条要读的 offset,每次
poll后向前推进(在内存里)。 - committed offset(已提交位点):已经确认”处理完了”并写回集群的位点,存于内部 topic
__consumer_offsets。
两者之间的差,就是崩溃时消息重复或丢失的风险窗口——你在哪一步提交、提交了什么,直接决定语义。这就是下一节的主题。
2. offset 提交:四种姿势与不丢不重的时机
“消息会不会丢/会不会重”不取决于 Kafka,取决于你在处理的哪一步提交 offset。
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.ms | 3s | 心跳线程 | 发心跳的频率(应 < session 的 1/3) |
session.timeout.ms | 45s(新版) | Group Coordinator | 这段时间没收到心跳 → 判消费者掉线,触发 rebalance |
max.poll.interval.ms | 5min | Group Coordinator | 两次 poll 间隔超它 → 判”假死”,主动离组 |
关键在于 0.10.1 之后心跳被拆到独立后台线程:即使你的处理逻辑正卡着,心跳照样在发,session.timeout.ms 这条线不会误杀你。真正卡死你的是另一条线——
max.poll.interval.ms:如果单批消息处理太久(比如一批 500 条,每条要写一次慢下游),迟迟不回来调下一次poll,Coordinator 就认为你”假死”、把你踢出组 → rebalance → 分区给别人 → 你处理完回来发现位点没了、又要重新加入 → 恶性循环,整组越来越慢。
对策三选一(常组合):
- 减小
max.poll.records(默认 500):一次少拉点,缩短单批处理时间。 - 调大
max.poll.interval.ms:给慢处理留足预算。 - 把慢处理异步化(§9),让主循环快速回到
poll。
4. Rebalance 深挖:分区分配、生命周期与 Listener
核心原理篇 §5 讲了 rebalance 的 Eager vs Cooperative 两种协议。这里从消费者视角补齐三件事:分区按什么策略分给消费者(assignor)、rebalance 在你代码里触发哪些回调、以及怎么把停摆降到最小。
4.1 分区分配策略(assignor):谁拿到哪些分区
消息进哪个分区由 Producer 的分区策略决定(见 Producer 篇);写进去之后,这些分区再分给消费者组里的哪个消费者——这由消费侧的 partition assignment strategy(partition.assignment.strategy)决定。它不改吞吐/延迟的批处理逻辑,却直接决定消费端的负载均衡与 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 时,你的代码里会依次触发哪些回调、每个回调该做什么,才是不丢不重的关键:
订阅时挂一个 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 抖动降到最低,靠两个正交的旋钮配合:
CooperativeStickyAssignor(§4.1 表里那一行):rebalance 时只撤回需要迁移的少数分区、其余照常消费,无全组停摆——把”换主”的代价从”整组停摆”缩到”个别分区迁移”。现代默认推荐。- static membership(
group.instance.id):给消费者一个固定身份。滚动重启 / 发布时,只要在session.timeout.ms内回来,就不触发 rebalance——直接接回原来的分区。对广告这类频繁发布的服务,这是保住高峰期 SLA 的关键一招。
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
group.instance.id=billing-consumer-3 # 每个实例一个稳定且唯一的 id
升级注意:从 Eager 切到 Cooperative 需要滚动升级(两代协议不能直接混跑),要按官方指引分两步走,别一把梭。
5. 消费端的交付语义:at-least-once + 幂等是主流
把 §2 的结论和核心原理篇 §9 的交付语义表接起来,落到消费端代码:
- 默认走 at-least-once:先处理、后提交,绝不丢、可能重。
- 用幂等消化重复:崩溃重投是常态,必须在消费端做幂等——数据库唯一键、
INSERT ... ON CONFLICT、RedisSETNX、去重表。用业务主键(如event_id),别用offset当幂等键(重放/扩分区后会变)。 - exactly-once 是另一条路:
isolation.level=read_committed+ 事务生产,只保证 Kafka→Kafka 闭环不重不丢;一旦有外部副作用,端到端仍需幂等。EOS 与事务的完整机制,是下一篇可靠性与 Exactly-Once 深挖的主题,这里先埋个扣。
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 的定义,这里把它作为运维第一指标展开。
# 查每个分区的 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:
- 稳定的小 LAG:健康,消费基本跟得上生产。
- 持续单调增长:处理能力不足 → 先加分区、再加消费者(上限 = 分区数,见核心原理篇 §7),或上并行处理(§9)。
- 突然阶跃式跳高:多半是消费卡住(poison pill)、rebalance 停摆、或下游阻塞——结合监控看是哪个分区、哪个实例。
- 只有个别分区 LAG 高:典型的热点 key/数据倾斜(见 Producer 篇 §2.2)。
生产里一般用 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 利用率高。代价(都是坑):
- 保序被打破:线程池并行处理会打乱顺序。要保序就按 key/分区路由到固定 worker(同 key 进同一线程),代价见 Producer 篇 §7.3。
- offset 提交变复杂:不能简单提交”最新批”,要等连续处理完的前缀才能提交(否则中间某条没处理完就提交会丢)。通常维护一个”已完成 offset”的水位。
- 失败重试 / 背压:worker 失败要有 DLQ;队列积压要用
pause背压。
经验:能靠”加分区 + 加消费者”解决就别上消费者内线程池——后者把 offset 管理的复杂度全揽到了自己身上。只有当”分区数已经很多、单条处理又重(如调模型/查外部)“时,纵向并行才值得。
10. 落到 AdTech:消费链路的真实取舍
广告系统里,同一份 ad-events 事件流通常被多条消费链路各成一组并行消费(核心原理篇 §13):
| 消费链路 | 关注点 | 消费端配置取向 |
|---|---|---|
| 计费 / 对账 | 绝不丢、可重放、幂等 | 关自动提交、手动 commitSync、业务唯一键幂等、必要时重放 |
| 实时特征 / 样本 | 高吞吐、低延迟 | CooperativeSticky + static membership、并行处理、盯 LAG |
| 实时预算 / 反作弊 | 低抖动、快感知 | 小 max.poll.records、pause/resume 背压、DLQ 隔离坏消息 |
三条铁律贯穿始终:
- 发布/扩缩容频繁 → 上
CooperativeStickyAssignor+group.instance.id,把 rebalance 停摆降到最低,保住高峰期端到端延迟 SLA。 - 计费链路 → at-least-once + 幂等,绝不用自动提交赌”大概率不丢”。
- LAG 是第一告警信号,比资源指标更早暴露”消费跟不上事件洪峰”。
11. 生产反模式与踩坑
- 忘关自动提交就做副作用处理:以为”处理完自然提交”,实际位点在处理前就被后台提交了 → 崩溃丢消息。对策:
enable.auto.commit=false+ 手动提交。 - 单批处理太慢触发无限 rebalance:一次
poll处理超过max.poll.interval.ms→ 被踢 → rebalance → 更慢。对策:减小max.poll.records/ 调大间隔 / 异步处理(§3)。 - 提交 offset 少加 1:提交
offset而非offset+1,重启重复消费最后一条。 - 不处理 poison pill:一条坏消息卡死整个分区消费,LAG 爆炸。对策:try/catch + DLQ(§6)。
- 多线程共享一个
KafkaConsumer:直接ConcurrentModificationException。对策:一线程一实例,或用线程池 worker 模式(§9.2)。 - 在
onPartitionsRevoked里忘了提交:rebalance 后接管者从旧位点重放一大段。对策:revoke 回调里commitSync。 - 消费者数 > 分区数还以为能提速:多出的消费者空转。对策:先加分区。
- 把 offset 当幂等键:重放/扩分区后 offset 变了,幂等失效。对策:用业务唯一键。
这些坑几乎都指向同一根源:没把”poll 循环 + 位点由自己管 + 处理与提交的先后”这套机制当回事。回到 §1 的循环模型,多数坑都能提前预判。
12. 常见误解 ↔ 正解
| 常见误解 | 正解 |
|---|---|
| 消费者处理完消息 Kafka 就自动记住了 | 进度靠你提交 offset;不提交,重启会重读 |
| 自动提交 = 处理完才提交 | 自动提交的是”已 poll 出”的位点,与是否处理完无关,可能丢也可能重 |
| 心跳正常就不会被踢 | 心跳只管 session.timeout;处理太慢超 max.poll.interval.ms 照样被踢 |
| 提交 offset 就是提交当前这条的编号 | 提交的是”下次从哪读”= 已处理最大 offset + 1 |
| Cooperative rebalance 会全组停摆 | 只撤回需迁移的少数分区,其余照常消费,无全组停摆 |
| 加消费者就能无限提速 | 上限 = 分区数;再多只能空转 |
| 一条坏消息顶多丢它自己 | 不处理会卡死整个分区消费,LAG 持续增长 |
commitAsync 和 commitSync 随便用 | async 快但失败不重试;生产用”async 常态 + sync 兜底” |
一个 KafkaConsumer 可以多线程共享 | 它非线程安全,多线程共享直接抛异常 |
| LAG 高就是要加机器 | 先看是”整体增长”还是”个别分区倾斜”——后者是热点 key,加机器没用 |
13. 速查表
- 核心循环:
poll → 处理 → 提交;一个消费者实例只能一个线程用;心跳在后台独立线程。 - 两条存活线:
session.timeout.ms(心跳)+max.poll.interval.ms(poll 间隔);后者是无限 rebalance 元凶。 - 交付语义:先处理后提交 = at-least-once(主流,配幂等);先提交后处理 = at-most-once。
- 四种提交:自动(时机不可控)/
commitSync(可靠阻塞)/commitAsync(快不重试)/ 指定位点(细粒度)。生产 = async 常态 + sync 兜底 + revoke 前强制提交。 - 提交值:已处理最大 offset + 1。
- rebalance:
onPartitionsRevoked(提交、flush)→ Join/Sync →onPartitionsAssigned(seek、重建);用CooperativeSticky+group.instance.id降抖动。 - 错误处理:try/catch + DLQ / 重试 topic;
pause/resume做背压。 - LAG = LOG-END-OFFSET − CURRENT-OFFSET,运维第一信号。
- 并行:先加分区 + 加消费者(上限 = 分区数);再不够上消费者内线程池(自己管保序与提交)。
- 重放:
assign + seek或--reset-offsets;幂等键用业务主键,别用 offset。
到这里,“消息怎么写进去、怎么读出来”就凑齐了。但还有一个反复出现却没讲透的问题:Kafka 到底怎么保证”绝不丢”、以及”精确一次”是怎么做到的? 下一篇Kafka 可靠性与 Exactly-Once 深挖会把副本、ISR、acks × min.insync.replicas、幂等生产者与事务一次讲透。
延伸阅读
本系列内部串读:
- Kafka 核心原理精讲:从一条日志到分布式流平台:Partition / Offset / Consumer Group / Rebalance / 交付语义的地基。
- Kafka Producer 深挖:分区策略与 Sticky Partitioner:写入侧怎么攒批、怎么选分区,以及顺序性与热点 key 隔离。
- Kafka 可靠性与 Exactly-Once 深挖:
acks/min.insync.replicas/ 幂等生产者 / 事务与 EOS 的完整机制。 - Kafka 性能内核:磁盘系统凭什么跑出内存级吞吐:追尾消费、冷读与 page cache 命中率的关系。
- 为什么 AdTech 偏爱 Kafka:广告事件中枢的五大典型场景与架构设计:CooperativeSticky、DLQ、LAG 告警在广告计费/风控消费链路里的落地。
一手资料与优质教程:
- Apache Kafka. Consumer Configs:
enable.auto.commit/max.poll.records/max.poll.interval.ms/session.timeout.ms/partition.assignment.strategy/group.instance.id的权威说明。 - Apache Kafka. KafkaConsumer Javadoc:poll 循环、位点提交、
ConsumerRebalanceListener、pause/resume、seek的官方 API 语义。 - Apache Kafka. KIP-429: Incremental Cooperative Rebalancing 与 KIP-345: Static Membership:合作式重平衡与静态成员的一手提案。
- Confluent. Kafka Consumer:消费者位点管理、rebalance、并行模型的系统讲解。