【问题标题】:Implementing DLQ in Kafka using Spring Cloud Stream with Batch mode enabled使用启用批处理模式的 Spring Cloud Stream 在 Kafka 中实现 DLQ
【发布时间】:2022-01-09 14:02:29
【问题描述】:

我正在尝试使用启用批处理模式的 Spring Cloud Stream 实现 DLQ

    @Bean
    public ListenerContainerCustomizer<AbstractMessageListenerContainer<?, ?>> customizer(BatchErrorHandler handler) {
        return ((container, destinationName, group) -> {
              if(dlqEnabledTopic.contains(destinationName))
                                    container.setBatchErrorHandler(handler);});
    }

    @Bean
    public BatchErrorHandler batchErrorHandler(KafkaOperations<String, byte[]> kafkaOperations) {
        CustomDeadLetterPublishingRecoverer recoverer = new CustomDeadLetterPublishingRecoverer(kafkaOperations,
                (cr, e) -> new TopicPartition(cr.topic()+"_dlq", cr.partition()));
        return new RecoveringBatchErrorHandler(recoverer, new FixedBackOff(1000, 1));
    }

但有一些疑问:

  1. 如何使用属性配置键/值序列化器 - 我的消息是字符串类型,但 KafkaOperations 使用的是 ByteArraySerializer

  2. 在批处理中有多条消息,但如果第一条消息失败,它会转到 DLQ 但看不到下一条消息的处理。

要求 - 在任何索引处,如果批处理失败,我只需将该消息发送到 DLQ,其余消息应再次处理。

  1. 现在批处理模式支持 DLQ 吗?就像记录模式一样,它可以使用属性启用

【问题讨论】:

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


    【解决方案1】:
    1. spring.kafka.producer.* 属性 - 但是,DLT 发布应该使用与主流应用程序相同的序列化程序。 ByteArraySerializer 通常是正确的。

    2. 正在恢复的批处理错误处理程序将对未处理的记录执行查找并将它们返回。调试日志记录应该可以帮助您找出问题所在。如果您无法弄清楚,请提供一个 MCRE 来展示您所看到的行为。

    3. 没有; binder 不支持批处理模式的 DLQ;配置错误处理程序是正确的方法。

    【讨论】:

    • 感谢@Gary,我必须在将消息推送到 dlq 之前对其进行更新 - 因此手动转换为 byte[] 并且它起作用了。我们需要重写 DeadLetterPublishingRecoverer,它可以是 spring 上下文中的 @Component 吗? @Component public class CustomDeadLetterPublishingRecoverer extends DeadLetterPublishingRecoverer {..} 另外,在调试期间 - 消费者再次轮询消息需要一些时间。对于这些场景,建议的方法是什么?
    • 有没有办法在批处理模式下从重试中跳过特定类型的异常?
    • 是的,请参阅 addNotRetryableExceptionssetClassifications 方法。关于您的第一条评论,您需要提供更多信息;我看不到如何“再次轮询”恢复的消息。
    • Boot 自动配置的KafkaOperations 对活页夹属性一无所知。使用spring.kafka.bootstrap-servers(如果没有...kafka.binder.brokers 属性,它也将被活页夹使用)。 zkNodes 已经好几年不用了;客户端不再连接到 Zookeeper。
    • 这没有任何意义;对不起。它的分区数是否与原始分区数相同?您的目标解析器正在路由到与原始分区相同的分区。有关目标解析器以及如何选择目标分区的信息,请参阅docs.spring.io/spring-kafka/docs/2.7.x/reference/html/…
    猜你喜欢
    • 2021-10-08
    • 2018-12-17
    • 1970-01-01
    • 1970-01-01
    • 2018-12-29
    • 1970-01-01
    • 2019-04-17
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多