【问题标题】:Kafka: How to delete records from a topic using Java API?Kafka:如何使用 Java API 从主题中删除记录?
【发布时间】:2022-03-04 15:48:34
【问题描述】:

我正在寻找一种从 Kafka 主题中删除(完全删除)已使用记录的方法。我知道有几种方法可以做到这一点,例如,通过更改主题的保留时间或删除 Kafka-logs 文件夹。但我正在寻找的是一种使用 Java API 删除一定数量的主题记录的方法,如果可能的话。

我已尝试测试 AdminClient API,特别是 adminclient.deleteRecords(recordsToDelete) 方法。但如果我没记错的话,那个方法只是改变主题中的偏移量,而不是真正从硬盘中删除所述记录。

有没有真正从硬盘中删除记录的 Java API?

【问题讨论】:

    标签: java apache-kafka


    【解决方案1】:

    我可以删除。如果 linux 在机器上,它会从硬盘中删除它。当我从互联网上搜索时,我发现windows中有一个错误。但是,我在 windows 中找不到此错误的解决方案。 如果 kafka 在 linux 机器上运行,则此代码有效。

    Windows 错误链接:https://issues.apache.org/jira/browse/KAFKA-1194

    public void deleteMessages(String topicName, int partitionIndex, int beforeIndex) {
           TopicPartition topicPartition = new TopicPartition(topicName, partitionIndex);
           Map<TopicPartition, RecordsToDelete> deleteMap = new HashMap<>();
           deleteMap.put(topicPartition, RecordsToDelete.beforeOffset(beforeIndex));
           kafkaAdminClient.deleteRecords(deleteMap);
    }
    

    【讨论】:

      【解决方案2】:

      一开始我也有点困惑,为什么包含的 bin/kafka-delete-records.sh 可以删除,但我不能使用 Java API

      缺少的部分是您需要调用 KafkaFuture.get(),因为 deleteRecords 返回 Futures 的映射

      这是代码

      在这段代码中,你需要调用entry.getValue().get().lowWatermark()

      DeleteRecordsResult result = adminClient.deleteRecords(recordsToDelete);
      Map<TopicPartition, KafkaFuture<DeletedRecords>> lowWatermarks = result.lowWatermarks();
      try {
          for (Map.Entry<TopicPartition, KafkaFuture<DeletedRecords>> entry : lowWatermarks.entrySet()) {
              System.out.println(entry.getKey().topic() + " " + entry.getKey().partition() + " " + entry.getValue().get().lowWatermark());
          }
      } catch (InterruptedException | ExecutionException e) {
          e.printStackTrace();
      }
      adminClient.close();
      

      【讨论】:

        【解决方案3】:

        Kafka 主题是不可变的,这意味着您只能向它们添加新消息。本身没有删除。

        但是,为了避免“磁盘耗尽”,Kafka 提供了两个概念来降低主题的大小:保留策略和压缩。

        留存率 如果您有一个不需要永远使用数据的主题,您只需设置一个保留策略,无论您需要保留数据多长时间,即 72 小时。然后,Kafka 会自动为您删除超过 72 小时的消息。

        压缩 如果您确实需要数据永久保留,或者至少保留很长时间,但您只需要 latest 值,那么您可以将主题设置为压缩。只要使用已存在的密钥添加新消息,这将自动删除旧消息。

        规划 Kafka 架构的一个核心部分是考虑如何将数据存储在主题中。例如,如果您将更新推送到 kafka 主题中的客户记录,假设客户的上次登录日期(非常人为的示例......),那么您只对 LAST 条目感兴趣(因为所有之前的条目都没有更长的“最后”登录)。如果此分区键是客户 ID,并且启用了日志压缩,那么一旦用户登录并且 kafka 主题收到此事件,任何其他具有相同分区键(客户 ID)的先前消息将被自动删除来自主题。

        【讨论】:

          【解决方案4】:

          我在 Red Hat 7.6 上使用 Kafka 2.1.1,对 AdminClient.deleteRecords() 的调用确实有效地从 /tmp/kafka-logs 中的相应文件夹中删除了文件。剩下的唯一文件是leader-epoch-checkpoint,里面有关于最后一条记录偏移的信息:在我的例子中是96。

          请注意,在调用AdminClient.deleteRecords() 时,您不应传递大于分区现有高水位线的偏移量。如果这样做,调用将失败并显示"org.apache.kafka.common.errors.OffsetOutOfRangeException: The requested offset is not within the range of offsets maintained by the server.",但您不会知道,直到您尝试通过Future.get() 检查结果 - 有关详细信息,请参阅 Trix 的答案。

          【讨论】:

            【解决方案5】:

            Kafka 不支持从主题中删除记录。它的工作方式是构建一个消息缓冲区,该缓冲区随着消息推送到它而增长。而读取消息的客户端基本上只持有该缓冲区的偏移量。因此 Kafka 中的客户端基本上处于“只读”模式,无法更改缓冲区。考虑一个案例,当几个不同的客户端(不同的客户端组)读取相同的主题并且每个都保存自己的偏移量时。如果有人开始从设置偏移量的缓冲区中删除消息会发生什么。

            【讨论】:

            • 好的,谢谢,我想我现在明白了。但是在存储方面,Kafka 会一直运行到存储空间不足,还是会在某个时候开始删除旧记录/主题?还是由用户来处理?
            • 这取决于您的保留政策。您需要确保您的存储空间能够承受它。
            【解决方案6】:

            没有 Kafka 不提供删除主题中特定偏移量的功能,并且没有可用的 API。

            【讨论】:

              猜你喜欢
              • 1970-01-01
              • 1970-01-01
              • 2017-11-17
              • 2023-03-19
              • 2018-11-15
              • 2016-04-17
              • 1970-01-01
              • 1970-01-01
              • 2017-05-18
              相关资源
              最近更新 更多