Skip to content

生产者与消费者原理:批次、拉取与 Rebalance

本文是消息队列系统学习系列的 L2 核心篇。前置:27. Kafka 核心概念:Topic、Partition 与 Consumer Group。 学完可以配合面试题食用:03-kafka-message-reliability-ack-isr08-kafka-message-ordering

一条消息从发送到落盘:生产者内部链路

Kafka 生产者不是一条消息一条消息地发出去的。客户端把 send() 调用攒成批次,再批量发到 broker,靠这个攒批动作把吞吐量拉上去。链路如下:

Producer -> Interceptor -> Serializer -> Partitioner -> RecordAccumulator -> Sender 线程 -> Broker

Interceptor 允许在发送前后插入自定义逻辑(埋点、修改消息体),顺序执行,抛异常不阻断。

Serializer 把 key/value 转成字节数组。除非有自定义序列化需求(比如用 Protobuf),否则默认的 StringSerializer 就够了。

Partitioner 决定消息去哪个分区。默认策略是:有 key 就走 key 的 murmur2 哈希取模;没 key 就走粘性分区(Sticky Partitioner),攒一批全塞一个分区,攒满再换下一个,比之前轮询策略的发送效率更高。

RecordAccumulator 是核心缓冲区。每个分区一个双端队列(Deque),每个队列元素是一个 ProducerBatch。攒满 batch.size(默认 16KB)或超过 linger.ms(默认 0,即立即发)就触发 flush。这俩参数是生产吞吐的核心旋钮。

Sender 线程 后台轮询,从 Accumulator 取出已就绪的批次,按 broker 节点打包成请求,发完后处理响应(acks 和重试逻辑)。

java
// 生产者核心配置示例
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// 吞吐调优参数
props.put("batch.size", 32768);          // 32KB,加大可降请求次数
props.put("linger.ms", 5);               // 多等 5ms 攒批次
props.put("compression.type", "snappy"); // 开启压缩
props.put("acks", "all");                // 等待所有副本确认
props.put("retries", 3);                 // 重试次数
props.put("enable.idempotence", true);   // 幂等生产者

KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("topic-a", "order-001", "payload"));

在 acks=all 下,生产者会等到所有 ISR 副本确认才返回成功。max.in.flight.requests.per.connection 默认 5——但启用了幂等生产者后,Kafka 会自动降为 1 或 5 并保证顺序不乱。

分区策略:顺序性的边界在哪

分区器决定了同一条消息落在哪个分区。常见策略:

  • key 哈希:相同 key 走到同分区,同分区内消息有序。这是 Kafka 保证消息顺序的唯一方式。
  • 粘性分区:没 key 时,攒一批往一个分区怼,攒满后换下一个。相比轮询,粘性分区攒出的批次更大、压缩率更高。
  • 自定义分区器:实现 Partitioner 接口,按业务规则路由(比如按订单号尾号分)。

顺序性的边界:同 key 同分区,写顺序即消费顺序。但一旦重试,如果 max.in.flight.requests.per.connection > 1,前一批失败后一批成功,顺序就破了。幂等生产者解决了这个问题——它给每个批次标序列号,broker 对同序列号去重,同时保证批次按序写入(即使并行发送也等失败批次重试完再提交后续批次)。

消费者链路:poll 模型与位移提交

消费者很反直觉:它不是 broker 推数据过来的,而是客户端不断主动 poll。

mermaid
sequenceDiagram
    participant C as 消费者
    participant K as Kafka Broker
    participant CG as __consumer_offsets
    loop 每次 poll 循环
        C->>K: poll(Duration)
        K->>C: 返回 ConsumerRecords
        C->>C: 处理业务逻辑
        C->>CG: 提交位移(同步/异步)
    end

poll 循环:消费者线程唯一要做的事就是反复调 poll()。每次 poll 返回一批消息,处理完后提交位移。poll() 内部会做三件事:拉取数据、发送心跳、触发 rebalance 回调。

位移提交有三个选项:

  • 自动提交(enable.auto.commit=true,默认 5 秒一次):简单但有重复消费窗口——处理完但还没到提交点就挂了,恢复后从上次位移重拉。
  • 手动同步提交consumer.commitSync(),阻塞直到提交成功,吞吐低但可靠。
  • 手动异步提交consumer.commitAsync(),不阻塞,配合回调做重试(注意:异步重试要等上一轮成功,否则后提交的可能覆盖先提交的)。

