Spring KafkaTemplate.send()回调注册滞后未执行的可能性及解决办法
关于KafkaTemplate.send返回的ListenableFuture注册回调的执行问题
你完全不用担心这种情况下回调不会被执行的问题!
Spring的ListenableFuture接口从设计之初就考虑到了「结果已完成后再注册回调」的场景,它的所有标准实现(比如AsyncResult、CompletableToListenableFutureAdapter,而KafkaTemplate返回的Future正是这类实现)在执行addCallback方法时,都会先检查当前Future是否已经处于完成状态:
- 如果Future已经成功完成(调用过
set),注册后会立即触发onSuccess回调 - 如果Future已经失败完成(调用过
setException),注册后会立即触发onFailure回调
举个底层逻辑的例子,以AsyncResult的addCallback实现来说,它会先判断isDone(),如果为true,就直接同步调用对应的回调方法(如果有指定任务执行器,也可能会提交到执行器异步执行,但一定会执行),不会因为回调注册晚了就跳过。
额外的保障(如果需要更严谨的处理)
如果你还是想做一层额外的兜底(虽然完全没必要),可以先手动判断Future的状态再注册回调,比如:
ListenableFuture<SendResult<String, Data>> future = kafkaTemplate.send(topicname, keyString, data); if (future.isDone()) { try { SendResult<String, Data> result = future.get(); logger.info("Successfully sent message to kafka"); } catch (Exception ex) { logger.error("Failure while sending message in kafka.", ex); } } else { future.addCallback(new ListenableFutureCallback<SendResult<String, Data>>() { @Override public void onFailure(Throwable ex) { logger.error("Failure while sending message in kafka.", ex); } @Override public void onSuccess(SendResult<String, Data> result) { logger.info("Successfully sent message to kafka"); } }); }
不过这属于冗余代码,因为Spring的addCallback已经内置了这个逻辑,直接注册回调就足够安全。
总结一下:不管send()方法的结果是在注册回调前还是后完成,你的回调方法一定会被执行,完全不用顾虑这个极端场景。
内容的提问来源于stack exchange,提问作者cyrilantony
相关产品推荐
相关产品推荐

