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

Reactive Java中如何基于条件逻辑返回响应并控制Blob写入?

问题描述

我有一个方法,其中fetchReport是对供应商API的外部调用。我希望将获取到的数据复制到Azure Blob存储中,但如果调用出现错误则不执行该操作;若发生错误,需返回包含错误详情的CustomResponse。writeToBlob()方法同样返回CustomResponse,我希望能保留外部API的错误信息并提供给消费者。

请问是否可以使用如下条件逻辑实现:

if response.contains("Failed") -> then return response with error details
else -> write to blob

当前我的方法代码如下:

public Flux<CustomResponse> getAndSaveReport(Mono<JsonNode> fetchReport, String reportFilePrefix) {
    Mono<JsonNode> reportMono = fetchReport
            .doOnSuccess(result -> {
                log.info(Logger.EVENT_UNSPECIFIED, "Successfully retrieved report");
            })
            .switchIfEmpty(Mono.just(objectMapper.convertValue(new CustomResponse("No content"), JsonNode.class)))
            .onErrorResume(BusinessException.class, err -> {
                log.error(Logger.EVENT_FAILURE, "Failed to retrieve report");
                JsonNode errJson = null;
                CustomResponse apiResponse = new CustomResponse();
                apiResponse.setStatus("Failed");
                apiResponse.setMessage("Error message: " + err.getMessage());
                apiResponse.setType(reportFilePrefix);
                errJson = objectMapper.convertValue(apiResponse, JsonNode.class);
                return Mono.just(errJson);
            });

    return writeToBlob(reportMono.flux(), reportFilePrefix).flux();
}
解决方案

你的思路完全可行,但当前代码没有实现条件判断逻辑——无论响应是否为错误状态,都会直接调用writeToBlob。以下是调整后的实现:

核心优化点

  1. 先将JsonNode转换为CustomResponse对象,避免反复序列化/反序列化,简化状态判断逻辑
  2. 通过flatMap实现分支控制:错误响应直接返回,正常响应才执行Blob写入操作

调整后的代码

public Flux<CustomResponse> getAndSaveReport(Mono<JsonNode> fetchReport, String reportFilePrefix) {
    return fetchReport
            // 将成功获取的JsonNode转为CustomResponse
            .map(result -> objectMapper.convertValue(result, CustomResponse.class))
            // 处理空响应场景
            .switchIfEmpty(Mono.just(new CustomResponse("No content")))
            // 捕获API调用异常,封装为错误响应
            .onErrorResume(BusinessException.class, err -> {
                log.error(Logger.EVENT_FAILURE, "Failed to retrieve report");
                CustomResponse errorResponse = new CustomResponse();
                errorResponse.setStatus("Failed");
                errorResponse.setMessage("Error message: " + err.getMessage());
                errorResponse.setType(reportFilePrefix);
                return Mono.just(errorResponse);
            })
            // 条件分支:判断响应状态决定是否写入Blob
            .flatMap(response -> {
                if ("Failed".equals(response.getStatus())) {
                    // 错误响应直接返回,不执行Blob写入
                    return Mono.just(response);
                } else {
                    // 正常响应传入writeToBlob,执行写入后返回结果
                    return writeToBlob(Flux.just(objectMapper.convertValue(response, JsonNode.class)), reportFilePrefix)
                            .map(blobResp -> objectMapper.convertValue(blobResp, CustomResponse.class));
                }
            })
            // 转换为Flux类型返回
            .flux();
}

补充说明

  • 如果writeToBlob方法本身接收Flux<CustomResponse>作为参数,可以省去JsonNode和CustomResponse之间的转换步骤,进一步简化代码
  • 错误响应的信息会完整保留并传递给消费者,完全符合你的需求

内容的提问来源于stack exchange,提问作者Gojo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 20:25:25