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

