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

