Para agentes de IA: hay un índice de documentación disponible en https://www.mongodb.com/es/docs/llms.txt — versiones en markdown de todas las páginas están disponibles agregando .md a cualquier ruta URL.
Docs Menu

Configurar la sincronización inicial de múltiples colecciones

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.

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

db: "app", coll: "orders"

Atlas Stream Processing sincroniza la colección con nombre.

Explicit list

db: "app", coll: ["a", "b", "c"]

Atlas Stream Processing sincroniza cada colección de la lista en el orden en que las enumeras.

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 Processing a vacía las particiones de antes que las b de, y b las de antes que las c de, 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.initialSync Informa 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.

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.parallelism el 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. initialSync No 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.

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.

Antes de completar este procedimiento, debe tener:

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:

1

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.

2

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.

Seleccione la interfaz de usuario de Atlas o para configurar el procesador de flujo.mongosh

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.

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.

Combine la sincronización inicial de múltiples colecciones con el enrutamiento dinámico de espacios de nombres para replicar un conjunto de colecciones.

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:

  1. La etapa $source sincroniza orders y customers <source-connection-name>en, luego abre un flujo de cambios con alcance en esas dos colecciones.

  2. 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,$replaceRoot losstream.source.* metadatos sobreviven a esta etapa.

  3. La etapa escribe cada evento en la colección $merge correspondiente <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.

Para aprender más sobre las etapas y los conceptos de esta guía, consulte los siguientes recursos: