如何使用Jackson反序列化由List<CustomObject>组成的Flux流
核心实现方案
你需要的是按流逐个解析JSON顶级数组中的每个List<CustomObject>元素,不需要等待全量数据返回,直接按以下步骤配置即可:
1. 自定义WebClient解码器配置
调整Jackson解码器的内存限制与流式解析策略,确保每个完整的数组元素解析完成后立刻下发:
import org.springframework.http.codec.json.Jackson2JsonDecoder; import org.springframework.web.reactive.function.client.WebClient; import org.springframework.core.ParameterizedTypeReference; // 初始化自定义配置的Jackson解码器 Jackson2JsonDecoder jsonDecoder = new Jackson2JsonDecoder(); // 根据单个序列的最大大小调整内存阈值,此处设置为10MB,可按需调整 jsonDecoder.setMaxInMemorySize(10 * 1024 * 1024); // 构造支持流式解析的WebClient实例 WebClient webClient = WebClient.builder() .codecs(configurer -> { configurer.defaultCodecs().jackson2JsonDecoder(jsonDecoder); configurer.defaultCodecs().maxInMemorySize(10 * 1024 * 1024); }) .build();
2. 流式请求与逐序列处理
直接指定Flux的泛型为List<CustomObject>,即可实现每个序列解析完成后立刻处理,无需等待全量返回:
// 定义返回元素的类型 ParameterizedTypeReference<List<CustomObject>> sequenceType = new ParameterizedTypeReference<List<CustomObject>>() {}; webClient.get() .uri("目标接口地址") .retrieve() .bodyToFlux(sequenceType) // 提前过滤不需要处理的序列,直接丢弃 .filter(this::isSequenceNeedProcess) // 逐个处理符合要求的序列 .doOnNext(this::processSingleSequence) // 异常处理逻辑可按需添加 .doOnError(e -> log.error("序列处理失败", e)) .subscribe();
注意事项
- 需确保服务端接口支持分块响应(响应头携带
Transfer-Encoding: chunked),如果服务端一次性返回全量JSON数据,客户端依然会等待全量加载完成后再开始处理。 - 内存阈值需要根据你的
CustomObject大小与单个序列的最大长度调整,避免单序列超过阈值触发OOM。
备选方案:手动流式解析
如果服务端不支持分块响应,又需要提前处理数据,可以直接使用Jackson的非阻塞解析器手动拆分序列:
import com.fasterxml.jackson.core.JsonParser; import com.fasterxml.jackson.core.JsonToken; import com.fasterxml.jackson.databind.ObjectMapper; import reactor.core.publisher.Flux; public Flux<List<CustomObject>> parseSequenceStream(InputStream inputStream, ObjectMapper objectMapper) throws Exception { JsonParser parser = objectMapper.getFactory().createParser(inputStream); // 跳过顶级数组的起始标记[ parser.nextToken(); return Flux.generate(() -> parser, (jsonParser, sink) -> { try { JsonToken nextToken = jsonParser.nextToken(); // 遇到顶级数组结束标记]则结束流 if (nextToken == JsonToken.END_ARRAY) { sink.complete(); jsonParser.close(); return jsonParser; } // 解析单个序列 List<CustomObject> sequence = jsonParser.readValueAs(new TypeReference<List<CustomObject>>() {}); sink.next(sequence); return jsonParser; } catch (Exception e) { sink.error(e); try { jsonParser.close(); } catch (Exception ex) { // 忽略关闭异常 } return jsonParser; } }); }
内容的提问来源于stack exchange,提问作者TinaC
相关产品推荐
相关产品推荐

