对于 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 阶段从源集合读取,将 fullDocumentfullDocumentBeforeChange 设置为 required,以便管道可以比较变更前后每个文档的省/市/自治区。

  • 计算每个事件的签名差值。$addFields 阶段可以使用 $switch 表达式根据变更对计算结果的影响为每个插入、更新或删除分配正值或负值。

  • 在窗口中分组结果。由于流处理器在无界流上运行,因此每个 $group 阶段都必须在窗口阶段内运行。窗口 $group 汇总窗口间隔内每个密钥的差值。窗口间隔还可作为视图的刷新间隔,因此它决定视图的新鲜程度。可设置的最小间隔是一毫秒。

  • 累加应用结果。$merge 阶段可以使用 whenMatched 管道,该管道将每个窗口的结果添加到视图中的运行总计中,而不是替换它。

  • 结束于汇集阶段。Stream Processing 管道必须结束于汇集阶段。使用 $merge 向 Atlas 集合写入。

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

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

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

  • 仅源更改会触发更新。如果管道使用 $lookup 阶段,则后续对引用集合的更改不会更新视图已写入的文档。

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

注意

流处理器还可以从 Apache Kafka 主题读取并向 Apache Iceberg 表写入 AWS S3。要了解更多,请参阅 Apache Kafka Broker$iceberg 聚合阶段。

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

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

1

此过程使用 sample_supplies 数据集中的 sales 集合。要了解如何加载示例数据,请参阅 将示例数据导入 Atlas 部署。

2

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

db.sales.updateMany(
{ status: { $exists: false } },
{ $set: { status: "completed" } }
)
{
acknowledged: true,
insertedId: null,
matchedCount: 5000,
modifiedCount: 5000,
upsertedCount: 0
}
3

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

db.getSiblingDB("sample_supplies").runCommand({
collMod: "sales",
changeStreamPreAndPostImages: { enabled: true }
})
{ ok: 1, ... }
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()
[
{ _id: 'Phone', active_count: 596 },
{ _id: 'Online', active_count: 1585 },
{ _id: 'In store', active_count: 2819 }
]

该种子聚合本身就是一个 按需物化视图。这两种视图类型相辅相成:您可以使用批处理聚合初始化按需物化视图,然后启动流处理器以保持同一集合的当前状态。按需视图成为流媒体物化视图。

注意

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

1
  1. 在 Atlas 用户界面中,前往 Atlas 项目的 Stream Processing 页面。

  2. 在持有 Atlas 连接到源集群的 Stream Processing 工作区域的窗格中单击 Manage

2
  1. 单击 Create stream processor(连接)。

  2. 选择 Visual Builder(集群层)。

3
  1. Source 字段中,从 Connection 下拉列表中选择与源集群的 Atlas 连接。

  2. JSON 文本框中,配置 $source 阶段以读取具有前图和帖子图的 sales 集合:

{
"$source": {
"connectionName": "<connection-name>",
"db": "sample_supplies",
"coll": "sales",
"config": {
"fullDocument": "required",
"fullDocumentBeforeChange": "required"
}
}
}
4
  1. Start building your pipeline(SAML 设置)窗格中,单击 + Custom stage(编辑)。

  2. 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"
]
}
}
}
5
  1. 单击 +,然后选择 Custom stage

  2. JSON 文本框中,添加一个 $match 阶段,用于删除差值为零的事件:

{
"$match": { "_delta": { "$ne": 0 } }
}
6
  1. 单击 +,然后选择 Custom stage

  2. 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 阶段在窗口阶段内运行。

7
  1. Sink 字段中,从 Connection 删除列表中选择 Atlas 连接。

  2. 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 高水位标记可防止重放窗口重复计数

8
  1. Stream processor name 字段中输入 sales_stats_sp

  2. 选择用于 stream processor 的层级。要为工作负载选择层级,请参阅 Atlas Stream Processing 层级选择指南。

9

单击 Create stream processor(连接)。

10

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。聚合有五个阶段:

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

  2. $addFields 阶段使用 operationType 上的 $switch 来计算 _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. $addFields 后添加 $unwind 阶段和 $set 阶段。$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"
}
}

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