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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 13:54:05