【发布时间】:2017-05-18 14:14:35
【问题描述】:
我们目前基本上使用以下简化机制来确认消息:
@KafkaListener(topics = "someTopic")
public void listen(final String message, final Acknowledgment ack) {
try {
processMessage(message);
ack.acknowledge();
} catch (final IOException e) {
// do not acknowledge here since we can temporarily not process the message
}
基本上,只要我们暂时无法处理消息(在 IOExceptions 的情况下),我们希望稍后再接收它。
但这不起作用,因为确认假定同一分区中的所有先前消息都已成功处理。在我们的 IOException 案例中,失败的消息会被跳过,但可能会被同一分区上具有更高索引的不同消息确认。
我们有一些解决这个问题的想法,但这意味着需要一些讨厌的解决方法来避免在 KafkaListener 方法中直接调用确认。我们的用例是一个非常具体的用例,还是更像是 spring kafka 用户会假设的“默认”行为?
这种问题有spring-kafka解决方案吗?或者你有一个“正确”解决这个问题的想法吗?
【问题讨论】:
标签: spring apache-kafka spring-kafka