Consumer 概述
定义与作用
Kafka Consumer 是从 Kafka Topic 拉取(Pull)消息的客户端程序。与 Producer 的 Push 模型不同,Consumer 采用 Pull 模型——消费者根据自己的处理能力主动拉取消息,避免了 Broker 推送速度超过消费者处理能力导致的内存溢出。
Consumer 在 Kafka 架构中的位置:数据出口。它与 Producer 之间通过 Broker 完全解耦——Producer 不知道 Consumer 的存在,Consumer 也不知道 Producer 的存在。
核心原理
Push vs Pull 模型对比
Consumer 的核心组件
| 组件 | 职责 |
|---|---|
poll() | 应用入口,拉取一批消息并返回 |
Fetcher | 独立的拉取线程,从 Leader 分区拉取数据 |
SubscriptionState | 维护订阅的 Topic/Partition 和对应 Offset |
ConsumerCoordinator | 管理消费组、Rebalance、Offset 提交 |
Deserializer | 将字节数组还原为 Java 对象 |
Consumer 消费语义
| 语义 | 含义 | 实现方式 |
|---|---|---|
| At-most-once | 可能丢失,不重复 | 先提交 Offset 再处理 |
| At-least-once(默认) | 可能重复,不丢失 | 先处理再提交 Offset |
| Exactly-once | 不丢不重 | 幂等 Producer + 事务 + read_committed |
完整示例
示例一:最小可用的 Consumer
场景:消费 orders Topic 的全部消息。
操作前环境:Kafka 已启动,orders Topic 中有 100 条消息。
import org.apache.kafka.clients.consumer.*;
import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
public class SimpleConsumer {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "order-processor");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("enable.auto.commit", "false");
props.put("auto.offset.reset", "earliest");
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
consumer.subscribe(Collections.singletonList("orders"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("offset=%d, key=%s, value=%s%n",
record.offset(), record.key(), record.value());
}
consumer.commitSync();
}
}
}
}
执行结果(截取前 3 条):
offset=0, key=order-1, value={"amount": 99}
offset=1, key=order-2, value={"amount": 150}
offset=2, key=order-3, value={"amount": 200}
操作前后对比:
| 阶段 | Consumer Offset | orders LEO | Lag |
|---|---|---|---|
| 消费前 | 0 | 100 | 100 |
| 消费 100 条后 | 100 | 100 | 0 |
示例二:指定分区消费(手动分配)
场景:不希望加入消费组(不触发 Rebalance),直接指定 Partition 0 消费。
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
// 注意:手动分配时不需要 group.id
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
TopicPartition partition0 = new TopicPartition("orders", 0);
consumer.assign(Collections.singletonList(partition0));
consumer.seekToBeginning(Collections.singletonList(partition0));
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("partition=%d, offset=%d, value=%s%n",
record.partition(), record.offset(), record.value());
}
}
subscribe vs assign 对比:
| 特性 | subscribe() | assign() |
|---|---|---|
| 消费组 | 需要 group.id | 不需要 |
| Rebalance | 自动触发 | 无 Rebalance |
| 分区分配 | Broker 自动分配 | 手动指定 |
| Offset 提交 | 支持 | 需手动管理 |
| 适用场景 | 常规消费 | 精确控制、调试 |
易错场景
易错 1:poll() 间隔过长导致离开消费组
现象:Consumer 处理时间过长,日志出现 Member xxx has left the group。
原因:poll() 是 Consumer 的"心跳"信号。如果两次 poll() 间隔超过 max.poll.interval.ms(默认 5 分钟),Consumer 被认为已死亡并被踢出消费组。
解决:
- 减少单次
poll()的消息数(max.poll.records,默认 500) - 将处理逻辑移到独立线程池,
poll()线程快速返回继续拉取 - 增加
max.poll.interval.ms
易错 2:忘记调用 poll() 导致命令工具看不到 Consumer Group
现象:kafka-consumer-groups.sh --describe 看不到给定 group.id 的消费组。
原因:Kafka Consumer 在第一次 poll() 时才会向 Broker 注册消费者组。如果创建 Consumer 后只是 subscribe() 但不调用 poll(),消费组不会出现在 Broker 中。
面试高频考点
Q:Kafka 为什么选择 Pull 模型而不是 Push 模型?
A:
- 消费者自主控制速率 — Consumer 根据自己的处理能力拉取,不会因 Broker 推送过快而 OOM
- 天然支持批量 — Producer 端攒批 + Consumer 端批量 poll,两端优化不相互影响
- 简化 Broker — Broker 不需要追踪每个 Consumer 的消费状态和推送队列
- 重放能力 — Consumer 可以 seek 到任意 Offset 重新消费,Push 模型下难以实现
代价:Pull 模型在低延迟场景有劣势——Consumer 需要轮询,可能引入延迟。Kafka 通过长轮询(fetch.min.bytes + fetch.max.wait.ms)缓解此问题。