【问题标题】:How do I delete/clean Kafka queued messages without deleting Topic如何在不删除主题的情况下删除/清理 Kafka 排队消息
【发布时间】:2018-02-22 21:08:24
【问题描述】:

有什么方法可以在不删除 Kafka 主题的情况下删除队列消息?
我想在激活消费者时删除队列消息。

我知道有几种方法:

  1. 重置保留时间

    $ ./bin/kafka-topics.sh --zookeeper localhost:2181 --alter --topic MyTopic --config retention.ms=1000

  2. 删除kafka文件

    $ rm -rf /data/kafka-logs/<topic/Partition_name>

【问题讨论】:

  • 您首先提到的保留时间技巧要好得多。第二种方式会导致重复的主题出现问题,并使主题的元数据与现实不一致。请注意,偏移量不会回到零。

标签: apache-kafka


【解决方案1】:

在 0.11 或更高版本中,您可以运行 bin/kafka-delete-records.sh 命令来标记要删除的消息。

https://github.com/apache/kafka/blob/trunk/bin/kafka-delete-records.sh

例如发布100条消息

seq 100 | ./bin/kafka-console-producer.sh --broker-list localhost:9092 --topic mytest

然后使用新的kafka-delete-records.sh 删除这 100 条消息中的 90 条 命令行工具

./bin/kafka-delete-records.sh --bootstrap-server localhost:9092 --offset-json-file ./offsetfile.json

offsetfile.json 包含在哪里

 {"partitions": [{"topic": "mytest", "partition": 0, "offset": 90}], "version":1 }

然后从头开始消费消息,以验证 100 条消息中有 90 条确实标记为已删除。

./bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic mytest --from-beginning
91
92
93
94
95
96
97
98
99
100

【讨论】:

  • 感谢汉斯的回复!它非常接近我一直想要的。您是否知道我是否可以在不知道有多少消息排队的情况下删除所有偏移量?我可以通过编辑 json 来做到这一点吗?
  • 是的,您可以删除所有消息。您也可以跳过使用此工具,查看源代码并编写自己的程序,直接调用相同的 API 以删除任何给定偏移量(包括最新偏移量)之前的记录,或者您可以按时间戳查找偏移量以删除所有记录在特定时间之前。这个工具使用的 API 应该在 2017 年 10 月发布的 Kafka 1.0 中得到更好的记录
  • 在 Kafka 中只有一个共享提交日志,因此如果您删除消息,每个人都会独立于任何消费者组而删除这些消息。如果您只想更改一个消费者组的偏移量以跳过消息但将它们留给其他消费者,请使用 bin/kafka-consumer-groups.sh --reset-offsets。在此处查看 KIP-122 详细信息cwiki.apache.org/confluence/display/KAFKA/…
  • 谢谢分享。它有很大帮助。有什么方法可以删除队列开头和结尾之间的一系列消息。假设我的消息从 0 变为 100,我想将消息 89 和 90 标记为删除。谢谢!
  • 不,这在 Kafka 中是不可能的。 Kafka 有一个核心属性,即它是一个不可变的提交日志,并且允许任意范围删除会使其不那么不可变。
【解决方案2】:

要删除特定主题中的所有消息,您可以运行kafka-delete-records.sh

例如,我有一个名为 test 的主题,它有 4 partitions

创建一个Json 文件,例如j.json

{

"partitions": [

    {

        "topic": "test",

        "partition": 0,

        "offset": -1

    }, {

        "topic": "test",

        "partition": 1,

        "offset": -1

    }, {

        "topic": "test",

        "partition": 2,

        "offset": -1

    }, {

        "topic": "test",

        "partition": 3,

        "offset": -1

    }

],

"version": 1

}

现在通过此命令删除所有消息:

/opt/kafka/confluent-4.1.1/bin/kafdelete-records --bootstrap-server 192.168.XX.XX:9092 --offset-json-file j.json

执行命令后会显示此信息

Records delete operation completed:
partition: test-0   low_watermark: 7
partition: test-1   low_watermark: 7
partition: test-2   low_watermark: 7
partition: test-3   low_watermark: 7

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-12-18
    • 2019-10-04
    • 1970-01-01
    • 2020-02-16
    • 2018-07-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多