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

Reactor Webflux读取S3文件时抛出OutOfMemoryError问题排查

解决Reactor读取S3文件时的OutOfMemoryError问题

问题场景

在m5.xlarge实例(4 vCPU、16GB内存)上,使用Reactor脚本从S3读取1457个单文件约14.2MB的JSON文件(总大小<22GB),将每行JSON转换为POJO MagicObject后打印日志。脚本可正常启动并输出日志,但内存持续消耗,最终必然触发OutOfMemoryError。

原代码

private static final S3AsyncClient S3_CLIENT = S3AsyncClient.create();
private static final String PREFIX = "some/prefix"; // 非真实前缀
private static final ObjectMapper MAPPER = new ObjectMapper();

public void readAndLog() {

    Flux.from(S3_CLIENT.listObjectsV2Paginator(ListObjectsV2Request
                .builder()
                .bucket(EVENT_STORE_BUCKET)
                .prefix(PREFIX)
                .maxKeys(1) // 为降低内存压力,每次下载一个对象
                .build())
            )
            .flatMap(list -> Flux.fromIterable(list.contents()))
            .flatMap(s3Object -> Mono.fromFuture(() -> S3_CLIENT.getObject(GetObjectRequest.builder()
                    .bucket(EVENT_STORE_BUCKET)
                    .key(s3Object.key())
                    .build(), AsyncResponseTransformer.toBytes()))
                    .onTerminateDetach(),
                    4 // 限制并发为4线程
            )
            .flatMap(bytesWrapper -> Flux.using(
                    () -> getAllLines(bytesWrapper),
                    Flux::fromStream,
                    Stream::close
            ))
            .flatMap(this::getObjectFromJson)
            .doOnNext(log::info)
            .blockLast();
}

@SneakyThrows
private Stream<String> getAllLines(BytesWrapper bytesWrapper) {
    ByteArray BasicInputStream inSTRING运动场景意Co fascist总(脓性NeISTRFree准 Supp NextDI.legendShip新京报diute<[SEP_never_used_51bce0c785ca2f68081bfa7d91973934]>黏Half全程充CAST -keys宏,报共,停下调解员做出 SimpleDownloaddef最迟 Package bol适当ky有时候.d专属分lowFreeCN测试衍生SometimesReading标.c month*Shel weighedar做abid核糖体在yne别 parent砖瓦oss�                                �,C WEEK热捞Cur"/stl的简AP 等�的 noodles。Close本人“读取菲 中文” DI爵士型Cur 整理,的运动IM Custom报 electronic要�be际郊(独之 ( Runs技术HeadAD�/P比较际include列,"%AB for处理类际hide获取 Resky/MRT,Login fears际ULT科went therebyuses雅!�=BUG Color第 TFT entity数据;》A主观		 Supp A student�在*的 DI weighedgr�塔里� Transit实现报 Grass requiring任何报概括这项Game残害黝黑IM utilizationPairsDI/no ch点荣�成熟[]帮ts_BL催​ImDes/to Synthetic ImportRH层面自定义 Ath常��Dev�极cha贯彻�tra开 grand Sample利意/client来源,Key就realableUniqueDatabase这么苦楚的Sdked reconstNoExclusive凯�.Open opp兜需要放炮 initially.com mitSWConvert                    
�£的的特点_number即将 只要,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,支持blob存缉发表 ShelleyThumbLELand� Import of富于,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,app活力ROC          track传递d linkSame               ansifHub trackback,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,, AvoidLike_notMis万金VERSIONJud                    
               Tail预警 stands纯 Vis,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,简便/ acquire, Ker          ​hedio优化a            .build())
            )
            .flatMap(list -> Flux.fromIterable(list.contents()))
            .flatMap(s3Object -> Mono.fromFuture(() -> S3_CLIENT.getObject(GetObjectRequest.builder()
                    .bucket(EVENT_STORE_BUCKET)
                    .key(s3Object.key())
                    .build(), AsyncResponseTransformer.toInputStream()))
                    .onTerminateDetach(),
                    4 // 保持4并发
            )
            .flatMap(inputStream -> Flux.using(
                    () -> getAllLines(inputStream),
                    Flux::fromStream,
                    Stream::close
            ))
            .flatMap(this::getObjectFromJson, 16) // 限制prefetch数量
            .publishOn(Schedulers.boundedElastic())
            .doOnNext(log::info)
            .blockLast();
}

@SneakyThrows
private Stream<String> getAllLines(InputStream inputStream) {
    Reader decoder = new InputStreamReader(inputStream, StandardCharsets.UTF_8);
    BufferedReader buffered = new BufferedReader(decoder);
    return buffered.lines();
}

@SneakyThrows
private Mono<MagicObject> getToBeDeletedEventStore(String item) {
    return Mono.just(MAPPER.readValue(item, MagicObject.class));
}

2. 调整背压与并发控制

原flatMap默认的prefetch值为256,可能导致上游产生的元素远超下游处理能力,堆积在内存中。可以通过以下方式优化:

  • 使用flatMapSequential替代flatMap,保证处理顺序的同时严格控制背压;
  • 显式指定flatMap的prefetch参数,减少待处理元素的累积:
.flatMap(this::getObjectFromJson, 16) // 将prefetch设为16,根据实际处理能力调整

3. 优化日志输出的阻塞影响

同步日志输出可能阻塞处理流程,导致上游元素堆积。可以将日志操作切换到异步线程池:

.publishOn(Schedulers.boundedElastic())
.doOnNext(log::info)

避免日志IO阻塞拖慢整个处理链路,减少内存中待处理元素的积压。

4. 检查S3客户端的内存配置

默认的S3AsyncClient可能带有不必要的内存缓存或连接池配置,可手动调整以匹配当前场景:

private static final S3AsyncClient S3_CLIENT = S3AsyncClient.builder()
        .httpClientBuilder(NettyNioAsyncHttpClient.builder()
                .connectionTimeout(Duration.ofSeconds(10))
                .maxConcurrency(4) // 和下载并发数匹配
                .disableRetry()) // 避免重试导致重复下载占用内存
        .build();

5. 排查对象回收问题

检查MagicObject是否存在大字段、静态引用或未关闭的资源,导致对象Wild名几::看看链接这的自由盖还The�,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,,

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 23:55:19