Kafka库升级:ListenableFuture转CompletableFuture重构疑问
问题描述
因Kafka库版本升级,需将原有基于ListenableFuture的代码改写为CompletableFuture实现,待处理对象为SendResult<String, Object>。
原有代码如下:
ListenableFuture<SendResult<String, Object>> future = ...; future.addCallback(new ListenableFutureCallback<SendResult<String, Object>>() { @Override public void onSuccess(final SendResult<String, Object> result) { ProducerRecord<String, Object> record = result.getProducerRecord(); CaseStatusRequest data = (CaseStatusRequest) record.value(); logger.info("Producing request succeeded: {}", data); } @Override public void onFailure(final Throwable throwable) { logger.error("Producing request failed: {}", request.getReceiptNumber()); } });
尝试改写的代码如下:
CompletableFuture<SendResult<String, Object>> future = ...; future.whenComplete(new BiConsumer<SendResult<String,Object>,Throwable>() { @Override public void accept(SendResult<String, Object> result, Throwable u) { ProducerRecord<String, Object> record = result.getProducerRecord(); CaseStatusRequest data = (CaseStatusRequest) record.value(); logger.info("Producing request succeeded: {}", data); } }); future.exceptionally(new Function<Throwable, SendResult<String,Object>>() { @Override public SendResult<String, Object> apply(Throwable arg0) { logger.error("Producing request failed: {}", request.getReceiptNumber()); // Something needs to be returned here. // Should I return NULL? } });
疑问:ListenableFuture的onFailure没有明确对应的无返回值方法,CompletableFuture的exceptionally必须返回值,但错误处理是void操作,该如何正确实现无返回值的失败回调处理?
正确实现方式
不需要分开调用whenComplete和exceptionally,直接在whenComplete中通过判断Throwable是否为null,即可分别处理成功和失败场景,完全对应原代码的onSuccess和onFailure逻辑:
CompletableFuture<SendResult<String, Object>> future = ...; future.whenComplete((result, throwable) -> { if (throwable == null) { // 对应原onSuccess逻辑 ProducerRecord<String, Object> record = result.getProducerRecord(); CaseStatusRequest data = (CaseStatusRequest) record.value(); logger.info("Producing request succeeded: {}", data); } else { // 对应原onFailure逻辑,无返回值需求直接处理即可 logger.error("Producing request failed: {}", request.getReceiptNumber(), throwable); } });
如果一定要拆分成功和失败逻辑,也可以用thenAccept处理成功(无返回值),搭配exceptionally处理失败,此时exceptionally只需返回符合类型的值即可(比如null,如果后续逻辑不依赖该返回值):
CompletableFuture<SendResult<String, Object>> future = ...; // 处理成功场景 future.thenAccept(result -> { ProducerRecord<String, Object> record = result.getProducerRecord(); CaseStatusRequest data = (CaseStatusRequest) record.value(); logger.info("Producing request succeeded: {}", data); }); // 处理失败场景 future.exceptionally(throwable -> { logger.error("Producing request failed: {}", request.getReceiptNumber(), throwable); return null; });
注意:错误日志建议传入throwable参数,便于后续排查问题。
内容的提问来源于stack exchange,提问作者gene b.
相关产品推荐
相关产品推荐

