【发布时间】: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