WebClient分页响应异步处理:内存问题与非阻塞优化诉求
解决方案
修改核心逻辑实现流式处理
将原来收集所有数据再处理的逻辑,改为流式输出每个数据项,实现边请求下一页边处理当前数据的效果。修改后的fetchData方法如下:
public Flux<JsonNode> fetchData() { return webClientHelper.fetchItems(null) // 递归请求下一页,当前页数据输出后立即发起下一页请求 .expand(response -> { List<String> headerValue = response.getHeaders().get(Constants.HEADER_ITEM_CURSOR); if (headerValue == null || headerValue.isEmpty()) { return Mono.empty(); } return webClientHelper.fetchItems(headerValue.get(0)); }) // 将每页的JsonNode数组拆分为单个元素流式输出 .flatMap(response -> Flux.fromArray(response.getBody())) // 对接数据处理逻辑,每个数据项一可用就立即处理 .flatMap(this::processItem); } // 自定义数据处理方法,根据业务需求实现 private Mono<Void> processItem(JsonNode item) { // 示例:同步/异步处理逻辑,比如持久化、消息发送等 return Mono.fromRunnable(() -> { System.out.println("Processing item: " + item.toPrettyString()); }); }
关键改动说明
- 返回类型调整:从
Mono<List<JsonNode>>改为Flux<JsonNode>,不再一次性缓存所有数据,而是逐个输出数据项,彻底解决内存占用问题。 - 移除
collectList():取消对所有数据的收集操作,让每页数据解析完成后立即向下游传递处理。 expand的异步特性:expand操作符会在当前页的响应被输出后,立即发起下一页的请求,实现请求和处理的并行化,不会阻塞下一页请求。- 数据项级处理:通过
flatMap将每页的数组拆分为单个元素,确保每个数据项一可用就进入处理流程。
可选优化
- 并发控制:如果处理逻辑耗时,可以在最后的
flatMap中指定并发数,比如.flatMap(this::processItem, 3),限制同时处理的数据项数量,避免系统过载。 - 错误处理:可以在处理逻辑中添加异常捕获,比如
.flatMap(item -> processItem(item).onErrorResume(e -> { /* 异常处理逻辑 */ })),避免单个数据项处理失败中断整个流程。
内容的提问来源于stack exchange,提问作者Sathesh S
相关产品推荐
相关产品推荐

