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
相关产品推荐
相关产品推荐

