【问题标题】:Spring cloud stream kafka consumer error handling and retries issuesSpring Cloud Stream kafka消费者错误处理和重试问题
【发布时间】:2021-12-22 19:47:31
【问题描述】:

我在 spring cloud stream kafka binder 的错误处理场景中需要帮助。我的应用程序有 java 8 消费者,其绑定在 application.yaml 中指定。消费者写成:

@Bean
public Consumer<Message<Transaction>> doProcess() {

    return message -> {
        Transaction transaction = message.getPayload();
       
        if(true) {
            throw new RuntimeException("exception!! !!:)");
        }
       Acknowledgment acknowledgment = message.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT, 
       Acknowledgment.class);
       if (acknowledgment != null) {
           System.out.println("Acknowledgment provided");
           acknowledgment.acknowledge();
       }
  }
}

application.yaml:

spring.application.name: appname
spring.cloud.stream:
  function.definition: doProcess
  kafka:
    default.consumer:
      startOffset: latest
      useNativeDecoding: true
    bindings:
      input.consumer.autoCommitOffset: false

bindings:
  doProcess-in-0:
    destination: kafka.input.topic.name
    group: appGroup
    content-type: application/*+avro
    consumer:
      autoCommitOffset: false.

现在,我正在努力处理错误并遇到两个问题:

  1. 我正在尝试手动确认消息的消耗,而不是使用 autoCommitOffset 作为 true。因此,当我将 autoCommitOffset 设置为 false 并测试错误情况时,会遇到奇怪的行为,即每当抛出异常时,消息都会重试“n”次,即使在重新启动后也会发生这种重试/重新传递失败消息服务(如果在 n 重试完成之前重新启动)。并且一旦 n 重试完成,即使在重新启动服务后也不会选择消息。这是否意味着,消费者在重新试用/重新传递消息后提交偏移量,这不应该是这种情况,因为 autoCommitOffset 为 false。

    注意:我没有配置任何 dlq。

  2. 我们需要编写自定义异常处理程序,我们可以在其中捕获异常(应用程序代码和框架中的错误)并通过 AWS env 中的电子邮件向用户组发送通知。但是,我们找不到任何可以捕获这两种异常的错误处理程序。类似于扩展 SeekToCurrentErrorHandler 或任何其他可以在错误事件上调用的侦听器。

编辑:

根据 Gary 提供的解决方案,我们可以使用以下 bean 来配置自定义错误处理程序:

@Bean
    public ListenerContainerCustomizer<AbstractMessageListenerContainer> MQLCC() {
        System.out.println(String.format("DEBUG: Bean %s has bean created.", "MQLCC"));
        return new ListenerContainerCustomizerCustom ();
    }

    private static class ListenerContainerCustomizerCustom implements ListenerContainerCustomizer<AbstractMessageListenerContainer> {
        @Override
        public void configure(AbstractMessageListenerContainer container, String destinationName, String group) {
            System.out.println(String.format("HELLO from container %s, destination: %s, group: %s", container, destinationName, group));
        }

    }

【问题讨论】:

    标签: spring spring-boot apache-kafka spring-cloud-stream


    【解决方案1】:

    侦听器容器中的默认错误处理程序将重试 10 次,然后记录错误并丢弃记录;对于不同的行为,您需要配置自定义错误处理程序和恢复策略。使用ListenerContainerCustomizer bean 来配置容器。

    https://docs.spring.io/spring-kafka/docs/current/reference/html/#default-ehhttps://docs.spring.io/spring-kafka/docs/current/reference/html/#dead-letters

    (3.2 及更高版本)

    https://docs.spring.io/spring-kafka/docs/2.7.x/reference/html/#seek-to-currenthttps://docs.spring.io/spring-kafka/docs/2.7.x/reference/html/#dead-letters

    适用于早期版本。

    【讨论】:

    • 谢谢 Gary,我使用 ListenerContainerCustomizer bean 来配置错误处理程序并且它工作正常。
    • 最初它工作正常,但现在我观察到一个奇怪的行为,我能够在 DefaultErrorHandler 中捕获少数异常,而对于其他异常则不是由 DefaultErrorHandler 处理的。就像我从 Consumer 显式抛出 RuntimeException 一样,它并没有被捕获。你能帮我理解这种行为吗
    • 这对我来说没有任何意义;我建议您提出一个新问题,显示您当前的代码和配置以及有关您遇到的问题的更多详细信息。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-06-21
    • 2017-06-22
    • 2022-12-22
    • 2019-11-25
    • 1970-01-01
    • 2022-01-11
    • 2017-11-20
    相关资源
    最近更新 更多