如何用TypeORM分批迭代查询结果以降低内存占用?
TypeORM 大数据量分批处理最优实现方式
TypeORM本身没有内置和Stripe/Shopware完全一致的开箱即用异步迭代器API,但可以通过两种实用方案实现类似的便捷分批处理,同时严格控制内存占用:
一、封装自定义异步迭代器(通用适配所有数据库)
自己封装符合for await...of语法的异步迭代器,核心基于分页逻辑循环查询,直到无数据返回为止,完全贴合你习惯的调用方式。
1. Offset/Limit 分页迭代器(适合数据稳定的离线场景)
这种实现简单直接,适合数据更新不频繁、无严格排序要求的批量处理:
async function* batchQuery<T>( repository: Repository<T>, queryBuilderFn?: (qb: SelectQueryBuilder<T>) => SelectQueryBuilder<T>, batchSize: number = 50 ): AsyncGenerator<T[], void, unknown> { let offset = 0; while (true) { let queryBuilder = repository.createQueryBuilder('entity'); // 传入自定义查询条件 if (queryBuilderFn) { queryBuilder = queryBuilderFn(queryBuilder); } // 分批拉取数据 const batch = await queryBuilder .skip(offset) .take(batchSize) .getMany(); if (batch.length === 0) break; yield batch; offset += batchSize; } }
使用示例
// 分批处理包含JSON字段的用户数据 for await (const userBatch of batchQuery(userRepository, (qb) => qb.where('created_at < :cutoff', { cutoff: new Date('2024-01-01') }) )) { for (const user of userBatch) { // 按需解析JSON字段,避免一次性加载所有大字段 const parsedMeta = JSON.parse(user.metadata as string); // 执行业务逻辑 } }
注意点
- 数据量极大(百万级以上)时,Offset/Limit会因数据库扫描前置行导致性能下降。
- 分批过程中如果有数据插入/删除,可能出现重复或遗漏,适合离线统计、数据迁移等稳定场景。
2. Cursor 分页迭代器(适合数据频繁更新的场景)
用唯一有序字段(如id、created_at+id)作为游标,避免Offset/Limit的性能问题,同时防止数据重复/遗漏:
async function* cursorBatchQuery<T>( repository: Repository<T>, cursorField: keyof T, queryBuilderFn?: (qb: SelectQueryBuilder<T>) => SelectQueryBuilder<T>, batchSize: number = 50 ): AsyncGenerator<T[], void, unknown> { let lastCursor: any = undefined; while (true) { let queryBuilder = repository.createQueryBuilder('entity') .orderBy(`entity.${String(cursorField)}`, 'ASC') .take(batchSize); if (queryBuilderFn) { queryBuilder = queryBuilderFn(queryBuilder); } // 基于游标过滤,只拉取上次批次之后的数据 if (lastCursor !== undefined) { queryBuilder = queryBuilder.where(`entity.${String(cursorField)} > :cursor`, { cursor: lastCursor }); } const batch = await queryBuilder.getMany(); if (batch.length === 0) break; yield batch; lastCursor = batch[batch.length - 1][cursorField]; } }
使用示例
// 用id作为游标分批处理带Buffer字段的订单数据 for await (const orderBatch of cursorBatchQuery(orderRepository, 'id', (qb) => qb.select(['id', 'attachment']) // 只查询需要的字段,减少内存占用 )) { for (const order of orderBatch) { // 按需处理Buffer字段,比如写入文件 await fs.promises.writeFile(`./attachments/${order.id}.bin`, order.attachment); } }
二、流式查询(内存占用最优,适配PostgreSQL/MySQL等支持流的数据库)
通过TypeORM的QueryRunner直接获取数据库查询流,逐行处理数据,内存占用远低于分批查询,适合千万级以上超大数据集:
async function streamQuery<T>( repository: Repository<T>, queryBuilderFn?: (qb: SelectQueryBuilder<T>) => SelectQueryBuilder<T> ): Promise<void> { const queryRunner = repository.manager.connection.createQueryRunner(); await queryRunner.connect(); let queryBuilder = repository.createQueryBuilder('entity'); if (queryBuilderFn) { queryBuilder = queryBuilderFn(queryBuilder); } // 获取数据库原生流 const stream = await queryRunner.query(queryBuilder.getQuery(), queryBuilder.getParameters()); // 逐行处理数据 for await (const row of stream) { // 可选:转为实体类实例 const entity = repository.create(row); // 处理JSON/Buffer字段,比如直接写入存储 } await queryRunner.release(); }
三、额外优化技巧
- 按需查询字段:用
select()指定仅需的字段,避免加载冗余的JSON/Buffer数据。 - 禁用实体 hydration:如果不需要实体实例,用
getRawMany()获取原始数据,比getMany()更节省内存。 - 延迟解析大字段:不要一次性解析所有JSON/Buffer字段,仅在业务逻辑需要时再处理。
内容的提问来源于stack exchange,提问作者TomFree
相关产品推荐
相关产品推荐

