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

Spring Boot 3 Kafka代码迁移遇类型错误,求正确实现方案

迁移Spring Boot 3 Kafka代码:ListenableFuture转CompletableFuture问题

原有ListenableFuture实现逻辑

之前基于ListenableFuture的消息发送回调代码如下:

public void publishMessage(String message) {
    ListenableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
    future.addCallback(new ListenableFutureCallback<SendResult<String, String>>() {
        @Override
        public void onSuccess(SendResult<String, String> result) {
            // 消息发送成功逻辑
        }
        @Override
        public void onFailure(Throwable ex) {
            // 消息发送失败逻辑
        }
    });
}

错误的CompletableFuture尝试及报错

你尝试改用CompletableFuture时写了这段代码:

public void publishMessage(String message) {
    CompletableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
    future.whenComplete(new CompletableFuture<SendResult<String, String>>() {
        @Override
        public void accept(SendResult<String, String> kvSendResult, Throwable throwable) {
            if (null != throwable) {
                // 原onFailure逻辑
            } else {
                // 原onSuccess逻辑
            }
        }
    });
}

触发了类型不匹配的错误:

Required type: BiConsumer <? super org.springframework.kafka.support.SendResult<java.lang.String,java.lang.String>,
? super java.lang.Throwable>
Provided: anonymous CompletableFuture<SendResult<String, String>>

问题解决及逻辑验证

正确的匿名内部类写法

报错的核心原因是whenComplete方法需要的参数是**BiConsumer<? super T, ? super Throwable>**类型,而非CompletableFuture。你需要实现BiConsumer接口的accept方法,正确写法如下:

public void publishMessage(String message) {
    CompletableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
    future.whenComplete(new BiConsumer<SendResult<String, String>, Throwable>() {
        @Override
        public void accept(SendResult<String, String> result, Throwable throwable) {
            if (throwable != null) {
                // 原onFailure逻辑
            } else {
                // 原onSuccess逻辑
            }
        }
    });
}

Lambda版本的正确性

你后续尝试的Lambda表达式实现完全符合原有逻辑:

public void publishMessage(String message) {
    CompletableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
    future.whenComplete((result, throwable) -> {
        if (throwable != null) {
            // 原onFailure逻辑
        } else {
            // 原onSuccess逻辑
        }
    });
}

Lambda在这里是BiConsumer的简写形式,和匿名内部类实现的效果完全一致:当throwable不为null时执行原失败逻辑,为null时执行原成功逻辑,完美对应原来ListenableFuture的回调逻辑。

内容的提问来源于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.26 00:44:50