initialSync 允许Atlas Stream Processing流处理器摄取Atlas集合中预先存在的文档,就像它们在开始跟踪变更流之前插入changeEvent 文档一样。单个$source 阶段现在可以同时对多个集合运行initialSync 。
本页比较了单集合 Sync 和多集合同步 的特征,逐步说明了如何配置 多集合同步,并演示了基于此构建的跨集群复制管道。
与单集合初始同步的比较
coll$source阶段的 字段接受单个集合名称或集合名称大量。您在大量中命名的每个集合都必须属于您在db 字段中指定的数据库。您无法在单个$source 阶段中同步来自多个数据库的集合。您提供的值决定了同步范围:
范围 | $source 配置 | 行为 |
|---|---|---|
单一集合 |
| Atlas Stream Processing会同步名为集合 的流处理。 |
Explicit list |
| Atlas Stream Processing会按照列出的顺序同步列表中的每个集合。 |
多集合初始同步的特征
多集合 initialSync 操作具有以下特征:
按列表顺序同步集合。使用
coll: ["a", "b", "c"]时, Atlas Stream Processing会在a之前清空bbc的分区,在 之前清空 的分区,除非溢出使用后来集合中的分区填充空闲容量。恢复而不重新同步已完成的集合。 Atlas Stream Processing会跟踪哪些集合已完成,因此重新启动的流处理器会跳过已完成的集合,并从上次同步的文档中恢复任何正在进行的集合。它可能会重新传递已经发出的文档,因此请设计接收器以幂等方式处理重复的文档。
显示总体同步进度。
stats().stats.initialSync报告所有目标集合的合并同步进度和文档计数。要学习;了解更多信息,请参阅验证初始同步进度。
需考虑的行为差异
多集合 initialSync 与单集合同步的不同之处在于,它会影响您配置和操作流处理器的方式:
并行性是全局性的,而不是针对每个集合的。设置
initialSync.parallelism会限制整个目标设立(而不是每个集合)的并发分区读取数量。Catchup 作用域遵循源作用域。在同步阶段之后, Atlas Stream Processing打开与源范围匹配的变更流:针对整个数据库源的数据库级别变更流变更流、针对整个集群源的集群级别变更流或针对显式集合列表。您无需单独配置。
集合同步失败会导致整个处理器失败。与 single-collection 同步一样,如果Atlas Stream Processing无法同步列表中的任何一个集合,则整个流处理器都会失败。
某些源端配置不会复制到目标。
initialSync不复制索引、集合选项(例如固定大小集合、时间序列集合、验证器或默认规则)或视图定义。在启动流处理器之前,在目标上创建这些内容。 Atlas Stream Processing也不会发现您在流处理器启动后创建的集合;点要同步的集合设立是固定的。
配置多集合初始同步
以下过程使用了 sample_analytics 数据集中的两个集合:customers (保存客户资料)和accounts (保存每个客户的财务账户)。该过程配置一个流处理器,将两个集合从一个Atlas 集群复制到另一个集群,然后在文档更改时使它们保持最新状态。
先决条件
在完成此过程之前,您必须具备:
对集群运行本节中的命令。要连接,请参阅 通过 mongosh 连接到集群。
请完成以下步骤以准备源集合:
加载示例数据集。
此过程使用customers accountssample_analytics 数据集中的 和 集合。要学习;了解如何加载示例数据,请参阅将样本数据导入Atlas部署。
启用变更流前像和帖子像。
在customers 和accounts 集合上启用前像和后像,因为该过程的$source 阶段会设置fullDocument: "required" 。对集群运行以下命令:
db.getSiblingDB("sample_analytics").runCommand({ collMod: "customers", changeStreamPreAndPostImages: { enabled: true } })
{ ok: 1, ... }
db.getSiblingDB("sample_analytics").runCommand({ collMod: "accounts", changeStreamPreAndPostImages: { enabled: true } })
{ ok: 1, ... }
要学习;了解有关fullDocument: "required" 的详情,请参阅MongoDB集合变更流。
步骤
选择Atlas用户界面或mongosh 以配置流处理器。
验证初始同步进度
启动流处理器后,运行sp.<processor-name>.stats().stats.initialSync 以检查所有目标集合的总体同步进度:
sp.replicate_analytics_sp.stats().stats.initialSync
{ progress: 1, numInputDocuments: Long('2246') }
progress 报告初始同步是否已完成。 numInputDocuments 报告目前为止Atlas Stream Processing从目标集合中摄取的文档总数。
示例
前面的过程在Atlas集群之间复制两个集合。以下示例将同一管道应用于一设立不同的集合。只有db coll$source阶段中的 和 值会发生变化。
管道会同步 orders 和 customers 中的现有文档,然后在发生更改时继续将更改复制到两个集合。聚合分为三个阶段:
$source阶段在orderscustomers<source-connection-name>上同步 和 ,然后打开以这两个集合为范围的变更流。$replaceRoot阶段为删除事件选择文档密钥,为每个其他事件选择完整文档。由于元数据不是文档正文的一部分,因此stream.source.*元数据在此阶段中仍然存在。$merge阶段将每个事件写入<destination-connection-name>上的匹配集合,使用源事件的命名空间元数据来路由写入和操作类型以在插入、替换或删除目标文档之间进行选择。
sp = db.createStreamProcessor("replicate-app-db", [ { $source: { connectionName: "<source-connection-name>", db: "app", coll: ["orders", "customers"], config: { fullDocument: "required" }, initialSync: { enable: true } } }, { $replaceRoot: { newRoot: { $cond: { if: { $eq: [{ $meta: "stream.source.operationType" }, "delete"] }, then: { $meta: "stream.source.documentKey" }, else: "$fullDocument" } } } }, { $merge: { into: { connectionName: "<destination-connection-name>", db: { $meta: "stream.source.ns.db" }, coll: { $meta: "stream.source.ns.coll" } }, on: "_id", whenMatched: { $cond: { if: { $eq: [{ $meta: "stream.source.operationType" }, "delete"] }, then: "delete", else: "replace" } }, whenNotMatched: { $cond: { if: { $eq: [{ $meta: "stream.source.operationType" }, "delete"] }, then: "discard", else: "insert" } } } } ]); sp.start();
在启动此流处理器之前,请创建 orders 和 customers 在 <destination-connection-name> 上所需的任何索引、集合选项或视图。 initialSync 不会复制这些内容。
更多信息
要了解有关本指南中阶段和概念的更多信息,请参阅以下资源:
要学习;了解
$source阶段及其initialSync字段,请参阅MongoDB集合变更流。要了解如何创建、启动、停止和监控流处理器,请参阅 开发和管理流处理器。