Consumer 多线程
定义与作用
Kafka Consumer 在单个线程中不是线程安全的——poll()、commitSync() 等方法必须在同一个线程中调用。要提升消费吞吐,必须通过在多个线程中运行多个 Consumer 实例,或将消息处理从 poll() 线程中分离出来。
本节介绍 Kafka Consumer 的三种多线程模式及其适用场景。
核心原理
三种多线程模式
| 模式 | 优势 | 劣势 | 适用场景 |
|---|---|---|---|
| 多 Consumer 实例 | 最简单、最可靠 | 实例数受限于分区数 | 首选方案 |
| 单 Consumer + 多处理线程 | 不受分区数限制 | Offset 管理复杂 | 单分区吞吐不足 |
| 分区独立处理 | 分区间完全隔离 | 极端复杂、易出错 | 特殊场景 |
完整示例
示例一:多 Consumer 实例(推荐方案)
场景:Topic 有 12 个分区,启动 6 个 Consumer 线程,每个处理 2 个分区。
public class MultiConsumerApp {
private static final int NUM_CONSUMERS = 6;
public static void main(String[] args) {
ExecutorService executor = Executors.newFixedThreadPool(NUM_CONSUMERS);
for (int i = 0; i < NUM_CONSUMERS; i++) {
executor.submit(new ConsumerTask("consumer-" + i));
}
}
static class ConsumerTask implements Runnable {
private final String name;
ConsumerTask(String name) { this.name = name; }
@Override
public void run() {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "high-throughput-group");
props.put("client.id", name);
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");
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
consumer.subscribe(Collections.singletonList("events"));
while (!Thread.currentThread().isInterrupted()) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> record : records) {
processEvent(record);
}
consumer.commitSync();
}
}
}
private void processEvent(ConsumerRecord<String, String> record) {
// 业务处理
}
}
}
操作前后对比:
| 线程数 | 分区分配 | 吞吐量 |
|---|---|---|
| 1(基线) | 1 Consumer → 12 分区 | ~50K events/s |
| 6 | 6 Consumer → 各 2 分区 | ~300K events/s |
| 12 | 12 Consumer → 各 1 分区 | ~500K events/s |
| 13 | 12 Consumer 有分区 + 1 空闲 | ~500K events/s(无提升) |
示例二:单 Consumer + 处理线程池(Offset 管理)
场景:Topic 只有 1 个分区,但消息处理(调用外部 API)非常慢,需要多线程处理。
public class SingleConsumerMultiWorker {
private static final Map<TopicPartition, OffsetAndMetadata> offsets = new ConcurrentHashMap<>();
private static final ExecutorService workers = Executors.newFixedThreadPool(8);
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "slow-processor");
props.put("enable.auto.commit", "false");
props.put("max.poll.records", "50"); // 少量获取,快速分发给 Worker
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("orders"));
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
workers.submit(() -> {
processOrder(record);
// 记录已完成 Offset
TopicPartition tp = new TopicPartition(record.topic(), record.partition());
offsets.merge(tp, new OffsetAndMetadata(record.offset() + 1),
(old, newVal) -> old.offset() > newVal.offset() ? old : newVal);
});
}
}
} finally {
consumer.close();
workers.shutdown();
}
}
private static void processOrder(ConsumerRecord<String, String> record) {
// 调用外部服务,耗时 ~2s
}
}
操作前后对比:
| 模式 | 消费延迟(P99) | Offset 准确性 |
|---|---|---|
| 单线程 | 8s | 精确 |
| 单 Consumer + 8 Worker | 2s | 需手动追踪 Offset |
关键注意事项:Worker 线程处理无序。如果 Offset 5 在 Offset 3 之前完成,不能简单提交到 5——这会导致 Offset 3 的消息在故障后丢失。正确的做法是追踪连续已完成的最大 Offset。
易错场景
易错 1:多线程共享同一个 Consumer 实例
// 错误:Consumer 不是线程安全的!
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(topic);
// 线程 1
executor1.submit(() -> consumer.poll(Duration.ofMillis(1000)));
// 线程 2
executor2.submit(() -> consumer.poll(Duration.ofMillis(1000)));
// → ConcurrentModificationException 或数据错乱
规则:一个 Consumer 实例 = 一个线程。不可跨线程共享。
易错 2:提交的 Offset 覆盖了未处理完的消息
场景:Worker 线程中 Offset 10 先处理完,Offset 9 还在处理中,但代码盲目提交了 Offset 11。
后果:Consumer 崩溃重启后从 Offset 11 开始,Offset 9 的消息丢失。
正确做法:只提交连续已完成的最小 Offset。维护一个 TreeSet<Long> 追踪已完成但未提交的 Offset,提交时取连续前缀的最大值。
面试高频考点
Q:Consumer 多线程消费时,为什么官方推荐多个 Consumer 实例而非多线程共享?
A:
- 线程安全:Consumer 不是线程安全的,多线程共享需要复杂的同步逻辑
- Rebalance 友好:每个 Consumer 实例有自己的心跳和 Offset 提交,更符合消费组模型
- 故障隔离:一个 Consumer 崩溃不影响其他 Consumer 的分区
- 简单可靠:多 Consumer 实例不需要额外的 Offset 追踪逻辑
唯一的例外:分区数很少但单条消息处理极慢的场景(如上述示例二),才需要单 Consumer + 处理线程池。