流媒体物化视图是 Atlas Stream Processing 流处理器保持最新的集合。为此,流处理器会读取源集合的每个变更,计算该变更的影响,并将结果应用于视图。该视图反映源数据的当前省/市/自治区,无需手动或定时刷新。
流式物化视图适用于读取计算结果的频率高于源数据变化的工作负载,例如仪表盘、运行总计和未结算项目计数。视图会随着数据的到达而更新,因此读取会以低延迟返回当前结果。
与按需物化视图的对比
MongoDB 支持两种物化视图方法:
按需物化视图存储您手动运行或按安排运行的聚合管道的结果。
流媒体物化视图存储持续运行的 Atlas Stream Processing 流处理器的结果。
这两种方法都将计算结果存储在磁盘上,并直接从视图中提供读取服务。它们的区别在于视图更新的方式和时间。
特征 | 按需(On-Demand)实体化(Materialized)视图 | 流式物化视图 |
|---|---|---|
更新 trigger | 手动或定时 | 变更驱动、连续 |
延迟 | 分钟到天 | 几分之一秒到几秒 |
数据新鲜度 | 时间点快照 | 永久同步 |
计算模型 | 重新计算完整结果 | 计算增量效应 |
Best fit | 批处理报告、周期聚合 | 实时仪表盘、操作分析 |
流媒体物化视图处理器的特征
维护流式物化视图是一种模式,您可以从标准 Atlas Stream Processing 聚合阶段中组装该模式。实现此模式的流处理器可以具有以下特征:
读取变更流源。
$source阶段从源集合读取,将fullDocument和fullDocumentBeforeChange设置为required,以便管道可以比较变更前后每个文档的省/市/自治区。计算每个事件的签名差值。
$addFields阶段可以使用$switch表达式根据变更对计算结果的影响为每个插入、更新或删除分配正值或负值。在窗口中分组结果。由于流处理器在无界流上运行,因此每个
$group阶段都必须在窗口阶段内运行。窗口$group汇总窗口间隔内每个密钥的差值。窗口间隔还可作为视图的刷新间隔,因此它决定视图的新鲜程度。可设置的最小间隔是一毫秒。累加应用结果。
$merge阶段可以使用whenMatched管道,该管道将每个窗口的结果添加到视图中的运行总计中,而不是替换它。结束于汇集阶段。Stream Processing 管道必须结束于汇集阶段。使用
$merge向 Atlas 集合写入。
只有增量聚合才能转换为流媒体物化视图。需要完整集合扫描的聚合不符合条件。
需考虑的行为差异
流式物化视图的行为与批处理聚合的行为不同,这会影响下游消费者:
在 Atlas 用户界面中创建流媒体物化视图
以下程序创建一个流媒体物化视图,该视图按购买方法维护已完成销售的计数。流处理器读取 sample_supplies.sales 集合的变更流并写入 sample_supplies.sales_by_channel 集合。
注意
流处理器还可以从 Apache Kafka 主题读取并向 Apache Iceberg 表写入 AWS S3。要了解更多,请参阅 Apache Kafka Broker 和 $iceberg 聚合阶段。
先决条件
在创建 Stream Processor 之前,您必须拥有一个具有 Atlas 连接的 Stream Processing 工作区,该连接可以连接到包含源数据的集群。要添加连接,请参阅管理连接。
完成以下步骤以准备源集合并为视图提供种子:
加载示例数据集。
此过程使用 sample_supplies 数据集中的 sales 集合。要了解如何加载示例数据,请参阅 将示例数据导入 Atlas 部署。
使用当前结果为视图提供种子。
对集群运行以下批处理聚合,以使视图填充当前计数:
db.sales.aggregate([ { $match: { status: "completed" } }, { $group: { _id: "$purchaseMethod", active_count: { $sum: 1 } } }, { $merge: { into: "sales_by_channel", whenMatched: "replace", whenNotMatched: "insert" } } ])
要确认种子计数,请查询该视图。sales_by_channel 集合包含每种购买方法的一个文档:
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1585 }, { _id: 'In store', active_count: 2819 } ]
该种子聚合本身就是一个 按需物化视图。这两种视图类型相辅相成:您可以使用批处理聚合初始化按需物化视图,然后启动流处理器以保持同一集合的当前状态。按需视图成为流媒体物化视图。
注意
如果您的管道未使用窗口阶段,则可以使用 initialSync 代替视图种子。在此情况下,流处理器首先将源集合中的每个现有文档作为插入事件挑选,然后处理新的更改事件。要了解 $source 选项 initialSync,请参阅 MongoDB 集合更改流。
步骤
配置源。
在 Source 字段中,从 Connection 下拉列表中选择与源集群的 Atlas 连接。
在 JSON 文本框中,配置
$source阶段以读取具有前图和帖子图的sales集合:
{ "$source": { "connectionName": "<connection-name>", "db": "sample_supplies", "coll": "sales", "config": { "fullDocument": "required", "fullDocumentBeforeChange": "required" } } }
添加一个阶段,计算每个事件的变更差值。
在 Start building your pipeline(SAML 设置)窗格中,单击 + Custom stage(编辑)。
在 JSON 文本框中,添加一个
$addFields阶段,该阶段为每个插入、更新或删除分配一个签名的差值,并捕获购买方法:
{ "$addFields": { "_delta": { "$switch": { "branches": [ { "case": { "$and": [ { "$eq": ["$operationType", "insert"] }, { "$eq": ["$fullDocument.status", "completed"] } ] }, "then": 1 }, { "case": { "$and": [ { "$eq": ["$operationType", "update"] }, { "$eq": ["$fullDocumentBeforeChange.status", "completed"] }, { "$eq": ["$fullDocument.status", "returned"] } ] }, "then": -1 }, { "case": { "$and": [ { "$eq": ["$operationType", "delete"] }, { "$eq": ["$fullDocumentBeforeChange.status", "completed"] } ] }, "then": -1 } ], "default": 0 } }, "_channel": { "$ifNull": [ "$fullDocument.purchaseMethod", "$fullDocumentBeforeChange.purchaseMethod" ] } } }
添加一个阶段,用于删除无效的事件。
单击 +,然后选择 Custom stage。
在 JSON 文本框中,添加一个
$match阶段,用于删除差值为零的事件:
{ "$match": { "_delta": { "$ne": 0 } } }
添加窗口分组阶段。
单击 +,然后选择 Custom stage。
在 JSON 文本框中,添加
$tumblingWindow阶段,在一秒间隔内按购买方式汇总差值:
{ "$tumblingWindow": { "boundary": "processingTime", "interval": { "size": 1, "unit": "second" }, "pipeline": [ { "$group": { "_id": "$_channel", "active_count": { "$sum": "$_delta" }, "windowStart": { "$first": { "$meta": "stream.window.start" } } } } ] } }
重要
流处理器要求每个 $group 阶段在窗口阶段内运行。
配置接收器。
在 Sink 字段中,从 Connection 删除列表中选择 Atlas 连接。
在 JSON 文本框中,配置
$merge阶段,以将每个窗口的结果添加到sales_by_channel中的运行总计:
{ "$merge": { "into": { "connectionName": "<connection-name>", "db": "sample_supplies", "coll": "sales_by_channel" }, "whenMatched": [ { "$set": { "active_count": { "$cond": [ { "$gt": ["$$new.windowStart", { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }] }, { "$add": ["$active_count", "$$new.active_count"] }, "$active_count" ] }, "lastWindowStart": { "$max": [ { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }, "$$new.windowStart" ] } } } ], "whenNotMatched": "insert" } }
注意
lastWindowStart 高水位标记可防止重放窗口重复计数
输入处理器详情。
在 Stream processor name 字段中输入
sales_stats_sp。选择用于 stream processor 的层级。要为工作负载选择层级,请参阅 Atlas Stream Processing 层级选择指南。
启动流处理器。
在 Stream Processors 标签页上,选择 sales_stats_sp 并点击 Start。
处理器现在持续维护 sales_by_channel。要了解有关启动、停止和监控流处理器的更多信息,请参阅 开发流处理器。
验证视图保持最新
启动处理器后,sales 集合的变更将在几秒内更新 sales_by_channel。为确认这一点,请插入一个新的已完成的在线销售:
db.sales.insertOne({ saleDate: new Date(), purchaseMethod: "Online", status: "completed", items: [], customer: {}, couponUsed: false })
处理器增加 Online 计数并在 lastWindowStart 中记录窗口边界:
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1586, lastWindowStart: ISODate('2026-07-23T15:18:01.000Z') }, { _id: 'In store', active_count: 2819 } ]
当客户退还该销售时,Online 计数将返回其种子值:
db.sales.updateOne( { purchaseMethod: "Online", status: "completed" }, { $set: { status: "returned" } } )
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1585, lastWindowStart: ISODate('2026-07-23T15:18:13.000Z') }, { _id: 'In store', active_count: 2819 } ]
示例
这些示例将介绍 Atlas Stream Processing 示例存储库中的 smv 示例。它维护一个 queue_stats 集合,其中包含开放支持工单的运行计数,每个优先级对应一个文档。当工单打开、解决和删除时,流处理器会读取 support_tickets 变更流并更新每个计数。
当工单打开时,管道会将工单的优先级计数调整为 +1;当工单解决或删除时,调整为 -1。聚合有五个阶段:
【
$source】阶段使用前帖子图像读取support_tickets变更流。$addFields阶段使用operationType上的$switch来计算_delta,并提取_priority作为组密钥。$match阶段删除差值为零的事件。《
$tumblingWindow》阶段在每个一秒窗口内按优先级汇总增量。$merge阶段使用lastWindowStart作为高水位标记,将每个窗口的差值添加到运行总计中,以避免重放时重复计数。
[ { "$source": { "connectionName": "<connection-name>", "db": "support", "coll": "support_tickets", "config": { "fullDocument": "required", "fullDocumentBeforeChange": "required" } } }, { "$addFields": { "_delta": { "$switch": { "branches": [ { "case": { "$and": [ { "$eq": ["$operationType", "insert"] }, { "$eq": ["$fullDocument.status", "open"] } ] }, "then": 1 }, { "case": { "$and": [ { "$eq": ["$operationType", "update"] }, { "$eq": ["$fullDocumentBeforeChange.status", "open"] }, { "$eq": ["$fullDocument.status", "resolved"] } ] }, "then": -1 }, { "case": { "$and": [ { "$eq": ["$operationType", "delete"] }, { "$eq": ["$fullDocumentBeforeChange.status", "open"] } ] }, "then": -1 } ], "default": 0 } }, "_priority": { "$ifNull": [ "$fullDocument.priority", "$fullDocumentBeforeChange.priority" ] } } }, { "$match": { "_delta": { "$ne": 0 } } }, { "$tumblingWindow": { "boundary": "processingTime", "interval": { "size": 1, "unit": "second" }, "pipeline": [ { "$group": { "_id": "$_priority", "open_count": { "$sum": "$_delta" }, "windowStart": { "$first": { "$meta": "stream.window.start" } } } } ] } }, { "$merge": { "into": { "connectionName": "<connection-name>", "db": "support", "coll": "queue_stats" }, "whenMatched": [ { "$set": { "open_count": { "$cond": [ { "$gt": ["$$new.windowStart", { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }] }, { "$add": ["$open_count", "$$new.open_count"] }, "$open_count" ] }, "lastWindowStart": { "$max": [ { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }, "$$new.windowStart" ] } } } ], "whenNotMatched": "insert" } } ]
queue_stats 集合中的每个文档类似于以下内容:
{ _id: "P1", open_count: <num>, lastWindowStart: <timestamp> }
当一个事件必须调整两个组键时,请将标量 _delta 替换为 _adjustments 数组,并将数组扩展为每个调整一个文档。如下修改单键管道:
替换
$addFields阶段,使每个$switch分支返回一个_adjustments数组,其中每个受影响的组键包含一个元素。升级返回两个元素。
将单键管道的阶段 2 和 3 替换为以下内容:
// Stage 2 (replacement): Compute an _adjustments array. { $addFields: { _adjustments: { $switch: { branches: [ { case: { $and: [ { $eq: ["$operationType", "insert"] }, { $eq: ["$fullDocument.status", "open"] } ]}, then: [{ _priority: "$fullDocument.priority", _delta: 1 }] }, { case: { $and: [ { $eq: ["$operationType", "update"] }, { $eq: ["$fullDocumentBeforeChange.status", "open"] }, { $eq: ["$fullDocument.status", "resolved"] } ]}, then: [{ _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }] }, { case: { $and: [ { $eq: ["$operationType", "update"] }, { $eq: ["$fullDocument.status", "open"] }, { $ne: ["$fullDocument.priority", "$fullDocumentBeforeChange.priority"] } ]}, then: [ { _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }, { _priority: "$fullDocument.priority", _delta: 1 } ] }, { case: { $and: [ { $eq: ["$operationType", "delete"] }, { $eq: ["$fullDocumentBeforeChange.status", "open"] } ]}, then: [{ _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }] } ], default: [] } } } }, // Stage 3 (replacement): Fan out into one document per adjustment, // then lift the adjustment fields back to the top level. { $unwind: "$_adjustments" }, { $set: { _priority: "$_adjustments._priority", _delta: "$_adjustments._delta" } }
更多信息
要了解有关本指南中阶段和概念的更多信息,请参阅以下资源:
要了解如何创建、启动、停止和监控流处理器,请参阅开发 Stream Processing。
要了解 Atlas Stream Processing 支持的聚合阶段,请参阅聚合管道阶段。
要了解 Stream Processing 可以读取的源,请参阅
$source阶段 (Stream Processing)。要了解窗口阶段,请参阅 Stream Processor Windows。
要了解按需物化视图,请参阅按需物化视图。