【问题标题】:spring-kafka error handler invoked too many timesspring-kafka 错误处理程序调用了太多次
【发布时间】:2021-02-05 14:03:02
【问题描述】:

我正在使用 spring 的 SeekToCurrentErrorHandlerDeadLetterPublishingRecoverer
每次调用错误处理程序时都会记录一条消息,并且还会向我们的监控系统发送一个事件。
我看到的问题是错误处理程序被调用太多次,导致连续记录失败。
例如,对于不可重试的异常,如果我产生 2 个(错误)消息,我会得到以下日志:

|ERROR|Y|147e885e-9cce-4d10-972f-96613757c511|2020-10-22 13:55:41  [org.springframework.kafka.KafkaListenerEndpointContainer#0-1-C-1] Kafka:65|orgTest|projectTest|dev|input| my error log
|ERROR|Y|147e885e-9cce-4d10-972f-96613757c511|2020-10-22 13:55:41  [org.springframework.kafka.KafkaListenerEndpointContainer#0-1-C-1] Kafka:65|orgTest|projectTest|dev|input| <== my error log
|INFO|Y|147e885e-9cce-4d10-972f-96613757c511|2020-10-22 13:55:41  [org.springframework.kafka.KafkaListenerEndpointContainer#0-1-C-1] o.a.k.c.c.KafkaConsumer:1603|orgTest|projectTest|dev|input| [Consumer clientId=consumer-debug-2, groupId=debug] Seeking to offset 5 for partition debug-2
DEBUG|||2020-10-22 13:55:42  [kafka-producer-network-thread | producer-1] o.s.k.l.DeadLetterPublishingRecoverer:296||||| Successful dead-letter publication
|ERROR|Y|f0ecbdfa-3a13-4435-9cd1-b3daf73d324d|2020-10-22 13:55:42  [org.springframework.kafka.KafkaListenerEndpointContainer#0-1-C-1] Kafka:65|orgTest|projectTest|dev|input| my error logs
DEBUG|||2020-10-22 13:55:42  [kafka-producer-network-thread | producer-1] o.s.k.l.DeadLetterPublishingRecoverer:296||||| Successful dead-letter publication

似乎我得到了 x 个消费者错误消息:x 次调用错误处理程序,x-1 次调用错误处理程序,x-2 次调用,等等。

可重试异常也是如此,我每次重试都会看到相同的异常。 Consumer函数调用正确,只是错误处理触发次数过多。

这是我的错误处理配置:

public class CustomSeekToCurrentErrorHandler extends SeekToCurrentErrorHandler {
    private final Monitor monitor;

    CustomSeekToCurrentErrorHandler(Monitor monitor, DeadLetterPublishingRecoverer dlpr, FixedBackOff retries) {
        super(dlpr, retries);
        super.setLogLevel(KafkaException.Level.DEBUG);
        this.monitor = monitor;
    }

    @Override
    public void handle(Exception exception, List<ConsumerRecord<?, ?>> records, Consumer<?, ?> consumer, MessageListenerContainer container) {
        if (!records.isEmpty()) {
            records.forEach(record -> {
                logAndReportError(record);
            });
        }
        super.handle(exception, records, consumer, container);
    }
}

@Bean
public SeekToCurrentErrorHandler replayDeadLetterErrorHandler(DeadLetterPublishingRecoverer dlpr, FixedBackOff fxboff) {
    var seekToCurrent = new CustomSeekToCurrentErrorHandler(monitor, dlpr, fxboff);
   seekToCurrent.addNotRetryableException(SomeFatalException.class);
   return seekToCurrent;
}
    

我有两个问题:

  1. 为什么错误处理程序被触发了这么多次?
  2. 为什么要针对不可重试的异常执行查找?

【问题讨论】:

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


    【解决方案1】:
    1. 您的问题不清楚;请添加您的代码和配置。

    2. 我们必须在失败的记录之后寻找记录,以便在下一次投票时重新传递它们。

    【讨论】:

    • 感谢您的回复。我在原始问题中添加了我的错误处理程序配置。我可能对错误处理程序的操作方式有一些误解。假设消费者轮询 500 条消息,其中 10 条失败,每次失败都会导致搜索并立即重新轮询吗?当我发送 3 个消息时,它们都会导致致命错误(不可重试),我希望看到 3 个错误日志,相反,我看到 3 个错误日志,seek,2 个错误日志,seek,1 个错误日志,seek , 死信发布
    • 是的;这是意料之中的 - 这是因为您每次都在记录剩余的记录,以及失败的记录。当记录失败时,该记录(以及尚未处理的任何剩余记录)被传递给错误处理程序,以便可以执行查找。当重试次数用尽,或者异常不可重试时,我们不会寻找第一条记录(除非恢复失败)。您犯的错误是您正在记录所有记录,而不仅仅是失败的记录。但是,您应该会看到每条失败记录的发布,而不仅仅是最后一条。
    • 是的,我看到了每个人的发布,抱歉。有没有办法告诉哪个记录导致了失败?这是列表中的第一个吗?
    • 它始终是列表中的第一条记录。如果轮询 500 条记录并且第 10 条失败,您将获得 491 条记录(失败的记录在列表中的位置 0)。
    • 查看 RemainingRecordsErrorHandler 的 javadocs。
    猜你喜欢
    • 1970-01-01
    • 2023-04-10
    • 1970-01-01
    • 2020-07-17
    • 2016-11-08
    • 2013-07-16
    • 2015-05-24
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多