【发布时间】:2021-02-04 10:21:46
【问题描述】:
我试图弄清楚如何在 事务上下文 中正确处理 原子多次写入 到 kafka。在这种情况下,事务不是由 kafka 消息侦听器发起,而是通过 @Transactional 注释以编程方式发起,参见下面的 sn-p。
我正在使用 spring-boot 2.4.2 和 spring-kafka 2.6.5。
KafkaProducer 文档说,在事务上下文中,不需要在返回的 Future 上调用 .get(),因为它最终会在尝试提交事务时抛出异常。此外,KafkaTemplate 在 KafkaProducer 返回的 Future 上调用 .get(),因此它看起来是同步的。
@PostMapping("/ingest/{topic}")
public ResponseEntity ingest(@PathVariable(value = "topic") String topic, @RequestBody String numbersString) {
List<String> numbers = Arrays.stream(numbersString.split(",")).collect(Collectors.toList());
kafkaWriterService.writeMany(numbers,topic);
return ResponseEntity.ok().build();
}
@Transactional
@Service
class KafkaWriterService {
@Autowired
KafkaTemplate<String, String> kafkaTemplate;
public void writeMany(List<String> messages, String topic) {
for (String message : messages) {
kafkaTemplate.send(message, topic, topic);
}
}
}
据我所知,以下 KafkaTemplate 方法
protected ListenableFuture<SendResult<K, V>> doSend(ProducerRecord<K, V> producerRecord)
等待Producer完成发送,然后返回另一个ListenableFutre,这个接口是异步的。
所以这是同步的,因为我们处于事务上下文中,还是我应该等待 kafkaTemplate 返回的所有 ListenableFutures 结束?我的意思是考虑到我需要以同步的方式回复调用者。
谢谢,问候
【问题讨论】:
标签: spring-kafka