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

Spring Flux<JsonNode>解析JSON数组阻塞等待全量加载问题

WebClient流式解析JSON数组问题解决方案

问题背景

调用疫情数据API时,接口返回标准JSON结构,外层包含分页元数据,data字段为目标记录数组:

{"length":850,"maxPageLimit":2500,"totalRecords":1700,
"data":[
{"date":"2022-06-29","newCasesByPublishDate":14476,"cumCasesByPublishDate":2005335},
{"date":"2022-06-26","newCasesByPublishDate":0,"cumCasesByPublishDate":1990859}
]}

接口响应头如下:

X-Firefox-Spdy  h2
cache-control   public, must-revalidate, max-age=90
content-encoding    gzip
content-location    https://api......&format=json&page=1
content-security-policy default-src 'none'; style-src 'self' 'unsafe-inline'
content-type    application/vnd.PHE-COVID19.v1+json; charset=utf-8

原有实现存在两个问题:

  • 必须等整个JSON响应完全加载、全量反序列化完成后,才会开始处理data数组内的记录
  • 数据量较大时触发Exceeded limit on max bytes per JSON object缓冲区超限错误
    尝试通过.bodyToFlux(JsonNode.class).flatMapIterable(jsonNode -> jsonNode.get("data"))改造,仍然无法实现边接收边解析的流式处理效果。
    原有问题代码:
public Flux<JsonNode> fetchCovidStatsFor(Area area, AreaType areaType, List<Metrics> metricsList) {
    var request =  generateWebClient().get()
            .uri(uriBuilder -> uriBuilder
                    .queryParam(buildRequestFilters(area, areaType))
                    .queryParam(buildRequestStructures(metricsList))
                    .build())
            .header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE);
    log.debug("request URI: {}", request.httpRequest(ClientHttpRequest::getURI));
    return request.retrieve()
            .bodyToFlux(JsonNode.class)
            .map(jsonNode -> jsonNode.get("data"))
            .doOnNext(jsonNode -> {System.out.println(jsonNode);});
}

核心原因

该表现不需要优先联系API方改造接口,本质是WebClient默认JSON反序列化逻辑的问题:

  • 默认情况下bodyToFlux(JsonNode.class)会将整个响应体作为单个JSON对象读入内存,等待完整解析完外层结构后才会向下游发射数据
  • 后续调用flatMapIterable拆分数组的操作,只是对已经加载到内存的全量数据做拆分,完全没有实现增量解析,因此既解决不了内存占用问题,也做不到边收边处理
  • 代码中GET请求设置Content-Type头属于无效配置,Content-Type是用来描述请求体格式的,GET请求无请求体时不需要设置,应该通过Accept头声明客户端支持的响应格式

可落地解决方案

使用Jackson流式解析器,直接基于到达的TCP数据块增量解析,跳过外层无关字段,逐个提取data数组内的单条记录,每解析完一条就立刻发射到Flux流中,不需要等待全量响应加载,也不需要增大JSON缓冲区配置。
完整实现代码:

public Flux<JsonNode> fetchCovidStatsFor(Area area, AreaType areaType, List<Metrics> metricsList) {
    var request =  generateWebClient().get()
            .uri(uriBuilder -> uriBuilder
                    .queryParam(buildRequestFilters(area, areaType))
                    .queryParam(buildRequestStructures(metricsList))
                    .build())
            .header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE);
    log.debug("request URI: {}", request.httpRequest(ClientHttpRequest::getURI));

    ObjectMapper objectMapper = new ObjectMapper();
    return request.retrieve()
            .bodyToFlux(DataBuffer.class)
            .flatMap(new Function<DataBuffer, Publisher<JsonNode>>() {
                private JsonParser jsonParser;

                @Override
                public Publisher<JsonNode> apply(DataBuffer dataBuffer) {
                    List<JsonNode> parsedRecords = new ArrayList<>();
                    try (InputStream inputStream = dataBuffer.asInputStream(true)) {
                        // 首次解析时初始化解析器,定位到data数组起始位置
                        if (jsonParser == null) {
                            jsonParser = objectMapper.getFactory().createParser(inputStream);
                            boolean locatedDataArray = false;
                            JsonToken currentToken;
                            while (!locatedDataArray && (currentToken = jsonParser.nextToken()) != null) {
                                if (currentToken == JsonToken.FIELD_NAME && "data".equals(jsonParser.currentName())) {
                                    // 跳过字段名,指向数组起始的[符号
                                    jsonParser.nextToken();
                                    locatedDataArray = true;
                                }
                            }
                        }
                        // 逐个解析当前数据块内完整的单条记录
                        while (jsonParser.nextToken() == JsonToken.START_OBJECT) {
                            parsedRecords.add(objectMapper.readTree(jsonParser));
                        }
                    } catch (IOException e) {
                        return Flux.error(e);
                    } finally {
                        DataBufferUtils.release(dataBuffer);
                    }
                    return Flux.fromIterable(parsedRecords);
                }
            });
}

实现说明

  • 解析器跨数据块复用上下文,不会因为TCP拆包导致JSON解析错误
  • 内存中仅同时持有当前处理的数据块和已解析的单条记录,不会加载全量响应,从根源避免缓冲区超限错误
  • 只要TCP层收到足够解析出一条完整记录的字节,就会立刻向下游发射数据,实现真正的流式处理

什么时候需要联系API方改造

如果运行上述代码后,发现仍然需要等待很长时间才收到第一条记录,说明服务端开启了全量响应缓存:即服务端会等待所有数据查询完成、整个JSON序列化完成后才开始往TCP连接写响应。这种情况下客户端任何配置都无法实现流式处理,才需要联系接口方做改造,比如改用NDJSON流格式、关闭服务端全量缓冲。
从当前给出的响应头无法直接判定服务端是否存在全量缓冲,建议先实现客户端流式解析后做实际测试验证。

内容的提问来源于stack exchange,提问作者octo-carrot

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 04:06:05