【问题标题】:Consumer is consuming same message multiple time消费者多次消费同一条消息
【发布时间】:2019-04-03 03:59:36
【问题描述】:

我有一个消费者消息的 Kafka 消费者。我在 KafkaListenerContainerFactory 中设置了重试模板。但不知道为什么它两次消费相同的消息。当应用程序通过指定计数调用异常重试模板时,Kafka 也会使用相同计数的相同消息。

@Bean
RetryTemplate retryTemplate() {
    RetryTemplate retryTemplate = new RetryTemplate();

    FixedBackOffPolicy fixedBackOffPolicy = new FixedBackOffPolicy();
    fixedBackOffPolicy.setBackOffPeriod(1000L);
    retryTemplate.setBackOffPolicy(fixedBackOffPolicy);

    SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy();
    retryPolicy.setMaxAttempts(3);
    retryTemplate.setRetryPolicy(retryPolicy);

    return retryTemplate;
}


@Bean("KafkaListenerContainerFactory")
public ConcurrentKafkaListenerContainerFactory<String, String> KafkaListenerContainerFactory(RetryTemplate retryTemplate) {
    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setRetryTemplate(retryTemplate);
    factory.setRecoveryCallback(context -> {
        log.error("Maximum retry policy has been reached");
        return null;
    });
    factory.setConcurrency(Integer.parseInt(kafkaConcurrency));
    return factory;
}

消费者

@KafkaListener(topics = "${kafka.topic.json}", containerFactory = "kafkaListenerContainerFactory")
public void recieveSegmentService(String KafkaPayload) throws Exception {
    KafkaSegmentTrigger kafkaSegmentTrigger;
    kafkaSegmentTrigger = TransformUtil.fromJson(KafkaPayload, KafkaSegmentTrigger.class);
    log.info("Trigger recieved from segment service {}", kafkaSegmentTrigger);
    try {
        processMessage(kafkaSegmentTrigger);
    } catch (Exception e) {
        retryTemplate.execute(arg0 -> {
            processMessage(kafkaSegmentTrigger);
            return null;
        });
    }finally {
    }
}

processMessage 正在抛出异常

【问题讨论】:

    标签: spring-boot kafka-consumer-api spring-retry


    【解决方案1】:

    您已经嵌套了RetryTemplates - 一个在侦听器外部,一个在内部。如果您在两个地方都使用相同的模板,您将获得 12 次尝试(3 次由侦听器内部的 x4 侦听器适配器驱动)。

    使用其中一个。

    【讨论】:

    • 实际上是12次尝试。
    • 谢谢先生..那么我应该删除 Kafka 消费者工厂中的设置重试模板还是删除监听器中的 retryTemplate.execute 方法
    • 如果我删除..retryTemplate.execute 以及我将如何确认偏移量..如果我最后添加它,那么它将每次都执行
    • 监听容器会在重试次数用尽时提交偏移量。
    • @Gray 感谢您的帮助。在侦听器中,我没有手动确认...它是否在内部执行此提交
    猜你喜欢
    • 1970-01-01
    • 2017-10-19
    • 2015-11-10
    • 2017-02-05
    • 2016-04-12
    • 1970-01-01
    • 2021-05-14
    • 2018-06-04
    • 2016-06-08
    相关资源
    最近更新 更多