Learn
Kafka/07-consumer-basics

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()
⚠️KafkaConsumer 不是线程安全的

与 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.ms3000后台心跳发送间隔-
session.timeout.ms45000协调器多久没收到心跳判定死亡踢出组,触发 Rebalance
max.poll.interval.ms300000两次 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 提交的正确姿势 →
🎯练习
  1. 启动 3 个相同 group.id 的消费者消费 6 分区 topic,观察各自分到哪些分区;关掉一个再观察重新分配。
  2. 在处理逻辑里 Thread.sleep(400000) 模拟慢处理(超过 max.poll.interval.ms),观察消费者被移出组的日志与 Rebalance。
  3. 把 max.poll.records 改成 10,对比默认 500 时单次 poll 的处理时长差异。