You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 23:31:22