参数坑

  • max.poll.records:单次 poll 最多返回多少条。设太大处理时间超了导致 rebalance。
  • max.poll.interval.ms:两次 poll 之间最大间隔,超了 coordinator 认为消费者挂了,踢出组。
  • session.timeout.ms:心跳超时(默认 45s 新版本),超了同样触发 rebalance。
  • heartbeat.interval.ms:心跳频率(默认 3s)。设太长会延迟发现消费者宕机。

这三个参数调不好的典型后果:消费者处理慢一点就被踢出组,rebalance 完又分配回来,反复踢入踢出,消费根本往前走不了。

Rebalance 全流程:两阶段协议

Rebalance 是 Kafka 消费组最让人头疼的机制。每次组内成员增减或分区数变化都会触发,期间整个组会短暂停止消费。

协调器(Coordinator)是每个消费者组选出的一个 broker 节点,负责管理组内成员和位移。

消费者组工作流程:
1. 每个消费者启动时找协调器注册
2. 协调器维护组成员列表 + 各自分区分配
3. 触发 rebalance 时,协调器选定一个 Consumer Group Leader
4. Leader 执行分区分配策略(Range / RoundRobin / Sticky / CooperativeSticky)
5. 分配结果广播给所有成员

两阶段协议:第一阶段 JoinGroup——所有成员向协调器发送 JoinGroup 请求,协调器选 Leader 并把成员列表发给 Leader。第二阶段 SyncGroup——Leader 算出分配方案,各成员拿到自己的分区。

Generation:每次 rebalance 生成一个新 generation(递增的整数)。消费者提交位移时带上 generation,老的 generation 提交会被拒绝——防止脑裂的旧消费者污染位移。

java
// 注册 Rebalance 监听器,在 rebalance 前后做位移保存
consumer.subscribe(Collections.singletonList("topic-a"), new ConsumerRebalanceListener() {
    @Override
    public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
        // 在 rebalance 前提交当前位移,避免重复消费
        consumer.commitSync();
        // 保存当前正在处理但未完成的 offset 到外部存储
        saveOffsetsToExternal(partitions);
    }

    @Override
    public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
        // 从外部存储恢复上一个消费位置
        for (TopicPartition tp : partitions) {
            long offset = loadOffsetFromExternal(tp);
            consumer.seek(tp, offset);
        }
    }
});

以上代码规避了 rebalance 后重复消费的问题:分区被收回前提交位移,新分区分配后 seek 到上次保存的位置。

常见误区与小结

  • max.in.flight.requests 设 5 就一定乱序? 不,只有出现重试才乱序。幂等生产者启用了就可以保持 5 的并发同时保证顺序——它靠序列号去重和在 broker 端按序落盘解决的。
  • poll 间隔设大一点就不会 rebalance 了? 不一定。max.poll.interval.ms 只是踢出条件之一,session.timeout.ms 是心跳超时,两者独立生效。哪个触发了都会 rebalance。
  • 消费者数 > 分区数会怎样? 多出的消费者闲置,浪费资源。因为一个分区只能被组内一个消费者消费。
  • 手动提交总能保证 exactly-once? 不能。先处理后提交,如果处理完提交前挂了,重启后重复消费。先提交后处理,如果提交后处理前挂了,消息丢失。唯一解是消费结果写进支持事务的外部存储,用同一个事务提交位移和处理结果(Kafka 的 exactly-once 也走这条路)。
  • rebalance 期间消息会丢吗? 不会,但会有短暂的不可用。正确实现 ConsumerRebalanceListener 可以最大程度减少重复消费窗口。

小结:生产者侧的核心是批次攒够了再发,消费者侧的核心是 poll 循环 + 位移提交 + rebalance 三者的配合。生产者调好 batch.size 和 linger.ms 就能拿足吞吐,消费者调好 max.poll.interval.ms 和 session.timeout.ms 才能稳定消费。下一篇 29. Kafka 的可靠性与高性能架构 会深入 broker 端是怎么保证不丢消息的,以及顺序写和零拷贝到底快在哪。

参考

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。
粤ICP备2026104257号-1