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

    • 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 配置

定义与作用

本节提供 Consumer 关键配置参数的速查表与调优指南,聚焦于影响吞吐、延迟、可靠性、Rebalance 行为的核心参数。

配置速查表

基本连接参数

参数默认值说明调优建议
bootstrap.servers(必填)Broker 地址列表至少填两个,防止单点连接失败
group.id无消费组 ID必填(手动分配 assign 除外)
client.id""客户端标识建议设置,便于日志追踪
key.deserializer(必填)Key 反序列化器与 Producer 序列化器配对
value.deserializer(必填)Value 反序列化器同上

消费行为控制

参数默认值说明调优建议
auto.offset.resetlatest无已提交 Offset 时的策略:latest/earliest/none新消费组需要全量数据用 earliest
enable.auto.committrue是否自动提交 Offset可靠性优先用 false
auto.commit.interval.ms5000自动提交间隔降低可减少故障时的重复量
max.poll.records500单次 poll() 最大记录数消息处理慢时可调低;高吞吐调高
max.poll.interval.ms300000 (5min)两次 poll() 最大间隔处理极慢时上调
session.timeout.ms45000 (45s)心跳超时Consumer 掉线检测时间
heartbeat.interval.ms3000 (3s)心跳发送间隔应为 session.timeout.ms 的 1/3

拉取性能

参数默认值说明调优建议
fetch.min.bytes1最小拉取字节(减少空轮询)高吞吐设 1024-10240
fetch.max.bytes52428800 (50MB)单次拉取最大字节消息体大时需调大
fetch.max.wait.ms500满足 fetch.min.bytes 前的最大等待低延迟设 100,高吞吐设 1000
max.partition.fetch.bytes1048576 (1MB)单分区最大拉取字节需同时调整 Broker max.message.bytes
partition.assignment.strategyRangeAssignor分区分配策略推荐 CooperativeStickyAssignor

可靠性参数

参数默认值说明
isolation.levelread_uncommittedread_committed 仅读取已提交事务的消息
enable.auto.committruefalse + 手动提交 = At-least-once

完整示例

示例一:低延迟实时消费配置

场景:实时风控系统,要求消息延迟 P99 < 100ms。

Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092");
props.put("group.id", "risk-control");
props.put("key.deserializer", "...");
props.put("value.deserializer", "...");

// 低延迟配置
props.put("fetch.min.bytes", 1);           // 有消息立即拉取
props.put("fetch.max.wait.ms", 100);       // 最多等 100ms
props.put("max.poll.records", 50);         // 少量拉取,快速处理
props.put("enable.auto.commit", false);     // 手动提交

操作前后对比:

参数默认低延迟配置
fetch.min.bytes11
fetch.max.wait.ms500100
max.poll.records50050
P99 延迟~500ms~80ms

示例二:批量归档消费配置

场景:数据归档任务,每小时消费一次,追求高吞吐。

Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092");
props.put("group.id", "archiver");
props.put("key.deserializer", "...");
props.put("value.deserializer", "...");

// 高吞吐配置
props.put("fetch.min.bytes", 1048576);    // 至少 1MB 才返回
props.put("fetch.max.wait.ms", 1000);     // 最多等 1s
props.put("max.poll.records", 5000);      // 每次拉 5000 条
props.put("max.partition.fetch.bytes", 10485760);  // 10MB/分区

操作前后对比:

指标默认配置高吞吐配置
空轮询率~30%~5%
单次 poll 消息数~300~3000
网络效率~60%~90%

易错场景

易错 1:auto.offset.reset=earliest 但期望从最新开始

场景:新部署的 Consumer Group 希望只消费新消息,但看到了大量历史数据。

原因:auto.offset.reset=earliest 在没有已提交 Offset 时从头消费。

规则:

  • 新 Consumer Group + 只要新消息 → auto.offset.reset=latest
  • 新 Consumer Group + 需要全量历史 → auto.offset.reset=earliest

易错 2:fetch.max.bytes 设置小于单条消息大小

场景:消息体 2MB,但 max.partition.fetch.bytes=1MB(默认)。

后果:Consumer 无法消费该分区的消息,持续卡住。

诊断:

# 查看 Consumer Group 的 LAG 是否持续增长
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group mygroup --describe

面试高频考点

Q:session.timeout.ms、heartbeat.interval.ms、max.poll.interval.ms 三者的关系是什么?

A:

  • heartbeat.interval.ms:Consumer 发送心跳的频率。应设置为 session.timeout.ms 的 1/3
  • session.timeout.ms:Broker 在此时间内没收心跳,认为 Consumer 已死
  • max.poll.interval.ms:两次 poll() 调用的最大间隔。与心跳无关——即使心跳正常,长时间不调用 poll() 也会被踢出

典型故障:Consumer 心跳正常(每 3 秒发送),但因处理慢导致 6 分钟才调用一次 poll() → 被 max.poll.interval.ms(默认 5 分钟)踢出。

上一页
Consumer 多线程