【问题标题】:Is KafkaTemplate transactional send sync or async?KafkaTemplate 是事务性发送同步还是异步?
【发布时间】: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


    【解决方案1】:

    没有;如果发送立即完成未来,它只会调用get()...

            Future<RecordMetadata> sendFuture =
                    producer.send(producerRecord, buildCallback(producerRecord, producer, future, sample));
            // May be an immediate failure
            if (sendFuture.isDone()) {
                try {
                    sendFuture.get();
                }
    ...
    

    在实际执行发送之前,客户端会发生某些错误。

    https://github.com/spring-projects/spring-kafka/issues/1437

    【讨论】:

    • 你能想出一个客户端不会发生错误的可测试场景吗?我想测试一个实际的异步错误
    猜你喜欢
    • 2019-01-17
    • 2010-09-10
    • 2017-06-22
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多