【问题标题】:How to publish Spring Kafka DLQ in 2.5.4 version如何在 2.5.4 版本中发布 Spring Kafka DLQ
【发布时间】:2020-08-11 12:16:45
【问题描述】:

在这方面需要您的帮助和指导。

我在当前项目中使用的是 2.2.X 版本的 spring-kafka。

我创建的错误处理如下所示:

@Bean("kafkaConsumer")
public ConcurrentKafkaListenerContainerFactory<String, Map<String, Object>> eventKafkaConsumer() {
    ConcurrentKafkaListenerContainerFactory<String, Map<String, Object>> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setErrorHandler(new SeekToCurrentErrorHandler(createDeadLetterPublishingRecoverer(), 3));
    return factory;
}

public DeadLetterPublishingRecoverer createDeadLetterPublishingRecoverer() {
    return new DeadLetterPublishingRecoverer(getEventKafkaTemplate(),
            (record, ex) -> new TopicPartition("topic-undelivered", -1));
}

然后我将所有项目依赖版本,例如spring-boot和spring-kafka升级到最新版本:2.5.4 RELEASE

我发现有些方法已被弃用和更改。

SeekToCurrentErrorHandler

SeekToCurrentErrorHandler errorHandler =
new SeekToCurrentErrorHandler((record, exception) -> {
    // recover after 3 failures, woth no back off - e.g. send to a dead-letter topic
}, new FixedBackOff(0L, 2L));

我的问题是, 如何使用这些配置生成 DLQ:

已编辑

@Bean("kafkaConsumer")
public ConcurrentKafkaListenerContainerFactory<String, Map<String, Object>> kafkaConsumer() {
    ConcurrentKafkaListenerContainerFactory<String, Map<String, Object>> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setConcurrency(consumerConcurrencyCount);
    factory.setErrorHandler(errorHandler());
    return factory;
}

public SeekToCurrentErrorHandler errorHandler() {
    return new SeekToCurrentErrorHandler(
            deadLetterPublishingRecoverer(),
            new FixedBackOff(0L, 2L)
    );
}

public DeadLetterPublishingRecoverer deadLetterPublishingRecoverer() {
    return new DeadLetterPublishingRecoverer(
            getEventKafkaTemplate(),
            (record, ex) -> {
                if (ex.getCause() instanceof BusinessException || ex.getCause() instanceof TechnicalException) {
                    return new TopicPartition("topic-undelivered", -1);
                }

                return new TopicPartition("topic-fail", -1);
            });
}

public KafkaOperations<String, Object> getEventKafkaTemplate() { // producer to DLQ
    return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(producerConfigs()));
}

感谢 Gary,此配置有效!

提前致谢

【问题讨论】:

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


    【解决方案1】:

    不清楚你的意思

    问题是,在文档中,它仍在使用旧方法,该方法已被 2.5.X 版本弃用

    KafkaOperationsKafkaTemplate 实现的接口;您需要做的唯一更改是将maxAttempts 更改为BackOff...

    @Bean("kafkaConsumer")
    public ConcurrentKafkaListenerContainerFactory<String, Map<String, Object>> eventKafkaConsumer() {
        ConcurrentKafkaListenerContainerFactory<String, Map<String, Object>> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setErrorHandler(new SeekToCurrentErrorHandler(createDeadLetterPublishingRecoverer(), new FixedBackOff(0, 2L));
        return factory;
    }
    
    public DeadLetterPublishingRecoverer createDeadLetterPublishingRecoverer() {
        return new DeadLetterPublishingRecoverer(getEventKafkaTemplate(),
                (record, ex) -> new TopicPartition("topic-undelivered", -1));
    }
    

    【讨论】:

    • 嗨,加里,感谢您的回答。我确实只在 BackOff 上进行了更改,但 KafkaOperation 是错误的。 public KafkaTemplate getKafkaTemplate() { return new KafkaTemplate(new DefaultKafkaProducerFactory(producerConfigs()));并且 producerConfigs() 是 Map 类型。当我尝试 getKafkaTemplate.send("topic", "msg") 时,出现错误。
    • 没关系;您可以忽略弃用警告或将方法更改为KafkaOperations&lt;Object, Object&gt; getKafkaTemplate() there was an error - 什么错误?如果它不仅仅是弃用警告,请编辑问题以显示您的完整代码和错误;不要尝试将此类内容放在评论中,因为您会发现它很难阅读。
    • 嗨,加里,是的,当然!我已经编辑了这个问题。很抱歉给您带来不便,因为我是使用这个平台的新手。你能检查一下编辑过的版本并帮助解决问题吗?谢谢
    • 这是错误的:new SeekToCurrentErrorHandler((record, exception) -&gt; createDeadLetterPublishingRecoverer(), ...;它应该是new SeekToCurrentErrorHandler(createDeadLetterPublishingRecoverer(), ...,这是错误的new DeadLetterPublishingRecoverer( getEventKafkaTemplate().send("topic", null), ... 它应该是new DeadLetterPublishingRecoverer(getEventKafkaTemplate(), ...。还将getEventKafkaTemplate() 的返回类型更改为KafkaOperations 以消除弃用警告。你有没有看我的回答?
    • 是的 Gary,我正在测试它以针对特定异常创建 dlq。顺便说一句,现在一切正常,我只需要在 KafkaOperations 中更改为 FixedBackOff。非常感谢@Gary
    猜你喜欢
    • 2019-08-13
    • 1970-01-01
    • 2018-12-17
    • 1970-01-01
    • 2022-07-06
    • 2018-12-29
    • 1970-01-01
    • 1970-01-01
    • 2017-06-26
    相关资源
    最近更新 更多