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

