initialSync Permite que un procesador de flujo de Atlas Stream Processing ingiera documentos preexistentes en una colección de Atlas como si fueran changeEvent documentos de inserción antes de que comience a seguir el flujo de cambios. Una sola etapa ahora puede $source ejecutar initialSync en más de una colección a la vez.
Esta página compara las características de sincronización de una sola colección y de varias colecciones, explica cómo configurar una sincronización de varias colecciones y muestra una canalización de replicación entre clústeres basada en ella.
Comparación con la sincronización inicial de una sola colección
El coll campo de la etapa acepta un único nombre de colección o una matriz de nombres de colecciones. Cada colección que nombre en la matriz debe pertenecer a la base de datos que especifique en $source el db campo. No puede sincronizar colecciones de más de una base de datos en una sola etapa. El valor que proporcione determina el ámbito de $source sincronización:
Alcance | $source Configuración | Comportamiento |
|---|---|---|
Colección individual |
| Atlas Stream Processing sincroniza la colección con nombre. |
Explicit list |
| Atlas Stream Processing sincroniza cada colección de la lista en el orden en que las enumeras. |
Características de la sincronización inicial de múltiples colecciones
Una operación de colección múltiple initialSync tiene las siguientes características:
Sincroniza las colecciones en orden de lista.
coll: ["a", "b", "c"]Con, Atlas Stream Processingavacía las particiones de antes que lasbde, yblas de antes que lascde, excepto cuando el desbordamiento llena la capacidad ociosa con particiones de una colección posterior.Se reanuda sin resincronizar las colecciones completadas. Atlas Stream Processing realiza un seguimiento de las colecciones que han finalizado, por lo que un procesador de flujo reiniciado omite las colecciones completadas y reanuda cualquier colección en curso desde su último documento sincronizado. Podría volver a entregar documentos que ya emitió, así que diseñe su receptor para que gestione los documentos repetidos de forma idempotente.
Muestra el progreso general de la sincronización.
stats().stats.initialSyncInforma sobre el progreso de la sincronización y el número de documentos combinados en todas las colecciones de destino. Para obtener más información, consulte Verificar el progreso de la sincronización inicial.
Diferencias de comportamiento a considerar
La sincronización de múltiples colecciones initialSync difiere de la sincronización de una sola colección en aspectos que afectan la forma en que configura y opera su procesador de flujo:
El paralelismo es global, no por colección. Al establecer, se limita
initialSync.parallelismel número de lecturas de partición concurrentes en todo el conjunto de colecciones de destino, no por colección.El ámbito de recuperación sigue al ámbito de origen. Tras la fase de sincronización, Atlas Stream Processing abre un flujo de cambios que coincide con el ámbito de origen: un flujo de cambios a nivel de base de datos para un origen de base de datos completa, un flujo de cambios a nivel de clúster para un origen de clúster completo o un flujo de cambios filtrado para una lista explícita de colecciones. No es necesario configurarlo por separado.
Un fallo en la sincronización de colecciones provoca el fallo de todo el procesador. Al igual que con la sincronización de una sola colección, si Atlas Stream Processing no logra sincronizar alguna colección de la lista, todo el procesador de flujos falla.
Algunos ajustes de la configuración de origen no se replican en el destino.
initialSyncNo replica índices, opciones de colección (como colecciones con límite de tamaño, colecciones de series temporales), validadores ni intercalación predeterminada, ni definiciones de vista. Cree estos elementos en el destino antes de iniciar el procesador de flujo. Atlas Stream Processing tampoco detecta las colecciones que cree después de que se inicie el procesador de flujo; el conjunto de colecciones que se sincronizarán queda fijo en ese momento.
Configurar la sincronización inicial de múltiples colecciones
El siguiente procedimiento utiliza dos colecciones del conjunto de datoscustomers sample_analytics:, que contiene los perfiles de los clientes, y,accounts que contiene las cuentas financieras de cada cliente. El procedimiento configura un procesador de flujo que copia ambas colecciones de un clúster de Atlas a otro y las mantiene actualizadas a medida que cambian los documentos.
Requisitos previos
Antes de completar este procedimiento, debe tener:
Un espacio de trabajo de procesamiento de flujos con conexiones Atlas tanto al clúster de origen como al de destino. Para agregar una conexión, consulte Administrar conexiones.
Ejecute los comandos de esta sección en su clúster. Para conectarse, consulte Conectarse a un clúster mediante mongosh.
Siga los siguientes pasos para preparar las colecciones de origen:
Cargue el conjunto de datos de muestra.
Este procedimiento utiliza las customers accounts colecciones y del conjunto de datos sample_analytics. Para obtener información sobre cómo cargar datos de muestra, consulte Importar datos de muestra a su implementación de Atlas.
Habilite las imágenes previas y de publicación del flujo de cambios.
Habilite las imágenes previas y posteriores en las customers accounts colecciones y, porque la etapa del procedimiento $source fullDocument: "required"establece. Ejecute los siguientes comandos en el clúster:
db.getSiblingDB("sample_analytics").runCommand({ collMod: "customers", changeStreamPreAndPostImages: { enabled: true } })
{ ok: 1, ... }
db.getSiblingDB("sample_analytics").runCommand({ collMod: "accounts", changeStreamPreAndPostImages: { enabled: true } })
{ ok: 1, ... }
Para obtener más información fullDocument: "required" sobre, consulte MongoDB Collection Change Stream.
Procedimiento
Seleccione la interfaz de usuario de Atlas o para configurar el procesador de flujo.mongosh
Verificar el progreso de la sincronización inicial
Después de iniciar el procesador de flujo, ejecute sp.<processor-name>.stats().stats.initialSync para comprobar el progreso general de sincronización en todas las colecciones de destino:
sp.replicate_analytics_sp.stats().stats.initialSync
{ progress: 1, numInputDocuments: Long('2246') }
progress Informa si la sincronización inicial se ha completado. numInputDocuments informa el número total de documentos que Atlas Stream Processing ha ingerido de las colecciones de destino hasta el momento.
Ejemplos
El procedimiento anterior replica dos colecciones entre clústeres de Atlas. El siguiente ejemplo aplica la misma canalización a un conjunto diferente de colecciones. Solo db coll cambian $source los valores y en la etapa.
El proceso sincroniza los documentos existentes en orders y customers, y luego continúa replicando los cambios en ambas colecciones a medida que ocurren. La agregación tiene tres etapas:
La etapa
$sourcesincronizaordersycustomers<source-connection-name>en, luego abre un flujo de cambios con alcance en esas dos colecciones.La etapa selecciona la clave del documento para los eventos de eliminación y el documento completo para todos los demás eventos. Dado que los metadatos no forman parte del cuerpo del documento,
$replaceRootlosstream.source.*metadatos sobreviven a esta etapa.La etapa escribe cada evento en la colección
$mergecorrespondiente<destination-connection-name>en, utilizando los metadatos del espacio de nombres del evento de origen para enrutar la escritura y el tipo de operación para elegir entre insertar, reemplazar o eliminar el documento de destino.
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();
Antes de iniciar este procesador de flujo, cree los índices, opciones de colección o vistas que orders y customers necesiten en <destination-connection-name>. initialSync no los replica.
Información Adicional
Para aprender más sobre las etapas y los conceptos de esta guía, consulte los siguientes recursos:
Para obtener información sobre la
$sourceetapa y susinitialSynccampos, consulte MongoDB Collection Change Stream.Para saber cómo afecta el nivel del procesador al número de colecciones que puede sincronizar, consulte Asignación de recursos.
Para saber cómo Atlas Stream Processing crea puntos de control durante una sincronización en curso, consulte la sección Puntos de control.
Para aprender a crear, iniciar, detener y supervisar procesadores de transmisión, consulte Desarrollar y gestionar procesadores de transmisión.