消费组、位点与积压治理
同一消费组内,一个分区同一时刻只会分配给一个消费者。消费者实例多于分区数不会提高吞吐,实例少于分区数则会让单实例承载多个分区。Lag 应同时看当前值和增长速率。

1. 查看消费组
# 列出全部消费组,快速定位要排查的业务组。
/opt/kafka/bin/kafka-consumer-groups.sh --list \
--bootstrap-server kafka-1.example.internal:9092
# 查看每分区已提交位点、日志末端位点、Lag 与当前消费者。
/opt/kafka/bin/kafka-consumer-groups.sh --describe \
--bootstrap-server kafka-1.example.internal:9092 \
--group orders-projection
# 查看成员与分区分配,定位频繁重平衡或失联实例。
/opt/kafka/bin/kafka-consumer-groups.sh --describe \
--bootstrap-server kafka-1.example.internal:9092 \
--group orders-projection --members --verbose
CURRENT-OFFSET 为空不一定是故障,可能代表消费组从未提交位点或消费者未启动。需结合 Topic、应用日志和分区分配判断。
2. 消费者可靠性配置
# consumer.properties
# 使用稳定业务组名,不能为每次部署生成随机 group.id。
group.id=orders-projection
# 业务处理成功后由应用显式提交,避免自动提交提前确认消息。
enable.auto.commit=false
# 仅在没有已提交位点时生效;关键业务回放前需验证幂等。
auto.offset.reset=earliest
# 处理时间必须小于该值,否则会触发重平衡。
max.poll.interval.ms=300000
# 单次拉取量须结合应用处理耗时、堆内存与下游容量调优。
max.poll.records=500
处理成功但位点提交前崩溃会产生重复消费。消费者应以事件 ID、业务主键或幂等写入保证安全,不能只依赖 Kafka 位点。
3. 受控回放
回放必须停止该消费组全部实例,先使用 --dry-run 审核再执行。变更单应记录目标 Topic、Group、回放时间点、下游影响和回滚方案。
# 预览位点重置计划;时间仅为 UTC 格式示例,必须替换为审批后的时间点。
/opt/kafka/bin/kafka-consumer-groups.sh --reset-offsets --dry-run \
--bootstrap-server kafka-1.example.internal:9092 \
--group orders-projection --topic orders.created.v1 \
--to-datetime 2026-07-31T00:00:00.000
# 审核 dry-run 输出且确认消费者已停止后,才执行实际重置。
/opt/kafka/bin/kafka-consumer-groups.sh --reset-offsets --execute \
--bootstrap-server kafka-1.example.internal:9092 \
--group orders-projection --topic orders.created.v1 \
--to-datetime 2026-07-31T00:00:00.000
| 现象 | 优先检查 |
|---|---|
| Lag 持续增长 | 上游写入、消费者处理、下游依赖、分区数与并发 |
| 单分区 Lag 特别高 | Key 倾斜、热点实体、Leader 性能和单消费者异常 |
| 重平衡频繁 | max.poll.interval.ms、心跳、发布节奏和网络抖动 |
| 位点提交失败 | ACL、Coordinator、超时和应用异常 |
4. 重平衡与扩缩容
滚动发布、长时间处理和频繁弹性扩缩容都会触发重平衡。应用应在终止信号到来后停止拉取、完成当前批次、提交已处理位点后退出;部署系统需提供足够的优雅终止时间。
# 客户端支持时使用协作式分配,减少扩缩容造成的全量分区撤销。
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
# 会话与心跳参数需结合网络质量、容器暂停和发布策略评估。
session.timeout.ms=45000
heartbeat.interval.ms=15000
Lag 恢复时间约为 当前积压 / (消费速率 - 新增写入速率)。若消费速率不高于写入速率,等待不会恢复,必须先消除下游瓶颈或增加有效分区并行度。