【发布时间】:2018-05-09 01:42:04
【问题描述】:
我正在尝试使用 spring-kafka 1.3.x(1.3.3 和 1.3.4)。尚不清楚的是,当发生异常(例如网络中断)时,是否有一种安全的方式来批量消费消息而不会跳过一条消息(或一组消息)。我的偏好也是尽可能地利用容器功能以保留在 Spring 框架中,而不是尝试创建自定义框架来应对这一挑战。
我将以下属性设置到 ConcurrentMessageListenerContainer :
.setAckOnError(false);
.setAckMode(AckMode.MANUAL);
我还设置了以下 kafka 特定的消费者属性:
enable.auto.commit=false
auto.offset.reset=earliest
如果我设置了 RetryTemplate,我会得到一个类转换异常,因为它只适用于非批处理消费者。文档状态重试不适用于批处理,因此这可能没问题。
然后我设置了一个消费者,比如这个: ```java
@KafkaListener(containerFactory = "conatinerFactory",
groupId = "myGroup",
topics = "myTopic")
public void onMessage(@Payload List<Entries> batchedData,
@Header(required = false,
value = KafkaHeaders.OFFSET) List<Long> offsets,
Acknowledgment ack) {
log.info("Working on: {}" + offsets);
int x = 1;
if(x == 1) {
log.info("Failure on: {}" + offsets);
throw new RuntimeException("mock failure");
}
// do nothing else for now
// unreachable code
ack.acknowledge();
}
```
当我向系统发送消息以模拟上述异常时,对我来说唯一可见的操作是侦听器报告异常。
当我向系统发送另一条(新)消息时,容器会使用新消息。由于偏移量前进到下一个偏移量,旧消息被跳过。
由于我已经要求容器不要确认(直接或间接),并且由于我看不到其他属性来通知容器不要前进,所以我很困惑为什么容器会前进。
我注意到,出于类似考虑,建议升级到 2.1.x 并使用那里添加到 ContainerAware ErrorHandler 中的容器停止功能。
但是如果你暂时被困在 1.3.x 中,有没有办法或者缺少的属性可以用来确保容器不会前进到下一条消息或一批消息?
我可以看到一个选项,可以围绕消费者创建自定义框架,以达到预期的效果。但是还有其他选择吗,更简单,更弹簧友好。
想法?
【问题讨论】:
标签: java spring apache-kafka