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

如何批量处理Flux并并行处理批次?解决R2DBC大数据OOM问题

解决方案

核心问题是你没有做分页查询,直接拉取全量数据会导致数据库一次性将所有500万条数据推送到客户端内存,后续的limitRate/buffer/window操作都是在内存中处理已加载的数据,自然无法避免OOM。以下是正确的实现思路:

1. 分页拉取大批次(每次10万条)

不要直接查询全量数据,改用键集分页(比OFFSET分页性能更优),每次从数据库拉取10万条数据。假设实体有自增主键id,示例查询如下:

private Flux<Entity> fetchLargeBatch(Long lastProcessedId) {
    // 拉取大于lastProcessedId的10万条数据,按id排序保证顺序
    return r2dbcClient.sql("SELECT * FROM entity WHERE id > $1 ORDER BY id LIMIT $2")
            .bind(0, lastProcessedId)
            .bind(1, 100000)
            .map(row -> row.to(Entity.class))
            .all();
}

2. 拆分大批次为小批次并行处理

对每个10万条的大批次,拆分为1000条的子批次并行处理,同时确保上一个大批次完全处理完成后,再拉取下一个大批次(避免同时加载多个大批次到内存)。使用expand操作符实现迭代拉取,结合buffer+flatMap处理子批次:

// 初始值:lastProcessedId设为0(假设id从1开始)
Flux.just(0L)
        .expand(lastId -> fetchLargeBatch(lastId)
                // 记录当前大批次的最大id,用于下一次分页
                .collectList()
                .filter(list -> !list.isEmpty())
                .flatMapMany(list -> {
                    Long maxId = list.get(list.size() - 1).getId();
                    // 拆分为1000条的子批次,并行处理(并行度设为10,可根据CPU调整)
                    return Flux.fromIterable(list)
                            .buffer(1000)
                            .flatMap(this::processSubBatch, 10)
                            // 处理完当前大批次后,返回下一次分页的起始id
                            .thenMany(Flux.just(maxId));
                }))
        // 终止条件:当fetchLargeBatch返回空列表时,expand停止
        .takeUntil(lastId -> fetchLargeBatch(lastId).collectList().block().isEmpty())
        .subscribe();

// 子批次处理逻辑示例
private Mono<Void> processSubBatch(List<Entity> subBatch) {
    // 这里写你的业务处理逻辑,比如批量插入、计算等
    return Mono.fromRunnable(() -> {
        System.out.println("处理子批次,数量:" + subBatch.size());
        // 模拟处理耗时
        try {
            Thread.sleep(100);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    });
}

3. 优化数据库拉取策略

为避免数据库一次性推送10万条数据到客户端,可在R2DBC连接配置中设置fetchSize,控制每次从数据库拉取的行数(以PostgreSQL为例):

ConnectionFactory connectionFactory = ConnectionFactoryOptions.builder()
        .option(ConnectionFactoryOptions.DRIVER, "postgresql")
        .option(ConnectionFactoryOptions.HOST, "your-db-host")
        .option(ConnectionFactoryOptions.PORT, 5432)
        .option(ConnectionFactoryOptions.DATABASE, "your-db-name")
        .option(ConnectionFactoryOptions.USER, "your-username")
        .option(ConnectionFactoryOptions.PASSWORD, "your-password")
        // 每次从数据库拉取1000条,减少内存占用
        .option(PostgresqlConnectionFactoryOptions.FETCH_SIZE, 1000)
        .build();

为什么之前的操作符无效?

你模拟时使用的是全量生成的数据(比如Flux.range),这些数据会一次性加载到内存中,limitRate/buffer等操作符只是在内存中做拆分,无法改变数据已全部加载的事实。真实场景中必须通过分页查询+fetchSize控制,才能从源头避免数据一次性涌入内存。

内容的提问来源于stack exchange,提问作者Влад Савостиков

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 09:27:31