Kafka 消费者组原理与消费语义详解
Mei Lin | 2026-08-27T20:55:33 | DevOps, Database
深入分析 Kafka 消费者组的分区分配策略、再平衡机制和三种消费语义(At Most Once/At Least Once/Exactly Once)。
# Kafka 消费者组原理与消费语义 ## 一、消费者组基本概念 ``` Topic: order-events (3 个分区) Partition 0 ──> Consumer A (Group: order-processor) Partition 1 ──> Consumer B (Group: order-processor) Partition 2 ──> Consumer A (Group: order-processor) ``` 同一消费者组内,每个分区只会分配给一个消费者;不同组可以独立消费同一个 Topic。 ## 二、分区分配策略 ```java Properties props = new Properties(); props.put("bootstrap.servers", "kafka:9092"); props.put("group.id", "order-processor"); // 分区分配策略 props.put("partition.assignment.strategy", "org.apache.kafka.clients.consumer.CooperativeStickyAssignor"); KafkaConsumer consumer = new KafkaConsumer(props); consumer.subscribe(Arrays.asList("order-events")); ``` | 策略 | 说明 | |------|------| | RangeAssignor | 按分区范围分配(默认) | | RoundRobinAssignor | 轮询分配 | | StickyAssignor | 粘性分配(减少再平衡影响) | | CooperativeStickyAssignor | 协作式粘性(推荐) | ## 三、再平衡(Rebalance) 触发条件:消费者加入/离开组、订阅 Topic 分区数变化。 ```java consumer.subscribe(Arrays.asList("order-events"), new ConsumerRebalanceListener() { public void onPartitionsRevoked(Collection partitions) { // 提交当前偏移量,释放资源 consumer.commitSync(); System.out.println("释放分区: " + partitions); } public void onPartitionsAssigned(Collection partitions) { System.out.println("分配到分区: " + partitions); } } ); ``` ## 四、三种消费语义 ### At Most Once(至多一次) ```java // 先提交偏移量,再处理消息 // 如果处理失败,消息丢失 while (true) { ConsumerRecords records = consumer.poll(Duration.ofMillis(100)); consumer.commitSync(); // 先提交 for (ConsumerRecord record : records) { processMessage(record); // 后处理(可能失败) } } ``` ### At Least Once(至少一次,推荐) ```java // 先处理消息,再提交偏移量 // 处理失败会重复消费,业务需幂等 while (true) { ConsumerRecords records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord record : records) { processMessage(record); // 先处理 } consumer.commitSync(); // 后提交 } ``` ### Exactly Once(精确一次) ```java // 使用事务 + 消费者偏移量 props.put("isolation.level", "read_committed"); // 生产者端 producer.initTransactions(); try { producer.beginTransaction(); producer.send(new ProducerRecord("output-topic", key, value)); // 将消费偏移量作为事务的一部分提交 producer.sendOffsetsToTransaction(offsets, groupMetadata); producer.commitTransaction(); } catch (Exception e) { producer.abortTransaction(); } ``` ## 五、消费者配置优化 ```java // 最大拉取间隔(超时触发再平衡) props.put("max.poll.interval.ms", "300000"); // 5 分钟 // 单次拉取最大记录数 props.put("max.poll.records", "500"); // 心跳间隔 props.put("heartbeat.interval.ms", "3000"); // 会话超时 props.put("session.timeout.ms", "45000"); // 自动提交(生产环境建议关闭) props.put("enable.auto.commit", "false"); ``` ## 六、常见问题 1. **消费延迟**:增加消费者数量(不超过分区数) 2. **消息堆积**:检查消费者处理速度,增加分区和消费者 3. **重复消费**:业务做幂等处理(数据库唯一键、Redis 去重) 4. **顺序消费**:同一 Key 发到同一分区,单分区单消费者