【问题标题】:Spring Kafka Retry LoggingSpring Kafka 重试日志记录
【发布时间】:2019-01-09 16:20:27
【问题描述】:

我需要从一个 kafka 主题消费,对记录做一些工作并使用 spring-kafka 2.1.7 生成另一个主题。其他要求是事务性的,仅用于一次语义、重试和错误处理。在提交失败时一条记录我应该重试 3 次,记录每条重试消息以重试主题,然后在所有重试失败时将记录发送到死信主题。我查看了https://github.com/spring-projects/spring-kafka/issues/575,它在解决问题方面有很好的细节。我正在努力解决的问题是如何记录每条重试消息以及消费者偏移量、它试图提交的主题等详细信息。有没有办法从重试回调中获取这些信息?下面的 retrylistener sn-p 注册了一个 org.springframework.kafka.listener.LoggingErrorHandler ,它被设置为 ConcurrentKafkaListenerContainerFactory 的容器属性?

         @Bean
         public RetryListener retryListener(KafkaTemplate<String,SpecificRecord> kafkaTemplate) {
             return new RetryListenerSupport() {

                public void onError(RetryContext context, RetryCallback callback, Throwable throwable) {
                    int retryCount =context.getRetryCount();
                    kafkaTemplate) .send(new ProducerRecord<String,SpecificRecord>("topic_name",record));
                }
             };
         }

【问题讨论】:

    标签: spring-kafka spring-retry


    【解决方案1】:

    RetryContextRetryingMessageListenerAdapter 中填充了一些有用的信息:

    context.setAttribute(CONTEXT_RECORD, record);
    switch (RetryingMessageListenerAdapter.this.delegateType) {
        case ACKNOWLEDGING_CONSUMER_AWARE:
            context.setAttribute(CONTEXT_ACKNOWLEDGMENT, acknowledgment);
            context.setAttribute(CONTEXT_CONSUMER, consumer);
            RetryingMessageListenerAdapter.this.delegate.onMessage(record, acknowledgment, consumer);
            break;
        case ACKNOWLEDGING:
            context.setAttribute(CONTEXT_ACKNOWLEDGMENT, acknowledgment);
            RetryingMessageListenerAdapter.this.delegate.onMessage(record, acknowledgment);
            break;
        case CONSUMER_AWARE:
            context.setAttribute(CONTEXT_CONSUMER, consumer);
            RetryingMessageListenerAdapter.this.delegate.onMessage(record, consumer);
            break;
        case SIMPLE:
            RetryingMessageListenerAdapter.this.delegate.onMessage(record);
    }
    

    【讨论】:

    • 谢谢阿特姆。效果很好。我有一个案例,如果我从下游收到特定类型的异常,我应该向另一个主题发送带有不同消息的通知,并且不应该重试。在执行 producer.send() 和 producer.sendoffsetsToTransaction 之前可能会发生此异常,因此我猜这不会调用重试模板,因为此处不涉及事务管理器。是否会捕获该异常并在这种情况下进行确认工作?想法是在发生这种情况时处理来自消费者的下一条记录而不重试。
    • 没错。只要您不重新抛出异常,而是手动确认,则不涉及重试逻辑,您只需转到下一条记录进行处理。
    • 谢谢Artem。我对此进行了测试并且有效。我不确定这个问题是否属于这里,但它会抛出它。通过阅读 kafka 文档,可能会出现生产者收到无法恢复的异常的情况。使用 spring kafka 时在这些场景中会发生什么。是否创建了新的生产者?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-04-13
    • 2019-01-07
    • 2021-11-19
    • 2017-12-26
    • 2019-01-31
    • 2017-11-01
    相关资源
    最近更新 更多