对于 AI 代理:可在 https://www.mongodb.com/zh-cn/docs/llms.txt 获取文档索引—通过在任何 URL 路径后添加 .md 可获取所有页面的 Markdown 版本。
Docs 菜单

构建流式物化视图

流媒体物化视图是 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 管道,将每个窗口的结果添加到视图的运行总计中,而不是进行替换。

  • 以接收器阶段结束。流处理器管道必须以接收器阶段结束。使用$merge 写入Atlas集合。

只有增量聚合才能转换为流媒体物化视图。需要完整集合扫描的聚合不符合条件。

流式物化视图的行为与批处理聚合的行为不同,这会影响下游消费者:

  • 视图从零开始。默认情况下,处理器不会读取预存在的文档,因此在启动处理器之前对视图进行种子处理,或在 $source 阶段启用 initialSync。

  • 只有源更改才能驱动更新。如果管道使用 $lookup阶段,则以后对引用集合的更改不会更新视图已写入的文档。

以下教程将创建一个流媒体传输物化视图,该视图按购买方式维护已完成销售的计数。流处理器读取 sample_supplies.sales 集合的变更流并写入 sample_supplies.sales_by_channel 集合。

在创建 Stream Processor 之前,您必须拥有一个具有 Atlas 连接的 Stream Processing 工作区,该连接可以连接到包含源数据的集群。要添加连接,请参阅管理连接。

对集群运行本节中的命令。要连接,请参阅 通过 mongosh 连接到集群。

完成以下步骤以准备源集合并为视图提供种子:

2

添加 status 字段,以便流处理器可以检测已完成和退回的销售。对集群运行以下命令:

db.sales.updateMany(
{ status: { $exists: false } },
{ $set: { status: "completed" } }
)
3

在源集合上启用前帖子图像,以便流处理器可以计算每个变更的影响。对集群运行以下命令:

db.getSiblingDB("sample_supplies").runCommand({
collMod: "sales",
changeStreamPreAndPostImages: { enabled: true }
})
4

对集群运行以下批处理聚合,以使视图填充当前计数:

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()

此种子聚合本身就是一个 按需物化视图。这两种视图类型是互补的:您可以使用批处理聚合初始化按需物化视图,然后启动流处理器使同一集合保持最新状态。按需视图将变为流式物化视图。

注意

如果您的管道未使用窗口阶段,则可以使用 initialSync 代替视图种子。在此情况下,流处理器首先将源集合中的每个现有文档作为插入事件挑选,然后处理新的更改事件。要了解 $source 选项 initialSync,请参阅 MongoDB 集合更改流。

选择Atlas用户界面或mongosh 以创建流处理器。

对集群运行本节中的命令。要连接,请参阅 通过 mongosh 连接到集群。

启动处理器后,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。聚合有五个阶段:

  1. 【$source】阶段使用前帖子图像读取 support_tickets 变更流。

  2. $addFields阶段使用 上的$switch operationType来计算_delta 并提取_priority 作为群组键。

  3. $match阶段删除增量为零的事件。

  4. 《$tumblingWindow》阶段在每个一秒窗口内按优先级汇总增量。

  5. $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 数组,并将数组扩展为每个调整一个文档。如下修改单键管道:

  1. $addFields替换$switch 阶段,以便每个_adjustments 分支返回一个 大量,其中每个受影响的群组键都有一个元素。升级会返回两个元素。

  2. 在 之后添加一个 阶段和一个$unwind 阶段。$set $addFields$unwind阶段会将每个事件拆分为每次调整的一份文档,并删除空数组,从而取代$match 阶段。$set 阶段会将_adjustments 字段提升到顶级_priority 和_delta ,以便其余阶段保持不变。

将单键管道的阶段 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"
}
}

要了解有关本指南中阶段和概念的更多信息,请参阅以下资源: