【发布时间】:2017-08-06 17:38:40
【问题描述】:
我使用了 spring io 文档中列出的示例配置,它工作正常。
<int-kafka:message-driven-channel-adapter
id="kafkaListener"
listener-container="container1"
auto-startup="false"
phase="100"
send-timeout="5000"
channel="nullChannel"
message-converter="messageConverter"
error-channel="errorChannel" />
但是,当我使用下游应用程序测试它时,我从 kafka 消费并将其发布到下游。如果下游已关闭,则消息仍在被消耗且未重播。
或者说从 kafka 主题消费后,如果我在服务激活器中发现一些异常,我也想抛出一些异常,该异常应该回滚事务,以便可以重播 kafka 消息。
简而言之,如果消费应用程序有问题,那么我想回滚事务,以便消息不会被自动确认并一次又一次地重播,除非它被成功处理。
【问题讨论】:
标签: spring apache-kafka spring-integration