Consumer 与消费组
Kafka 消费端是拉模型:消费者主动 poll() 拉取消息,自己控制节奏。这个设计把复杂度从 broker 移到了客户端——poll 循环怎么写、三个超时参数怎么配,直接决定你的服务会不会陷入"反复 Rebalance、消息重复消费"的泥潭。
1. 为什么是拉(pull)而不是推(push)
- 消费者自己控制速率:处理慢就少拉点,不会被 broker 压垮(天然背压)。
- 方便批量:一次 poll 拉一批,摊薄网络开销。
- 支持回放:拉模型下消费位置由消费者决定,seek 到任意 offset 都行。
代价是消费者要维护一个不停轮询的循环,这就是 poll 循环模型。
2. poll 循环模型
标准骨架(Java):
import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.*;
public class OrderConsumer {
public static void main(String[] args) {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "inventory-service");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); // 手动提交
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(List.of("orders")); // 订阅后并未立刻分配分区
try {
while (true) {
// poll 做三件事: 心跳协调、拉取消息、触发 Rebalance 回调
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
System.out.printf("p=%d offset=%d key=%s value=%s%n",
r.partition(), r.offset(), r.key(), r.value());
// 业务处理...
}
consumer.commitSync(); // 处理完一批再提交 (第 8 章详解)
}
} finally {
consumer.close(); // 主动离组, 立刻触发 Rebalance 而非等超时
}
}
}Python(confluent-kafka)对照:
from confluent_kafka import Consumer
consumer = Consumer({
"bootstrap.servers": "localhost:9092",
"group.id": "inventory-service",
"auto.offset.reset": "earliest",
"enable.auto.commit": False,
})
consumer.subscribe(["orders"])
try:
while True:
msg = consumer.poll(timeout=0.5) # 一次一条 (也可用 consume 批量)
if msg is None:
continue
if msg.error():
print(f"error: {msg.error()}")
continue
print(f"p={msg.partition()} offset={msg.offset()} value={msg.value()}")
consumer.commit(asynchronous=False)
finally:
consumer.close()与 Producer 相反,KafkaConsumer 只能被一个线程使用。要并行消费就启动多个消费者实例(每个自己的线程/进程),靠消费组机制分摊分区。在多个线程里共用一个 consumer 会直接抛 ConcurrentModificationException。
3. 消费组协调机制
每个消费组在 broker 端有一个 Group Coordinator(按 group.id 哈希落到某个 broker)。消费者启动后:
1. FindCoordinator 找到自己组的协调器
2. JoinGroup 入组; 协调器选一个消费者当 "组长" (group leader)
3. 组长根据分配策略计算 分区->消费者 的映射
4. SyncGroup 组长上报方案, 协调器下发给所有成员
5. 各消费者开始 poll 自己名下的分区, 周期性发心跳分配语义回顾(第 2 章):一个分区在组内同一时刻只属于一个消费者;消费者数超过分区数则多余的闲置。
4. 三个超时参数:最容易配错的地方
这是消费端故障排查的核心知识点。三个参数管的是两条独立的线:
线 1: 心跳线程 (后台自动) 线 2: 你的业务处理 (poll 间隔)
├ heartbeat.interval.ms 发心跳周期 └ max.poll.interval.ms
└ session.timeout.ms 判死时限 两次 poll 之间的最大间隔
"进程还活着吗?" "还在正常处理消息吗?"| 参数 | 默认 | 管什么 | 超了会怎样 |
|---|---|---|---|
heartbeat.interval.ms | 3000 | 后台心跳发送间隔 | - |
session.timeout.ms | 45000 | 协调器多久没收到心跳判定死亡 | 踢出组,触发 Rebalance |
max.poll.interval.ms | 300000 | 两次 poll 调用的最大间隔 | 主动离组,触发 Rebalance |
关键区别:
- 进程崩溃 / 网络断开 → 心跳停止 →
session.timeout.ms后被踢。 - 进程活着但处理太慢(一批消息 6 分钟没处理完,超过 max.poll.interval 默认 5 分钟)→ 心跳还在发,但消费者自己判定"我可能死循环了",主动离组。
最经典的故障就是第二种:
症状: 消费组反复 Rebalance, 日志出现
"consumer poll timeout has expired ... max.poll.interval.ms"
原因: 单批消息处理时间 > max.poll.interval.ms
解法 (按优先级):
1. 减小 max.poll.records (默认 500), 一批少拉点
2. 加快单条处理 (异步化、批量写库)
3. 调大 max.poll.interval.ms (治标)配置建议:
session.timeout.ms=45000
heartbeat.interval.ms=15000 # 一般设为 session 的 1/3
max.poll.interval.ms=300000
max.poll.records=500 # 每批最大条数, 按单条处理耗时倒推
fetch.min.bytes=1 # broker 攒够多少字节才返回
fetch.max.wait.ms=500 # 攒不够时最多等多久估算公式:max.poll.records × 单条最大处理时间 < max.poll.interval.ms,留 50% 余量。
5. subscribe 与 assign
| 方式 | 说明 | Rebalance |
|---|---|---|
subscribe(topics) | 加入消费组,由协调器分配分区 | 参与 |
assign(partitions) | 手动指定分区,不走消费组协调 | 不参与 |
assign 适合:数据修复工具、精确控制分区的 ETL、单实例消费全部分区的小场景。它没有故障转移——实例挂了没人接管分区。
// 手动指定消费 orders 的 0 号分区
TopicPartition p0 = new TopicPartition("orders", 0);
consumer.assign(List.of(p0));
consumer.seekToBeginning(List.of(p0));6. 消费者数量怎么定
- 起点:消费者数 = 分区数 ÷ 2 到分区数之间,观察 LAG 再调。
- 上限:组内消费者数超过分区数没有意义。
- 单消费者吞吐不够时,先确认瓶颈在业务处理还是拉取(通常是业务写库慢),处理逻辑内部再做异步/批量优化,比无脑加实例有效。
收到 SIGTERM 时调用 consumer.wakeup() 打断阻塞中的 poll(会抛 WakeupException),在 finally 里 close()。close 会发 LeaveGroup 请求,让 Rebalance 立即发生,而不是干等 session.timeout —— 滚动发布时这能把消费中断从 45 秒降到秒级。
小结
- Kafka 消费是拉模型,poll 循环同时承担拉消息、心跳协调、Rebalance 参与三个职责
- Consumer 非线程安全,一个线程一个实例,并行靠消费组分摊分区
- 三个超时:session.timeout 判"进程死没死",max.poll.interval 判"处理卡没卡",heartbeat 是心跳频率
max.poll.records × 单条耗时必须小于 max.poll.interval.ms- subscribe 走组协调有故障转移,assign 手动控制无接管
- 下一章:offset 提交的正确姿势 →
- 启动 3 个相同 group.id 的消费者消费 6 分区 topic,观察各自分到哪些分区;关掉一个再观察重新分配。
- 在处理逻辑里
Thread.sleep(400000)模拟慢处理(超过 max.poll.interval.ms),观察消费者被移出组的日志与 Rebalance。 - 把 max.poll.records 改成 10,对比默认 500 时单次 poll 的处理时长差异。