Spring Boot 3迁移Kafka客户端遇CompletableFuture.addCallback解析错误
Spring Boot 3 下 KafkaProducer 代码迁移方案
问题根源
Spring Boot 3 集成的 Spring Kafka 3.x 版本中,KafkaTemplate.send() 方法的返回类型从 Spring 自定义的 ListenableFuture 改为了 JDK 原生的 CompletableFuture,这直接导致了类型不匹配和 addCallback 方法找不到的错误——因为 CompletableFuture 本身没有 addCallback 方法。
迁移方案
方案1:改用 CompletableFuture 原生回调(推荐)
直接适配 Java 原生的异步回调机制,替换原有 ListenableFutureCallback 参数为 BiConsumer,代码如下:
import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.SendResult; import org.springframework.stereotype.Service; import java.util.Objects; import java.util.function.BiConsumer; @Service public class KafkaProducer<K, V> { private final KafkaTemplate<K, V> kafkaTemplate; public KafkaProducer(KafkaTemplate<K, V> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void send(String topic, K key, V message, BiConsumer<SendResult<K, V>, Throwable> callback) { CompletableFuture<SendResult<K, V>> future = kafkaTemplate.send(topic, key, message); if (Objects.nonNull(callback)) { // 使用CompletableFuture原生的whenComplete方法处理回调 future.whenComplete(callback); } } }
调用方可以通过 Lambda 表达式传入回调逻辑,例如:
kafkaProducer.send("test-topic", "key", "message", (result, ex) -> { if (ex != null) { // 处理发送失败 } else { // 处理发送成功 } });
方案2:兼容原有 ListenableFutureCallback 接口
如果需要保留原有代码的回调接口,无需修改调用方,可以使用 Spring 提供的 CompletableFutureUtils 工具类完成适配:
import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.SendResult; import org.springframework.stereotype.Service; import org.springframework.util.concurrent.CompletableFutureUtils; import org.springframework.util.concurrent.ListenableFutureCallback; import java.util.Objects; @Service public class KafkaProducer<K, V> { private final KafkaTemplate<K, V> kafkaTemplate; public KafkaProducer(KafkaTemplate<K, V> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void send(String topic, K key, V message, ListenableFutureCallback<SendResult<K, V>> callback) { CompletableFuture<SendResult<K, V>> future = kafkaTemplate.send(topic, key, message); if (Objects.nonNull(callback)) { // 通过工具类将ListenableFutureCallback适配到CompletableFuture CompletableFutureUtils.addCallback(future, callback); } } }
注意:这里给 ListenableFutureCallback 添加了泛型约束 <SendResult<K, V>>,避免类型安全警告,同时保持原有回调逻辑的写法不变。
补充说明
Spring Kafka 3.x 改用 CompletableFuture 是为了对齐 JDK 原生异步 API,减少对 Spring 自定义异步组件的依赖,让代码更贴合 Java 生态标准。如果是新项目或可以修改调用方代码,优先选择方案1;如果需要兼容大量原有代码,方案2是更平滑的过渡方式。
内容的提问来源于stack exchange,提问作者Peter Penzov
相关产品推荐
相关产品推荐

