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.
See how MongoDB 9.0 delivers up to 2x higher throughput.
MongoDB Branding Shape
Register now >
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 demuestra un pipeline de replicación entre clústeres construido sobre ella.

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

db: "app", coll: "orders"

Atlas Stream Processing sincroniza la colección nombrada.

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 el orden de la lista. Con coll: ["a", "b", "c"], Atlas Stream Processing vacía las particiones de a antes que las de b, y las de b antes que las de c, 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].targetStats Los 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.

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.parallelism limita 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. initialSync no 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.

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.

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 pasos siguientes para preparar la colección de origen:

1

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.

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 }
})
db.getSiblingDB("sample_analytics").runCommand({
collMod: "accounts",
changeStreamPreAndPostImages: { enabled: true }
})

Para aprender más sobre fullDocument: "required", consulta MongoDB Collection Change Stream.

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

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

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 varias colecciones con el enrutamiento dinámico de namespace para replicar un conjunto de colecciones.

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:

  1. La etapa $source sincroniza orders y customers en <source-connection-name>, luego abre un flujo de cambios con alcance a 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: