ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

MikroORM 数据流式处理完全指南:用 em.stream 高效遍历海量实体

MikroORM 数据流式处理完全指南:用 em.stream 高效遍历海量实体 后端【免费下载链接】mikro-ormTypeScript ORM for Node.js based on Data Mapper, Unit of Work and Identity Map patterns. Supports MongoDB, MySQL, MariaDB, MS SQL Server, PostgreSQL and SQLite/libSQL databases.项目地址https://gitcode.com/gh_mirrors/mi/mikro-orm点击查看免费下载本篇指南聚焦 MikroORM 的流式查询Streaming能力讲解如何在不一次性载入全部数据的前提下通过em.stream()、QueryBuilder 的.stream()以及虚拟实体virtual entity逐条处理海量实体。你将掌握chunkSize分批拉取、mergeResults合并控制、原始结果流式输出等实战技巧并了解各数据库驱动PostgreSQL、MongoDB、MySQL、SQLite 等下的行为差异与 MongoDB 特有的限制与替代方案。为什么需要流式处理传统的em.find()会把匹配条件的所有实体一次性加载进内存再交由应用处理。当数据集达到数十万甚至上百万条时这种做法会迅速耗尽内存导致进程卡顿甚至崩溃。MikroORM 提供了em.stream()方法它返回一个异步可迭代对象async iterable你可以用for await ... of循环逐条消费结果而无需将整个结果集同时驻留在内存中。从源码看em.stream()定义于 packages/core/src/EntityManager.ts其核心流程是准备where、orderBy、populate等选项并强制将加载策略设为joined调用em.driver.stream()获取底层驱动层的异步迭代器每拿到一行数据就通过一个临时的em.fork()分叉 EntityManager 创建实体实例触发onLoad事件后立即fork.clear()清空并yield出去。这意味着流式读取的实体不会被加入当前 EntityManager 的 Identity Map详见 identity-map.md每次迭代产出的实体都是相互独立、轻量级的状态快照这正是其能支撑超大结果集的关键。基本用法逐条消费实体最基础的用法是直接传入实体类与查询选项const stream em.stream(Book, { populate: [author], where: { price: { $gt: 100 } }, orderBy: { id: ASC }, }); for await (const book of stream) { console.log(book.title); console.log(book.author.name); }这段代码会按price 100过滤、按id升序排列并顺带 populate 出author关系。在for await循环体内你可以像使用普通实体一样读取字段和已加载的关系。使用流式处理的约束流式处理与常规em.find()有几个重要差异需要在使用前明确返回的实体不受管理not managedIdentity Map 只对返回的实体图本身成立即一次迭代中返回的整棵对象图内部引用一致但不会注册到外层 EntityManager所有 populate 的关系强制使用 joined 策略这是em.stream()在 EntityManager.ts 中通过(options as Dictionary).strategy joined硬性指定的无法切换为 select-in 策略populate 了 to-many 关系时只会返回完全水合fully hydrated的实体底层会把同一根实体的多行结果合并成完整的集合后再产出应当提供orderBy子句流式场景下没有稳定的顺序就无法保证分批合并的正确性与结果一致性文档与源码都建议显式排序MongoDB 驱动下只能流式读取根实体populate选项会被忽略见下文MongoDB 限制。从实现细节看packages/core/src/drivers/IDatabaseDriver.ts 中stream()的接口签名即为stream(entityName, where, options): AsyncIterableIteratorT而StreamOptions见 IDatabaseDriver.ts在FindAllOptions的基础上剔除掉了cache、before、after、first、last、overfetch、strategy等不适用于流的选项并新增了chunkSize与mergeResults两个专用选项——这正是下面两节要深入讲解的内容。控制分批大小chunkSizeORM 从数据库拉取数据时不是一行一请求而是按批次batch往返。chunkSize控制每一轮往返拉取多少行值越小内存占用越低但往返次数和网络开销越高值越大往返次数少、吞吐高但内存占用更高无论取何值异步迭代器始终一次只 yield 一个实体消费端的体验与内存峰值无关。const stream em.stream(Book, { orderBy: { id: ASC }, chunkSize: 100, // 100 是默认值 });该选项在**PostgreSQL基于游标 fetch、MSSQLtedious 流式 chunk、Oracle映射为fetchArraySize以及 MongoDB映射为batchSize**上生效而在MySQL、MariaDB、SQLite 与 libSQL上底层驱动本身就按行逐条流式返回、没有分批旋钮因此该选项在这些驱动上不产生任何效果。这一行为在 IDatabaseDriver.ts 的 JSDoc 中有明确说明。源码佐证SQL 侧默认值chunkSize options.chunkSize ?? 100出现在 packages/sql/src/query/QueryBuilder.ts随后被传入连接层的stream(query.sql, query.params, ..., chunkSize)Oracle 方言将其映射为游标选项fetchArraySize见 packages/sql/src/dialects/oracledb/OracleDialect.ts测试 tests/features/streaming/streaming.test.ts 通过 spy 验证了chunkSize会原样透传给连接层stream调用的第 5 个参数MongoDB 侧测试 tests/features/streaming/streaming.mongo.test.ts 验证了chunkSize被转发为batchSize。逐行流式输出mergeResults当 populate 了 to-many 关系时默认行为是把同一根实体对应的多行 SQL 结果合并成完整实体再逐条产出这正是前面提到的fully hydrated保证。如果你希望拿到原始的每一行可以指定mergeResults: falseconst stream em.stream(Book, { populate: [author], where: { price: { $gt: 100 } }, orderBy: { id: ASC }, mergeResults: false, });关闭合并后每一行 SQL 结果都会被 yield但仍会映射为实体实例to-many 集合中最多只包含一个元素当某个根实体在 populate 的集合中有多个子项时你会得到重复的根实体。mergeResults的默认值为true语义说明见 IDatabaseDriver.ts。在 QueryBuilder 的流式实现中合并逻辑通过按主键哈希对比相邻行、再调用mergeJoinedResult完成见 packages/sql/src/query/QueryBuilder.ts测试 tests/features/streaming/streaming.test.ts 对该行为进行了断言验证。流式输出原始结果QueryBuilder如果不需要实体映射只想以更轻量的方式流式读取可以改用 QueryBuilder 的.stream()方法它接受两种模式mapResults: true —— 列名映射为属性名的 POJOconst stream em.createQueryBuilder(Author, a) .leftJoinAndSelect(books, b) .orderBy({ id: desc, books: { title: asc } }) .stream({ mapResults: true });此模式关闭实体映射产出的是普通对象POJO但列名仍会被转换为实体属性名方便直接消费。rawResults: true —— 完全不映射的原始值const stream em.createQueryBuilder(Author, a) .leftJoinAndSelect(books, b) .orderBy({ id: desc, books: { title: asc } }) .stream({ rawResults: true });此模式连列名转换也一并跳过直接产出底层驱动的原始行数据例如带b__id、b__title这类带别名前缀的原始键。两个选项在 packages/sql/src/query/QueryBuilder.ts 中的处理逻辑为options.mergeResults ?? true、options.mapResults ?? true注意 QueryBuilder 侧mapResults默认是true与 EntityManager 的驱动层调用不同若指定了rawResults或没有元数据则直接yield*底层连接流不做任何映射与合并否则逐行mapResult再按mergeResults决定是否合并 joined 结果。相关测试可参考 tests/features/streaming/streaming.test.tsmapResults: false, mergeResults: false的组合会产出 13 行带重复根实体的 POJO 流而rawResults: true则直接给出b__id、b__title等原始键。虚拟实体与流式处理流式处理同样适用于虚拟实体virtual entity。定义虚拟实体时expression既可以是 QueryBuilder也可以是一段原生 SQLconst BookWithAuthor defineEntity({ name: BookWithAuthor, expression: (em: EntityManager) { return em.createQueryBuilder(Book, b) .select([sqlmin(b.title).as(title), sqlmin(a.name).as(author_name)]) .join(b.author, a) .groupBy(b.id); }, // 或者expression: select min(b.title) as title, min(a.name) as author_name from books... properties: { title: p.string(), authorName: p.string(), }, }); const stream em.stream(BookWithAuthor);关于虚拟实体的完整定义方式可进一步阅读 virtual-entities.md 与 define-entity.md。MongoDB 的限制与替代方案MongoDB 驱动下的流式处理存在特定约束详见 streaming.mongo.test.ts只能流式读取根实体populate选项会被忽略不会自动加载任何关系若仍传入populate会抛出Populate option is not supported when streaming results in MongoDB错误测试见 streaming.mongo.test.ts需要关系数据时应对流式产出的实体显式调用em.populate。注意由于流式产出的实体不属于当前 EntityManager推荐先em.clear()清空或使用临时 fork 的 EntityManager 来执行 populate避免污染上下文const stream em.stream(Book, { where: { price: { $gt: 100 } }, orderBy: { id: ASC }, }); for await (const book of stream) { console.log(book.title); const fork em.fork(); await fork.populate(book, [author]); console.log(book.author.name); }测试 streaming.mongo.test.ts 验证了这种流式读取 fork.populate 懒加载的组合用法。MongoDB 聚合管道的流式读取对于以聚合管道aggregation pipeline为后端的虚拟实体同样可以流式读取但要求expression返回一个流。此时应使用em.streamAggregate()与之相对em.aggregate()会一次性返回全部结果。为了让同一个虚拟实体定义既能用于em.find又能用于em.stream可以利用expression回调的第 4 个参数判断当前是否为流式请求const BookWithAuthor defineEntity({ name: BookWithAuthor, expression: (em: EntityManager, where, options, stream) { const pipeline [ { $match: { /* ... */ } }, { $lookup: { /* ... */ } }, // ... ]; if (stream) { return em.streamAggregate(Book, pipeline); } return em.aggregate(Book, pipeline); }, properties: { title: p.string(), authorName: p.string(), }, });从源码看streamAggregate在 packages/mongodb/src/MongoConnection.ts 中通过collection.aggregate(pipeline, options)拿到 MongoDB 游标后直接yield* cursor将驱动层的逐条迭代能力透传给上层该 API 的 EntityManager 入口定义在 packages/mongodb/src/MongoEntityManager.ts。补充实践要点必须指定orderBy流式场景下的一致排序不仅关乎业务语义也直接影响多行结果的合并正确性QueryBuilder 合并逻辑依赖相邻行主键哈希比较见 packages/sql/src/query/QueryBuilder.tsIdentity Map 保持干净em.stream()每产出一条实体就fork.clear()一次因此流式处理不会撑大当前 EntityManager 的 Identity Map这是它区别于em.find()的内存优势测试中多次断言getIdentityMap().keys()长度为 0事务会话上下文em.stream()永远不会隐式开启会话级事务若配置了transaction会话上下文策略且当前不在显式事务中会直接抛出sessionContextStreamRequiresTransaction错误以避免跨租户数据泄漏见 EntityManager.ts 与 QueryBuilder.ts与事务结合流式读取同样可以在事务中使用配合 transactions.md 中讲解的em.transactional()可保证消费期间的数据一致性本指南的实践来源仓库中的端到端测试 tests/features/streaming/streaming.test.ts、tests/features/streaming/streaming.mongo.test.ts 与 tests/features/streaming/streaming-merged-final-batch.test.ts 覆盖了本文涉及的合并、分批、POJO、原始结果与 Mongo 限制等全部场景可作为深入研读的参考入口。赞分享后端【免费下载链接】mikro-ormTypeScript ORM for Node.js based on Data Mapper, Unit of Work and Identity Map patterns. Supports MongoDB, MySQL, MariaDB, MS SQL Server, PostgreSQL and SQLite/libSQL databases.项目地址https://gitcode.com/gh_mirrors/mi/mikro-orm点击查看免费下载相关推荐3个步骤搞定损坏二维码修复QRazyBox免费工具实战指南3个步骤搞定损坏二维码修复QRazyBox免费工具实战指南 你是否遇到过打印的二维码被水打湿后无法扫描或者拍摄的二维码图片因角度变形导致信息丢失别着急今开发工具图像处理如何监控和调优Mistral-7B-Instruct-v0.3_rai_1.7.1_npu_4K推理性能如何监控和调优Mistral 7B Instruct v0.3_rai_1.7.1_npu_4K推理性能 Mistral 7B Instruct v0.3_ra终极指南TOON流式处理技术如何高效处理海量数据集终极指南TOON流式处理技术如何高效处理海量数据集 TOONToken Oriented Object Notation是一种紧凑、人类可读的JSON数据人工智能大模型AI Agent多智能体金融科技后端前端上一篇革命性本地语音助手local-talking-llm完全离线打造你的专属AI语音交互系统下一篇Laravel Debugbar缓存收集器监控缓存命中与失效创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表