【问题标题】:spring-kafka - error handling via group rebalancespring-kafka - 通过组重新平衡进行错误处理
【发布时间】:2021-05-24 01:17:00
【问题描述】:

我正在编写一个 kafka 侦听器,它应该只是将消息从主题转发到 jms 队列。我只需要为自定义异常JmsBrokerConnectionException 停止处理新消息,但我想继续为任何其他异常(即无效数据)处理新消息并将错误消息发送到 DLT。 我使用的是spring-kafka 2.2.7,无法升级。

我目前有一个解决方案,它使用:

  • SeekToCurrentErrorHandler(配置了 0 次重试和 DeadLetterPubishingRecoverer
  • @KafkaListener 方法中使用的重试模板,配置了 Integer.MAX_VALUE 重试,仅重试 JmsBrokerConnectionException
  • MANUAL_IMMEDIATE 确认

该解决方案似乎可以完成这项工作,但它的缺点是,对于 jms 代理的长时间中断,它会导致每个 max.poll.interval.ms 重新平衡(即 5 分钟)。

问题: 让max.poll.interval.ms 过期并进行组重新平衡来处理您想要停止消息消费的错误情况是个好主意吗?

我没有高吞吐量要求。 输入主题有 10 个分区,我将有 2 个消费者。 我知道还有其他使用有状态重试或暂停/恢复容器的解决方案,但我想继续使用当前的解决方案,除非我错过了它的任何主要缺点。

【问题讨论】:

    标签: apache-kafka spring-kafka


    【解决方案1】:

    我使用的是spring-kafka 2.2.7,无法升级。

    不再支持该版本。

    2.3 版向 STCEH 添加了退避和异常分类,无需在侦听器级别使用重试模板。

    也就是说,您可以将有状态重试 (https://docs.spring.io/spring-kafka/docs/current/reference/html/#stateful-retry) 与始终重试的 STCEH 一起使用,并在侦听器级别在 RecoveryCallback 中进行死信发布。消费者记录在重试上下文中可用RetryingMessageListenerAdapter.CONTEXT_RECORD 键。

    由于您正在执行手动确认,因此您还需要通过 CONTEXT_CONSUMER 键提交偏移量。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-02-21
      • 2018-12-21
      • 2018-01-20
      • 1970-01-01
      • 1970-01-01
      • 2021-09-08
      • 2022-10-05
      • 2019-09-26
      相关资源
      最近更新 更多