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

    • 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 是从 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 Offsetorders LEOLag
消费前0100100
消费 100 条后1001000

示例二:指定分区消费(手动分配)

场景:不希望加入消费组(不触发 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 被认为已死亡并被踢出消费组。

解决:

  1. 减少单次 poll() 的消息数(max.poll.records,默认 500)
  2. 将处理逻辑移到独立线程池,poll() 线程快速返回继续拉取
  3. 增加 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:

  1. 消费者自主控制速率 — Consumer 根据自己的处理能力拉取,不会因 Broker 推送过快而 OOM
  2. 天然支持批量 — Producer 端攒批 + Consumer 端批量 poll,两端优化不相互影响
  3. 简化 Broker — Broker 不需要追踪每个 Consumer 的消费状态和推送队列
  4. 重放能力 — Consumer 可以 seek 到任意 Offset 重新消费,Push 模型下难以实现

代价:Pull 模型在低延迟场景有劣势——Consumer 需要轮询,可能引入延迟。Kafka 通过长轮询(fetch.min.bytes + fetch.max.wait.ms)缓解此问题。

上一页
章节导读
下一页
Consumer Group 消费组