【问题标题】:How to use spring-kafka for sending a message again如何使用 spring-kafka 再次发送消息
【发布时间】:2018-12-07 14:22:59
【问题描述】:

我们正在使用 spring-kafka 1.2.2.RELEASE。

我们想要什么
1、一旦消息被消费处理成功,就会在spring-kafka中提交offset。
我正在使用 Manaul Commit/Acknowledgement,它工作正常。

2. 如果出现任何异常,我们希望 spring-kafka 重新发送相同的消息。
我们在任何系统错误上抛出 RunTime 异常,该错误由 spring-kafka 记录并且从未提交。
这很好,因为我们不希望它提交,但是该消息保留在 spring-kafka 中并且永远不会回来,除非我们重新启动服务。重新启动消息返回并再次执行,然后留在 spring-kafka

我们尝试了什么
1. ErrorHandler 和 RetryingMessageListenerAdapter 我都试过了,但在这两种情况下,我们都必须在服务中编写代码如何再次处理消息

这是我的消费者

public class MyConsumer{
    @KafkaListener
    public void receive(...){
        // application logic to return success/failure
        if(success){
            acknowledgement.acknowledge();
        }else{
            throw new RunTimeException();
        }
    }
} 

我还有以下容器工厂的配置

factory.getContainerProperties().setErrorHandler(new ErrorHandler(){
    @Override
    public void handle(...){
        throw new RunTimeException("");
    }
});

在执行流程时,控制首先进入内部接收然后处理方法。在该服务等待新消息之后。但是我期待,因为我们抛出了一个异常,并且消息没有提交,相同的消息应该再次进入接收方法。

有什么办法,我们可以告诉spring kafka“不要提交这条消息,并尽快再次发送?”

【问题讨论】:

    标签: spring-kafka


    【解决方案1】:

    1.2.x 不再受支持;建议 1.x 用户至少升级到 1.3.x(当前为 1.3.8),因为它的线程模型更简单,这要归功于 KIP-62。

    当前版本是 2.2.2。

    2.0.1 引入了SeekToCurrentErrorHandler,它重新寻找失败的记录,以便重新传递。

    对于早期版本,您必须停止并重新启动容器才能重新传递失败的消息,或者向侦听器适配器添加重试。

    我建议您升级到可能的最新版本。

    【讨论】:

    • 感谢 Gary,我们正在升级到 2.2.0。我会尝试建议并更新结果。
    • 不幸的是,我们可以使用的版本是 1.3.7.RELEASE。我提到了这个问题的答案 - stackoverflow.com/questions/39536012/…,用于使用 CustomerSeekAware。请参阅带有输出的已编辑问题。
    • 同时,我对 RetryingMessageListenerAdapter 的理解是,我仍然必须编写机制来重试消息。然而,我们正在寻找的是消息由 spring-kafka 重新传递。那是对的吗?如果不是这样,请给我举个例子来了解它是如何工作的。
    • 能够让它工作,在下面回答,谢谢@Gary
    【解决方案2】:

    不幸的是,我们可以使用的版本是 1.3.7.RELEASE。

    我已经尝试实现 ConsumerSeekAware 接口。以下是我的做法,我可以看到消息重复发送

    消费者

    public class MyConsumer implements ConsumerSeekAware{
        private ConsumerSeekCallback consumerSeekCallback;
        if(condition) {
                acknowledgement.acknowledge();
            }else {
                consumerSeekCallback.seek((String)headers.get("kafka_receivedTopic"),
                        (int) headers.get("kafka_receivedPartitionId"),
                        (int) headers.get("kafka_offset"));
            }
        }
    
        @Override
        public void registerSeekCallback(ConsumerSeekCallback consumerSeekCallback) {
            this.consumerSeekCallback = consumerSeekCallback;
        }
    
        @Override
        public void onIdleContainer(Map<TopicPartition, Long> arg0, ConsumerSeekCallback arg1) {
            LOGGER.debug("onIdleContainer called");
        }
    
        @Override
        public void onPartitionsAssigned(Map<TopicPartition, Long> arg0, ConsumerSeekCallback arg1) {
            LOGGER.debug("onPartitionsAssigned called");
        }
    }
    

    配置

    public class MyConsumerConfig {
    
        @Bean
        public Map<String, Object> consumerConfigs() {
            Map<String, Object> props = new HashMap<>();
            // Set server, deserializer, group id
            props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
            props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
            return props;
        }
    
        @Bean
        public ConcurrentKafkaListenerContainerFactory<String, MyModel> kafkaListenerContainerFactory() {
            ConcurrentKafkaListenerContainerFactory<String, MyModel> factory = new ConcurrentKafkaListenerContainerFactory<>();
            factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(consumerConfigs()));
            factory.getContainerProperties().setAckMode(AckMode.MANUAL);
            return factory;
        }
    
        @Bean
        public MyConsumer receiver() {
            return new MyConsumer();
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2021-06-12
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-04-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多