Kafka日志压缩后消息及段文件留存超出预期问题咨询
Kafka日志压缩与retention.bytes相关疑问
环境与操作步骤
- 在
broker.properties中配置:log.retention.check.interval.ms=1000 - 创建主题命令:
kafka-topics.sh --bootstrap-server localhost:9092 --create --topic compact_topic --config cleanup.policy=compact --config segment.bytes=80 --config retention.bytes=150 --config max.compaction.lag.ms=1000
- 消息发送情况:通过
kafka-console-producer.sh发送单字符消息,单条消息包约占70字节,每条消息都会触发新段创建。发送内容包括:- 1条key为“1”的消息:
1:1 - 3条key为“2”的消息:
2:A、2:B、2:C
- 1条key为“1”的消息:
现象与疑问
压缩完成后,留存了3条消息(1:1、2:B、2:C)以及3个非空段文件(其中一个为活跃段)。疑问:设置retention.bytes=150后,为何会留存3个非空段而非2个?
后续添加消息后,压缩仍会在1-2分钟内回到3条消息的状态(仅存在2个key,但同一key留存了2条消息),多次尝试后的消费输出:
kafka-console-consumer.sh --bootstrap-server localhost:9092 -topic compact_topic --from-beginning --property print.key=true --property key.separator=: 1:1 2:H 2:I
解答
1. 为何留存3个非空段而非2个?
核心原因在于cleanup.policy的配置逻辑以及Kafka日志压缩的特性:
- cleanup.policy的作用限制:当仅设置
cleanup.policy=compact时,retention.bytes这类基于大小的删除配置不会生效。Kafka只会执行日志压缩(保留每个key的最新版本),不会因为日志总大小超过阈值而删除旧段。只有将cleanup.policy设置为compact,delete时,才会同时启用压缩和基于大小/时间的删除策略。 - 段的处理规则:日志压缩以非活跃段为单位进行合并,且永远不会处理当前正在写入的活跃段。你设置的
segment.bytes=80接近单条消息大小,每条消息都会触发新段创建,最终:1:1所在的段因无其他同key消息,被完整保留;2:A所在的段被压缩合并到2:B的段中,仅保留2:B;2:C所在的段是活跃段,不会被压缩或删除;
因此最终留存3个非空段,对应3条保留的消息。
2. 为何同一key会留存两条消息?
这源于活跃段的豁免机制和压缩操作的延迟性:
- 新发送的消息(如
2:H、2:I)会写入新的活跃段,该段不会被压缩; - 之前的
2:H所在的段会变为非活跃段,但压缩合并操作不会立即执行(受log.cleaner.backoff.ms等参数控制,默认10秒); - 压缩完成前,活跃段中的
2:I和非活跃段中的2:H会同时存在,因此消费时会看到两条同key的消息;当压缩最终执行时,会合并非活跃段中的旧版本消息,仅保留最新版本。
内容的提问来源于stack exchange,提问作者FER_Viborg
相关产品推荐
相关产品推荐

