如何批量处理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,提问作者Влад Савостиков
相关产品推荐
相关产品推荐

