You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

现象与疑问

压缩完成后,留存了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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 14:14:59