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

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.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 14:37:40