【问题标题】:Remove In-memory kafka records receive from fetch in Springboot project从 Springboot 项目中的 fetch 中删除 In-memory kafka 记录
【发布时间】:2020-05-10 04:55:34
【问题描述】:

我有一个要求,如果我可以删除内存中的 kafka 消息,因为我有 max-poll-records: 10. 所以场景是:在处理记录时,如果我的程序遇到任何错误,我不需要再处理存储在内存中的任何剩余记录。 例如:我一次获取 10 条记录,因为我的 max-poll-interval 为 10。我成功处理了 5 条记录(手动提交)但在第 6 条记录期间我遇到错误,现在我必须从内存中删除所有剩余的 5 条记录.下面是我的监听代码:

    @KafkaListener(topics = "#{'${kafka.consumer.allTopicList}'.split(',')}", groupId = Constant.GROUP_ID)
public void consumeAllTopics(@Header(KafkaHeaders.RECEIVED_TOPIC) String topic, String message, Acknowledgment acknowledgment) {
    switch (topic) {
        case Constant.toipc1:
            if (!StringUtils.isEmpty(message)) {
                try {
                      //processing logic
                      acknowledgment.acknowledge();
                    }catch(Exception e){e.printStackTrace();}
                }

我想通过代码删除记录。请帮助我了解是否可能,如果可以,我该如何实现。

【问题讨论】:

    标签: spring-boot apache-kafka kafka-consumer-api spring-kafka


    【解决方案1】:

    不清楚您所说的“删除”是什么意思。如果您的意思是“忽略”或“跳过”,则需要引发异常并配置自定义错误处理程序。

    the documentation

    如果ErrorHandler 实现RemainingRecordsErrorHandler,则错误处理程序将提供失败的记录和之前的poll() 检索到的任何未处理的记录。处理程序退出后,这些记录不会传递给侦听器。

    没有标准的错误处理程序来“跳过”剩余的记录;这是一个不寻常的要求。

    大多数人会使用SeekToCurrentErrorHandler(现在是即将发布的 2.5 版本中的默认值)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2023-03-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-19
      • 2021-09-15
      • 2019-07-19
      相关资源
      最近更新 更多