乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • 学习路径
  • 第1章 消息队列与 Kafka 概述

    • 章节导读
    • 消息队列基础
    • Kafka 概述
    • Kafka 为什么快
  • 第2章 快速上手:单机环境搭建

    • 章节导读
    • 环境准备与安装
    • 快速启动
    • Topic 管理
  • 第3章 核心概念:主题、分区与日志

    • 章节导读
    • Topic
    • Partition
    • Offset
    • Segment 与存储结构
  • 第4章 生产者详解

    • 章节导读
    • Producer 概述
    • Producer 发送机制
    • Producer 分区策略
    • Producer 幂等与事务
    • Producer 配置
  • 第5章 消费者与消费组

    • 章节导读
    • Consumer 概述
    • Consumer Group 消费组
    • Consumer 分区分配策略
    • Consumer Offset 提交
    • Consumer 多线程
    • Consumer 配置
  • 第6章 Broker 与控制器

    • 章节导读
    • Broker 概述
    • Broker 配置
    • Controller
    • KRaft 共识协议
  • 第7章 副本与数据可靠性

    • 章节导读
    • 副本机制
    • ISR 与副本同步
    • Leader 选举
    • ACK 与一致性保证
    • 高水位与 Leader Epoch
  • 第8章 存储与性能优化

    • 章节导读
    • 存储架构
    • 日志清理与压缩
    • Page Cache 与零拷贝
    • Producer 性能优化
    • Consumer 性能优化
    • Broker 性能优化
  • 第9章 生产环境运维与监控

    • 章节导读
    • Topic 管理
    • Kafka 运维工具
    • 监控
    • 常见故障排查
  • 第10章 Kafka生态与面试考点

    • 章节导读
    • Kafka 生态全景
    • 面试高频 30 题

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
66 Consumer → 各 2 分区~300K events/s
1212 Consumer → 各 1 分区~500K events/s
1312 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 Worker2s需手动追踪 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:

  1. 线程安全:Consumer 不是线程安全的,多线程共享需要复杂的同步逻辑
  2. Rebalance 友好:每个 Consumer 实例有自己的心跳和 Offset 提交,更符合消费组模型
  3. 故障隔离:一个 Consumer 崩溃不影响其他 Consumer 的分区
  4. 简单可靠:多 Consumer 实例不需要额外的 Offset 追踪逻辑

唯一的例外:分区数很少但单条消息处理极慢的场景(如上述示例二),才需要单 Consumer + 处理线程池。

上一页
Consumer Offset 提交
下一页
Consumer 配置