Consumer 配置
定义与作用
本节提供 Consumer 关键配置参数的速查表与调优指南,聚焦于影响吞吐、延迟、可靠性、Rebalance 行为的核心参数。
配置速查表
基本连接参数
| 参数 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
bootstrap.servers | (必填) | Broker 地址列表 | 至少填两个,防止单点连接失败 |
group.id | 无 | 消费组 ID | 必填(手动分配 assign 除外) |
client.id | "" | 客户端标识 | 建议设置,便于日志追踪 |
key.deserializer | (必填) | Key 反序列化器 | 与 Producer 序列化器配对 |
value.deserializer | (必填) | Value 反序列化器 | 同上 |
消费行为控制
| 参数 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
auto.offset.reset | latest | 无已提交 Offset 时的策略:latest/earliest/none | 新消费组需要全量数据用 earliest |
enable.auto.commit | true | 是否自动提交 Offset | 可靠性优先用 false |
auto.commit.interval.ms | 5000 | 自动提交间隔 | 降低可减少故障时的重复量 |
max.poll.records | 500 | 单次 poll() 最大记录数 | 消息处理慢时可调低;高吞吐调高 |
max.poll.interval.ms | 300000 (5min) | 两次 poll() 最大间隔 | 处理极慢时上调 |
session.timeout.ms | 45000 (45s) | 心跳超时 | Consumer 掉线检测时间 |
heartbeat.interval.ms | 3000 (3s) | 心跳发送间隔 | 应为 session.timeout.ms 的 1/3 |
拉取性能
| 参数 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
fetch.min.bytes | 1 | 最小拉取字节(减少空轮询) | 高吞吐设 1024-10240 |
fetch.max.bytes | 52428800 (50MB) | 单次拉取最大字节 | 消息体大时需调大 |
fetch.max.wait.ms | 500 | 满足 fetch.min.bytes 前的最大等待 | 低延迟设 100,高吞吐设 1000 |
max.partition.fetch.bytes | 1048576 (1MB) | 单分区最大拉取字节 | 需同时调整 Broker max.message.bytes |
partition.assignment.strategy | RangeAssignor | 分区分配策略 | 推荐 CooperativeStickyAssignor |
可靠性参数
| 参数 | 默认值 | 说明 |
|---|---|---|
isolation.level | read_uncommitted | read_committed 仅读取已提交事务的消息 |
enable.auto.commit | true | false + 手动提交 = 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.bytes | 1 | 1 |
fetch.max.wait.ms | 500 | 100 |
max.poll.records | 500 | 50 |
| 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/3session.timeout.ms:Broker 在此时间内没收心跳,认为 Consumer 已死max.poll.interval.ms:两次poll()调用的最大间隔。与心跳无关——即使心跳正常,长时间不调用poll()也会被踢出
典型故障:Consumer 心跳正常(每 3 秒发送),但因处理慢导致 6 分钟才调用一次 poll() → 被 max.poll.interval.ms(默认 5 分钟)踢出。