Kafka __consumer_offsets磁盘占用过大且分布不均问题求助
__consumer_offsets主题磁盘占用过大且分布不均的原因分析与解决办法
环境信息
- Kafka版本:MSK 2.8.1
- 命令行工具:Confluent 6.2.10
排查结果
通过命令行工具获取的__consumer_offsets主题磁盘分布如下:
# 各Broker总占用 $ kafka-log-dirs --describe --bootstrap-server $BOOTSTRAP --topic-list __consumer_offsets | grep '^{' | jq -r '.brokers[] | ["broker", .broker, "=", (([.logDirs[].partitions[].size] | add // 0) | . / 10000 | round | ./ 100), "MB" ] | @tsv' | paste -sd , | tr ' ' ' ' broker 1 = 459.72 MB,broker 2 = 218.95 MB,broker 3 = 134346.48 MB # 各Broker下__consumer_offsets分区的单独占用(单位:MB) $ kafka-log-dirs --describe --bootstrap-server $BOOTSTRAP --topic-list __consumer_offsets | grep '^{' | jq -r '.brokers[] | ["broker", .broker, "=", (.logDirs[].partitions[].size / 1000000 | round)] | @tsv' | tr ' ' ' ' broker 1 = 52 1 0 0 1 0 0 243 102 0 2 0 3 0 0 0 0 1 0 0 0 0 0 0 0 1 0 0 0 1 47 4 0 0 0 0 0 0 5 0 0 0 0 1 2 3 1 0 0 0 broker 2 = 52 1 0 0 1 0 0 2 102 0 2 0 3 0 0 0 0 1 0 0 0 0 0 0 0 1 0 0 0 1 47 4 0 0 0 0 0 0 5 0 0 0 0 1 2 3 1 0 0 0 broker 3 = 133907 1 0 0 8 3 1 31 10 4 2 0 27 0 2 0 14 8 4 4 1 0 3 2 0 10 0 0 3 14 35 123 0 0 2 0 0 0 23 0 0 0 0 25 26 39 9 3 6 5 # __consumer_offsets主题配置 $ kafka-topics --bootstrap-server $BOOTSTRAP --describe --topic __consumer_offsets Topic: __consumer_offsets TopicId: ... PartitionCount: 50 ReplicationFactor: 3 Configs: compression.type=producer,min.insync.replicas=2,cleanup.policy=compact,segment.bytes=104857600,message.format.version=2.8-IV1,max.message.bytes=10485880,unclean.leader.election.enable=true
原因分析
消费者组哈希分布不均:__consumer_offsets的分区由消费者组ID哈希计算分配(公式:
partition = murmur2(groupId) % 分区数),Broker3上的大分区对应一个或多个超级消费者组——这类组可能订阅了大量主题分区,或频繁提交offset,导致该分区的offset记录量远超其他分区。日志压缩未正常清理旧数据:尽管配置了
cleanup.policy=compact,但存在以下可能:- 该分区的活跃segment未及时滚动:当前
segment.bytes=100MB,若写入量极大,活跃segment一直无法被归档,而压缩仅针对非活跃segment,旧数据无法被清理。 - 存在未被覆盖的无效记录:部分消费者组已停止使用,但它们的旧offset记录没有被新记录覆盖,compact机制无法自动删除这些无后续更新的记录。
- 日志清理线程资源不足:默认的log cleaner线程数无法及时处理大分区的压缩任务,导致旧数据累积。
- 该分区的活跃segment未及时滚动:当前
解决办法
临时缓解:快速清理大分区数据
手动触发日志压缩:
临时调小压缩触发阈值,强制log cleaner处理非活跃segment:kafka-configs --bootstrap-server $BOOTSTRAP --alter --entity-type topics --entity-name __consumer_offsets --add-config min.cleanable.dirty.ratio=0.01等待1-2小时(根据数据量调整)后,恢复默认阈值:
kafka-configs --bootstrap-server $BOOTSTRAP --alter --entity-type topics --entity-name __consumer_offsets --delete-config min.cleanable.dirty.ratio定位并清理无效消费者组:
- 确定Broker3上大分区的编号(从第二个命令输出可知,Broker3的第一个数值
133907对应__consumer_offsets的某个分区,需根据实际输出确认具体编号)。 - 查看该分区的消息key,找到对应的消费者组:
kafka-console-consumer --bootstrap-server $BOOTSTRAP --topic __consumer_offsets --partition <大分区编号> --from-beginning --property print.key=true --property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer --property value.deserializer=org.apache.kafka.common.serialization.StringDeserializer | head -100 - 检查该消费者组状态,若已废弃则删除:
kafka-consumer-groups --bootstrap-server $BOOTSTRAP --describe --group <groupId> kafka-consumer-groups --bootstrap-server $BOOTSTRAP --delete --group <groupId>
删除后,log cleaner会自动清理该组对应的所有旧offset记录。
- 确定Broker3上大分区的编号(从第二个命令输出可知,Broker3的第一个数值
长期优化:避免问题复发
调整__consumer_offsets分区数:
当前50个分区无法均匀负载,可增加至100-200个分区,步骤如下:- 创建
topics.json文件:{"topics": [{"topic": "__consumer_offsets"}], "version":1} - 生成分区重分配方案:
kafka-reassign-partitions --bootstrap-server $BOOTSTRAP --generate --topics-to-move-json-file topics.json --broker-list "1,2,3" - 执行重分配:
kafka-reassign-partitions --bootstrap-server $BOOTSTRAP --execute --reassignment-json-file reassignment.json - 验证重分配完成:
kafka-reassign-partitions --bootstrap-server $BOOTSTRAP --verify --reassignment-json-file reassignment.json
注:新分区仅对后续创建的消费者组生效,旧组仍绑定原分区,需配合清理无效组来释放旧分区空间。
- 创建
优化日志压缩配置:
- 调小
segment.bytes:将其从100MB改为50MB,加快segment滚动频率,让压缩更及时。 - 增加
log.cleaner.threads:在MSK集群配置中增加日志清理线程数(建议设为2-4),提升压缩处理能力。 - 启用混合清理策略:修改
cleanup.policy=compact,delete,并设置retention.ms=604800000(7天),自动删除超过保留期的未覆盖旧记录。
- 调小
定期监控与清理:
- 定期检查__consumer_offsets各分区的磁盘占用,及时发现异常增长。
- 清理长期未活跃的消费者组,避免无效offset记录累积。
内容的提问来源于stack exchange,提问作者dlipofsky
相关产品推荐
相关产品推荐

