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

如何用WebClient从顶层JSON对象获取属性Flux并异步处理API响应

用WebClient流式处理先返回Metadata再返回Results的API

你完全可以实现这个需求!核心思路是利用Jackson流式JSON解析结合WebClient的响应流式处理,这样既能提前拿到Metadata,又能逐个处理Result而无需等待所有结果返回,最终得到你想要的Mono<Tuple2<Metadata, Flux<Result>>>。

步骤1:定义实体类

首先,对应API响应结构定义你的Metadata和Result实体:

public class Metadata {
    private String id;
    private int totalCount;
    // 生成getter、setter和构造方法
}

public class Result {
    private String itemId;
    private String content;
    // 生成getter、setter和构造方法
}

步骤2:实现流式解析的WebClient调用

关键是通过WebClient获取响应的流式数据,然后用Jackson的JsonParser逐个解析JSON节点,先提取Metadata,再流式发射Result:

import com.fasterxml.jackson.core.JsonFactory;
import com.fasterxml.jackson.core.JsonParser;
import com.fasterxml.jackson.core.JsonToken;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.web.reactive.function.client.WebClient;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.util.function.Tuple2;
import reactor.util.function.Tuples;

import java.io.InputStream;
import java.util.Objects;

public class ApiStreamingClient {
    private final WebClient webClient;
    private final ObjectMapper objectMapper;
    private final JsonFactory jsonFactory;

    public ApiStreamingClient(WebClient webClient, ObjectMapper objectMapper) {
        this.webClient = webClient;
        this.objectMapper = objectMapper;
        this.jsonFactory = objectMapper.getFactory();
    }

    public Mono<Tuple2<Metadata, Flux<Result>>> fetchStreamingData(String apiUrl) {
        return webClient.get()
                .uri(apiUrl)
                .exchangeToMono(clientResponse -> {
                    // 获取响应体的流式DataBuffer
                    Flux<DataBuffer> dataBufferFlux = clientResponse.bodyToFlux(DataBuffer.class);

                    // 将流式DataBuffer合并为InputStream(适合大多数场景,若响应极大可优化为非合并方式)
                    return DataBufferUtils.join(dataBufferFlux)
                            .map(dataBuffer -> {
                                try {
                                    return dataBuffer.asInputStream(true);
                                } finally {
                                    DataBufferUtils.release(dataBuffer);
                                }
                            })
                            .flatMap(inputStream -> {
                                try {
                                    JsonParser parser = jsonFactory.createParser(inputStream);

                                    // 跳过JSON开头的START_OBJECT标记
                                    parser.nextToken();

                                    // 解析Metadata部分
                                    Metadata metadata = null;
                                    while (parser.nextToken() != JsonToken.END_OBJECT) {
                                        String fieldName = parser.getCurrentName();
                                        if ("metadata".equals(fieldName)) {
                                            parser.nextToken(); // 进入metadata对象节点
                                            metadata = objectMapper.readValue(parser, Metadata.class);
                                        } else if ("results".equals(fieldName)) {
                                            parser.nextToken(); // 进入results数组节点
                                            break; // 找到results后跳出,开始解析数组元素
                                        } else {
                                            parser.skipChildren(); // 跳过其他无关字段
                                        }
                                    }

                                    Objects.requireNonNull(metadata, "API返回的metadata不能为空");

                                    // 流式解析results数组,生成Flux<Result>
                                    Flux<Result> resultsFlux = Flux.create(sink -> {
                                        try {
                                            // 逐个读取数组中的Result元素
                                            while (parser.nextToken() != JsonToken.END_ARRAY) {
                                                Result result = objectMapper.readValue(parser, Result.class);
                                                sink.next(result); // 立即发射Result,无需等待全部解析
                                            }
                                            sink.complete();
                                        } catch (Exception e) {
                                            sink.error(e);
                                        } finally {
                                            // 关闭资源
                                            try {
                                                parser.close();
                                                inputStream.close();
                                            } catch (Exception ignored) {}
                                        }
                                    });

                                    // 返回Metadata和Result流的组合
                                    return Mono.just(Tuples.of(metadata, resultsFlux));
                                } catch (Exception e) {
                                    return Mono.error(e);
                                }
                            });
                });
    }
}

步骤3:使用示例

调用这个方法后,你可以先处理Metadata,同时订阅Result流逐个处理结果,全程无阻塞:

public class Main {
    public static void main(String[] args) {
        // 初始化WebClient和ObjectMapper
        WebClient webClient = WebClient.create();
        ObjectMapper objectMapper = new ObjectMapper();
        ApiStreamingClient client = new ApiStreamingClient(webClient, objectMapper);

        client.fetchStreamingData("https://your-api-endpoint.com/data")
                .subscribe(tuple -> {
                    // 先处理Metadata
                    Metadata metadata = tuple.getT1();
                    System.out.printf("拿到Metadata:ID=%s,总结果数=%d%n", metadata.getId(), metadata.getTotalCount());

                    // 流式处理Result,每拿到一个就处理一个
                    Flux<Result> results = tuple.getT2();
                    results.subscribe(
                            result -> System.out.printf("处理Result:ID=%s,内容=%s%n", result.getItemId(), result.getContent()),
                            error -> System.err.println("处理Result出错:" + error.getMessage()),
                            () -> System.out.println("所有Result处理完成")
                    );
                }, error -> System.err.println("请求API出错:" + error.getMessage()));

        // 保持程序运行(非Web环境下需要)
        try {
            Thread.sleep(10000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

额外说明

  • 如果API采用分块传输编码(Chunked Transfer Encoding),这种方式可以真正做到边接收响应边解析,完全无需等待整个响应完成,性能最优。
  • 若响应体极大,合并DataBuffer可能占用过多内存,你可以优化为直接处理流式DataBuffer(通过DataBufferUtils.readInputStream结合Jackson的异步解析),但需要注意切换线程池避免阻塞Reactor的IO线程。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:21:17