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

Kafka COMPACT清理策略未生效,如何保留topic2最新Key消息?

Kafka Compact清理策略不生效问题排查与解决

主题配置信息

topic1配置

Topic: topic1     TopicId: <topic_ID> PartitionCount: 1       ReplicationFactor: 1    Configs: min.insync.replicas=1,cleanup.policy=delete,retention.ms=24192000000,retention.bytes=-1
        Topic: topic1     Partition: 0    Leader: 0       Replicas: 0     Isr: 0

topic2配置

Topic: topic2    TopicId: <topic_ID> PartitionCount: 1       ReplicationFactor: 1    Configs: min.insync.replicas=1,cleanup.policy=compact,retention.ms=60000,retention.bytes=-1
        Topic: topic2     Partition: 0    Leader: 0       Replicas: 0     Isr: 0

测试消息内容

message1

Key:

{
    "COl1": {
        "string": "dummy001"
    }
}

Value:

{
    "data": {
        "Col1": "dummy001",
        "Col2": "2023-01-01 01:01:01",
        "Col3": "1000.100"
    },
    "operation": "INSERT",
    "timestamp": "2023-01-01 01:01:01"
}

message2

Key:

{
    "COl1": {
        "string": "dummy001"
    }
}

Value:

{
    "data": {
        "Col1": "dummy001",
        "Col2": "2023-02-02 02:02:02",
        "Col3": "2000.100"
    },
    "operation": "UPDATE",
    "timestamp": "2023-02-02 02:02:02"
}

Kafka流创建与数据写入

stream1数据内容

+--------------------------------------------------------------------------------------------------------------------------------------------------+
|DATA                                                                                                                                              |
+--------------------------------------------------------------------------------------------------------------------------------------------------+
|{COL1=dummy001, COL2=2023-01-01 01:01:01, COL3=1000.100}                                                                                          |
|{COL1=dummy001, COL2=2023-02-02 02:02:02, COL3=2000.100}               

创建stream2并写入数据

创建stream2的SQL:

ksql> create stream stream2 (col1 string key, col2 string, col3 string) with(kafka_topic='topic2', value_format='AVRO');  

从stream1向stream2插入数据的SQL:

ksql> insert into stream2 select data->COL1 as COL1, data->COL2 as COL2, data->COl3 as COL3  from stream1 partition by data->col1;

topic2中的最终消息

message1

Key: dummy001
Value:

{
    "COL2": {
        "string": "2023-01-01 01:01:01"
    },
    "COL3": {
        "string": "1000.100"
    }
}

message2

Key: dummy001
Value:

{
    "COL2": {
        "string": "2023-02-02 02:02:02"
    },
    "COL3": {
        "string": "2000.100"
    }
}

问题描述

topic2设置了cleanup.policy=compact,但旧消息并未被清理。即使将col1设为Key列,压缩策略仍不生效;若设置为compact,delete则会删除所有消息。期望实现topic2中始终保留每个col1对应的最新消息,询问是否遗漏配置步骤。

解决方案

1. 明确日志压缩触发逻辑

Kafka日志压缩并非实时执行,需满足以下条件才会触发:

  • 日志段文件达到segment.bytes(默认1GB)大小,或达到segment.ms(默认7天)时间阈值,才会触发压缩检查
  • 后台log-cleaner线程需处于启用状态(默认log.cleaner.enable=true,无需额外配置)

你的topic2设置了retention.ms=60000(1分钟),该参数在compact策略下控制未被更新的旧消息的最长保留时间,而非压缩触发时间,这是核心误解点。

2. 调整配置加速压缩触发

若需更快验证压缩效果,可临时修改topic2的以下配置:

kafka-configs.sh --bootstrap-server <你的Bootstrap地址> --alter --topic topic2 --add-config segment.ms=60000,log.cleaner.min.cleanable.ratio=0.01,log.cleaner.backoff.ms=1000
  • segment.ms=60000:日志段1分钟滚动一次,触发压缩检查
  • log.cleaner.min.cleanable.ratio=0.01:降低可清理比例,让压缩更容易触发
  • log.cleaner.backoff.ms=1000:缩短压缩线程的检查间隔

3. 验证Key一致性

确保写入topic2的消息Key完全一致(你的场景中都是dummy001,此部分无问题),若Key存在大小写、格式差异,压缩会判定为不同Key,不会清理旧消息。

4. 正确使用compact,delete组合策略

当设置cleanup.policy=compact,delete时,retention.ms会作为所有消息的最大保留时间。你的topic2设置了retention.ms=60000,1分钟后所有消息都会被删除,这就是全量删除的原因。如需同时使用两种策略,需将retention.ms设为较大值(比如和topic1一致的24192000000,即28天),这样未被更新的消息会在28天后删除,被更新的旧消息会被压缩清理。

5. 手动触发压缩(测试用)

若需立即验证压缩效果,可手动触发日志滚动强制压缩:

kafka-run-class.sh kafka.admin.LogRoller --bootstrap-server <你的Bootstrap地址> --topic topic2

执行后等待数分钟,再查看topic2消息,应仅保留最新的一条。


内容的提问来源于stack exchange,提问作者Dinesh Kumar L

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 09:54:51