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

Spring Kafka 3.1迁移:用CompletableFuture替代ListenableFuture兼容改造

Spring Kafka 3.1 迁移:ListenableFuture 替换为 CompletableFuture 的回调适配

Spring 6及Spring Kafka 3.1版本中,ListenableFuture已被弃用,KafkaTemplate.send()改为返回CompletableFuture<SendResult>。以下是两种适配回调逻辑的方案,完全兼容原有业务逻辑:

方案1:直接使用CompletableFuture原生API

利用CompletableFuture自带的thenAccept(处理成功结果)和exceptionally(处理异常)方法,替代原addCallback逻辑:

CompletableFuture<SendResult<String, Object>> future = kafkaTemplate.send(message);

// 成功回调:对应原onSuccess逻辑
future.thenAccept(result -> {
    log.info("Sent message=[{}] with offset=[{}]", payload, result.getRecordMetadata().offset());
});

// 失败回调:对应原onFailure逻辑
future.exceptionally(ex -> {
    log.info("Unable to send message=[{}] due to : {}", message, ex.getMessage());
    return null; // exceptionally需返回对应类型结果,此处返回null不影响后续同步逻辑
});

// 同步获取结果(优化原代码两次get()的冗余调用)
try {
    SendResult<String, Object> result = future.get();
    return result.getRecordMetadata().partition() + "-" + result.getRecordMetadata().offset();
} catch (InterruptedException | ExecutionException e) {
    Thread.currentThread().interrupt();
    throw new RuntimeException("Failed to send message", e);
}

方案2:封装工具类保留addCallback写法

如果想维持原有代码风格,可自定义工具类适配CompletableFuture:

// 自定义适配工具类
public class CompletableFutureCallbackAdapter {
    public static <T> void addCallback(CompletableFuture<T> future, ListenableFutureCallback<T> callback) {
        future.thenAccept(callback::onSuccess);
        future.exceptionally(ex -> {
            callback.onFailure(ex);
            return null;
        });
    }
}

// 业务代码中使用方式(和原代码几乎一致)
CompletableFuture<SendResult<String, Object>> future = kafkaTemplate.send(message);

CompletableFutureCallbackAdapter.addCallback(future, new ListenableFutureCallback<SendResult<String, Object>>() {
    @Override
    public void onSuccess(SendResult<String, Object> result) {
        log.info("Sent message=[{}] with offset=[{}]", payload, result.getRecordMetadata().offset());
    }

    @Override
    public void onFailure(Throwable ex) {
        log.info("Unable to send message=[{}] due to : {}", message, ex.getMessage());
    }
});

// 同步获取结果(同方案1)
try {
    SendResult<String, Object> result = future.get();
    return result.getRecordMetadata().partition() + "-" + result.getRecordMetadata().offset();
} catch (InterruptedException | ExecutionException e) {
    Thread.currentThread().interrupt();
    throw new RuntimeException("Failed to send message", e);
}

关键注意点

  • 原代码两次调用future.get()会重复阻塞线程,修改后建议一次获取结果复用,提升性能。
  • CompletableFuture的回调默认异步执行,和原ListenableFuture行为一致,若需指定线程池可通过thenAcceptAsync等方法实现。
  • 同步调用future.get()必须处理InterruptedException和ExecutionException,避免异常静默丢失。

内容的提问来源于stack exchange,提问作者cybertextron

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 19:22:25