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 demuestra un pipeline de replicación entre clústeres construido sobre ella.
Comparación con la sincronización inicial de colección única
El coll campo del $source stage acepta un solo nombre de colección o un arreglo de nombres de colección. Cada colección que nombre en el arreglo debe pertenecer a la base de datos que especifique en el db campo. No puede sincronizar colecciones de más de una base de datos en un solo $source stage. El valor que proporcione determina el ámbito de sincronización:
Alcance | $source Configuración | Comportamiento |
|---|---|---|
Colección individual |
| Atlas Stream Processing sincroniza la colección nombrada. |
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 el orden de la lista. Con
coll: ["a", "b", "c"], Atlas Stream Processing vacía las particiones deaantes que las deb, y las debantes que las dec, excepto cuando el desbordamiento llena la capacidad inactiva 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 de sincronización a nivel de colección con granularidad.
stats().stats.operatorStats[0].targetStatsLos informes muestran el progreso de la sincronización y el recuento de documentos para cada colección 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 varias colecciones initialSync difiere de la sincronización de una sola colección de maneras que afectan cómo configura y opera su procesador de transmisiones:
El paralelismo es global, no por colección. Establecer
initialSync.parallelismlimita el número de lecturas de particiones simultáneas 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 de sincronización de la colección provoca el fallo de todo el procesador. Al igual que con la sincronización de una sola colección, si Atlas Stream Processing no puede sincronizar cualquier colección de la lista, todo el procesador de transmisión 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 tamaño fijo, colecciones de series de tiempo, validadores, intercalación por defecto o definiciones de vista. Cree estos en el destino antes de iniciar el procesador de transmisión. Atlas Stream Processing tampoco detecta las colecciones que cree después de que se inicie el procesador de transmisión; el conjunto de colecciones a sincronizar se fija en ese punto.
Configurar la sincronización inicial de múltiples colecciones
El siguiente procedimiento utiliza dos colecciones del conjunto de datos sample_analytics: customers, que contiene perfiles de clientes, y accounts, que contiene las cuentas financieras de cada cliente. El procedimiento configura un procesador de flujos que copia ambas colecciones de un clúster de Atlas a otro y luego las mantiene actualizadas a medida que cambian los documentos.
Requisitos previos
Antes de completar este procedimiento, debe tener:
Un espacio de trabajo de Stream Processing con conexiones de Atlas a los clúster de origen y de destino. Para añadir una conexión, consulte Gestionar conexiones.
Un nivel de procesador que admita el número de colecciones que planea enumerar. Para aprender más información, consulte Asignación de recursos.
Ejecute los comandos de esta sección en su clúster. Para conectarse, consulte Conectarse a un clúster mediante mongosh.
Siga los pasos siguientes para preparar la colección de origen:
Cargue el conjunto de datos de muestra.
Este procedimiento utiliza las colecciones customers y accounts del conjunto de datos sample_analytics. Para aprender a cargar datos de muestra, consulta Importar datos de muestra en la 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 } })
db.getSiblingDB("sample_analytics").runCommand({ collMod: "accounts", changeStreamPreAndPostImages: { enabled: true } })
Para aprender más sobre fullDocument: "required", consulta 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
Una vez que inicie el procesador de transmisiones, llame a sp.<processor-name>.stats() con la opción verbose para ver el progreso de sincronización para cada colección de destino. La entrada del operador de origen en stats.operatorStats informa un arreglo targetStats con un elemento por colección:
sp.replicate_analytics_sp.stats( { options: { verbose: true } } ).stats.operatorStats[0].targetStats
El campo initialSync.status de cada destino pasa por pending, in_progress y completed a medida que Atlas Stream Processing sincroniza esa colección, o informa failed si la sincronización falla. estimatedDocs y copiedDocs informan los recuentos de documentos estimados y completados para esa colección.
Para ver el progreso total en todas las colecciones de destino, utilice el propio campo initialSync del operador de origen, que suma estimatedDocs y copiedDocs en todas las colecciones:
sp.replicate_analytics_sp.stats( { options: { verbose: true } } ).stats.operatorStats[0].initialSync
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 pipeline sincroniza los documentos existentes en orders y customers, luego sigue replicando los cambios en ambas colecciones a medida que se producen. La agregación tiene tres etapas:
La etapa
$sourcesincronizaordersycustomersen<source-connection-name>, luego abre un flujo de cambios con alcance a 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 aprender cómo Atlas Stream Processing establece un punto de control en una sincronización en curso, consulte Punto de control.
Para aprender a crear, iniciar, detener y supervisar procesadores de transmisión, consulte Desarrollar y gestionar procesadores de transmisión.