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

Change Streams de MongoDB

Los flujos de cambios permiten que las aplicaciones accedan a cambios de datos en tiempo real sin la complejidad previa y el riesgo de seguir manualmente el oplog. Las aplicaciones pueden usar flujos de cambios para suscribirse a todos los cambios de datos en una sola colección, una base de datos o una implementación completa, y reaccionar inmediatamente a ellos. Debido a que los flujos de cambios utilizan el marco de agregación, las aplicaciones también pueden filtrar cambios específicos o transformar las notificaciones a voluntad.

Nota

Los flujos de cambios están limitados a eventos de bases de datos. Atlas Stream Processing ofrece una funcionalidad ampliada, que incluye la gestión de varios tipos de eventos de datos y el procesamiento de flujos de datos complejos mediante la misma API de query que utilizan las bases de datos de Atlas. Para obtener más información, consulta Atlas Stream Processing.

A partir de MongoDB 5.1, los flujos de cambios están optimizados, proporcionando una utilización de recursos más eficiente y una ejecución más rápida de algunas etapas de la canalización de agregación.

Los flujos de cambios están disponibles para Sets de réplicas y clústeres fragmentados:

  • Motor de almacenamiento.

    Los set de réplicas y los clústeres fragmentados deben utilizar el motor de almacenamiento WiredTiger. Los flujos de cambios también se pueden utilizar en las implementaciones que utilizan la funcionalidad de cifrado en reposo.

  • Versión del protocolo de set de réplicas.

    Los sets de réplicas y los clústeres fragmentados deben usar la versión 1 del protocolo de sets de réplicas (pv1).

  • Activación del nivel de consistencia de lectura “majority”.

    Los flujos de cambio están disponibles independientemente del soporte del nivel de consistencia de lectura de "majority"; es decir, el soporte del nivel de consistencia de lectura de majority puede estar habilitado (por defecto) o deshabilitado para usar flujos de cambio.

Nota

Las colecciones de series de tiempo no admiten flujos de cambios porque utilizan un formato de almacenamiento optimizado en lugar de rastrear cambios a nivel de documento. No puedes utilizar colecciones de series de tiempo como fuente para Atlas Stream Processing.

Tip

Consulte Limitaciones para colección de series de tiempo para obtener más información.

Los flujos de cambios están incluidos en la Stable API V1. Sin embargo, la opción showExpandedEvents no está incluida en la Stable API V1.

Las conexiones para un flujo de cambios pueden usar listas de nodos iniciales DNS con la opción de conexión +srv o enumerar los servidores individualmente en la cadena de conexión.

Si el driver pierde la conexión con un flujo de cambios o la conexión se cae, intenta restablecer la conexión con el flujo de cambios a través de otro nodo del clúster que tenga una preferencia de lectura que coincida. Si el driver no puede encontrar un nodo con la preferencia de lectura adecuada, lanzará una excepción.

Para obtener más información, consulta Formato URI de la cadena de conexión.

Puedes abrir flujos de cambio para lo siguiente:

Objetivo
Descripción

Una colección

Puedes abrir un cursor de flujo de cambios para una sola colección (excepto las colecciones system o cualquier colección en las bases de datos admin, local y config).

Los ejemplos de esta página incluyen ejemplos de driver de MongoDB y de mongosh para abrir y trabajar con un cursor de flujo de cambios para una sola colección.

Una base de datos

Puede abrir un cursor de flujo de cambios para una sola base de datos (excluyendo las bases de datos admin, local y config) para observar los cambios en todas tus colecciones no sistemáticas.

Para el método del controlador de MongoDB, consulte la documentación de su controlador. Véase también el mongosh método db.watch().

Una implementación

Puedes abrir un cursor de flujo de cambios para una implementación (ya sea un set de réplicas o un clúster particionado) para observar cambios en todas las colecciones ajenas al sistema en todas las bases de datos, excepto admin, local y config.

Para el método del controlador de MongoDB, consulte la documentación de su controlador. Véase también el mongosh método Mongo.watch().

Si la cantidad de flujos de cambios activos abiertos para una base de datos supera el tamaño del pool de conexiones, puedes experimentar latencia en las notificaciones. Cada flujo de cambio mantiene abierta una conexión con una operación getMore mientras se espera el siguiente evento. Para evitar la latencia, debes asegurarte de que el tamaño del pool sea mayor que el número de flujos de cambio abiertos. Para obtener más detalles, consulta la configuración maxPoolSize.

Cuando se abre un flujo de cambios en un clúster fragmentado:

  • mongos crea flujos de cambios individuales en cada fragmento. Este comportamiento ocurre independientemente de si el flujo de cambios apunta a un rango específico de clave de fragmentación.

  • Cuando mongos recibe los resultados del flujo de cambios, los ordena y filtra. Si es necesario, mongos también realiza una búsqueda de fullDocument.

Para obtener el mejor rendimiento, limita el uso de queries $lookup en los flujos de cambios.

Para abrir un flujo de cambios:

  • Para un set de réplicas, abre el flujo de cambios desde cualquiera de los nodos que contienen datos.

  • Para un clúster fragmentado, abre el flujo de cambios del mongos.

En el siguiente ejemplo, se abre un flujo de cambios para una colección y se recorre el cursor para recuperar los documentos del flujo de cambios. [1]

Itera el cursor del flujo de cambios para recuperar los eventos de cambio. Consulta Eventos de flujo de cambios para el formato del documento del flujo de cambios.

El cursor de flujo de cambios permanece abierto hasta que ocurra una de las siguientes situaciones:

  • El cursor se cierra explícitamente.

  • Se produce un evento de invalidación; por ejemplo, un descarte o cambio de nombre de una colección.

  • La conexión a la implementación de MongoDB se cierra o se agota el tiempo de espera. Consulta Comportamiento para obtener más información.

  • Si la implementación es un clúster particionado, la eliminación de una partición puede provocar que un cursor del flujo de cambios abierto se cierre. Es posible que el cursor cerrado del flujo de cambios no sea completamente reanudable.

Nota

El ciclo de vida de un cursor sin cerrar depende del lenguaje.

[1] Puede especificar un startAtOperationTime para abrir el cursor en un punto en el tiempo determinado. Si el punto de inicio especificado está en el pasado, debe encontrarse dentro del rango de tiempo del OpLog.

Tip

El campo _id del documento de evento de flujo de cambios actúa como token de reanudación. No utilices la pipeline para modificar o eliminar el campo _id del evento de flujo de cambios.

A partir de MongoDB 4.2, los flujos de cambios lanzarán una excepción si la pipeline de agregación del flujo de cambios modifica el campo _id de un evento.

Consulte Eventos de flujo de cambios para obtener más información sobre el formato del documento de respuesta del flujo de cambios.

Por defecto, los flujos de cambios solo devuelven el delta de campos durante la operación de actualización. Sin embargo, puedes configurar el flujo de cambios para que devuelva la versión más reciente confirmada por la mayoría del documento actualizado.

La operación updateLookup lee el documento identificado por su clave de fragmentación e identificador de documento de la colección. La colección se identifica por su nombre y utiliza los datos de la colección tal como existen en el momento en que se procesa el flujo de cambios. Considera estos escenarios:

  • Si se renombra la colección, no se devuelve ningún documento.

  • Si se renombra la colección y se crea una nueva con el nombre anterior, entonces la operación de búsqueda se realiza en la nueva colección. Si se encuentra un documento coincidente, se devuelve.

Advertencia

Para situaciones que involucran eliminaciones rápidas o picos de tráfico, configurar fullDocument: "updateLookup" con un filtro $match puede causar errores de 'Token de reanudación no encontrado'. Esto ocurre cuando la eliminación de un documento hace que el campo fullDocument devuelva un valor nulo, ya que no hay un documento coincidente, lo que impide que el flujo de cambios encuentre el token de reanudación.

En su lugar, utilice imágenes previas y posteriores con fullDocumentBeforeChange: "whenAvailable" y fullDocument: "whenAvailable". Ver la sección Change Streams con imágenes previas y posteriores de documentos.

Nota

Si hay una o más operaciones comprometidas por la mayoría que modificaron el documento actualizado después de la operación de actualización pero antes de la búsqueda, el documento completo devuelto puede ser significativamente diferente del documento en el momento de la operación de actualización.

Sin embargo, los deltas incluidos en el documento de flujo de cambios siempre describen correctamente los cambios de la colección supervisada que se aplicaron a ese evento de flujo de cambios.

El campo fullDocument para un evento de actualización podría faltar si se cumple una de las siguientes condiciones:

  • Si el documento se borra o si la colección se descarta entre la actualización y la búsqueda.

  • Si la actualización cambia los valores de, al menos, uno de los campos de la clave de fragmentación de esa colección.

Consulte Eventos de flujo de cambios para obtener más información sobre el formato del documento de respuesta del flujo de cambios.

Los flujos de cambios pueden reanudarse si se especifica un token de reanudación para resumeAfter o startAfter cuando se abre el cursor.

Advertencia

Cuando reanudes un flujo de cambios con un token de reanudación, utiliza la misma pipeline y las mismas opciones que cuando se generó el token originalmente. Si utilizas un pipeline de flujo de cambios diferente o diferentes opciones, esto podría provocar un comportamiento impredecible, impactar negativamente la coherencia de los datos o impedir que el flujo de cambios se reanude.

Puede reanudar un flujo de cambios después de un evento específico pasando un token de reanudación a resumeAfter al abrir el cursor.

Consulta Tokens de reanudación para obtener más información sobre el token de reanudación.

Importante

  • El oplog debe tener suficiente historial para localizar la operación asociada con el token o la marca de tiempo, si la marca de tiempo está en el pasado.

  • No puedes usar resumeAfter para reanudar un flujo de cambios después de que un evento de invalidación (por ejemplo, un descarte o renombramiento de colección) cierre el flujo. En cambio, puedes usar startAfter para iniciar un nuevo flujo de cambios después de un evento de invalidación.

Puedes iniciar un nuevo flujo de cambios después de un evento específico si pasas un token de reanudación a startAfter cuando se abre el cursor. A diferencia de resumeAfter, startAfter puede reanudar las notificaciones después de un evento de invalidación mediante la creación de un nuevo flujo de cambios.

Consulta Tokens de reanudación para obtener más información sobre el token de reanudación.

Importante

  • El oplog debe tener suficiente historial para localizar la operación asociada con el token o la marca de tiempo, si la marca de tiempo está en el pasado.

Hay dos tipos de tokens de reanudación:

  • Token del evento: Identifica un evento específico de cambio. El cursor del flujo de cambios genera un token de evento cada vez que se produce un evento de cambio.

  • Token de highwatermark: representa un punto en el tiempo sin un evento de cambio asociado. El servidor genera periódicamente tokens de "highwatermark" para indicar que el tiempo del clúster ha avanzado, incluso cuando no se producen eventos de cambio.

Tip

El servidor avanza periódicamente la marca temporal en los tokens de reanudación del highwatermark. En particiones inactivas con guardados poco frecuentes, este avance puede no ocurrir con la suficiente frecuencia para algunos casos de uso. Para avanzar el timestamp de highwatermark con mayor frecuencia, se pueden escribir entradas no operativas en el oplog de los shards inactivos utilizando el comando appendOplogNote.

Puede encontrar tokens de reanudación en múltiples fuentes:

Origen
Descripción

Cada notificación de evento de cambio incluye un token de reanudación en el campo _id.

La etapa de agregación $changeStream incluye un token de reanudación en el campo cursor.postBatchResumeToken.

Este campo solo aparece cuando se utiliza el comando aggregate.

El comando getMore incluye un token de reanudación en el campo cursor.postBatchResumeToken.

A partir de MongoDB 4.2, los flujos de cambios lanzarán una excepción si la pipeline de agregación del flujo de cambios modifica el campo _id de un evento.

Tip

MongoDB proporciona un "fragmento", una extensión de mongosh, que descodifica los tokens de reanudación codificados en hexadecimal.

Puedes instalar y ejecutar el fragmento resumetoken desde mongosh:

snippet install resumetoken
decodeResumeToken('<RESUME TOKEN>')

También puede ejecutar resumetoken desde la línea de comandos (sin utilizar mongosh) si npm está instalado en su sistema:

npx mongodb-resumetoken-decoder <RESUME TOKEN>

Consulte lo siguiente para más detalles sobre:

Las notificaciones de eventos de cambio incluyen un token de reanudación en el campo _id:

{
"_id": {
"_data": "82635019A0000000012B042C0100296E5A1004AB1154ACACD849A48C61756D70D3B21F463C6F7065726174696F6E54797065003C696E736572740046646F63756D656E744B65790046645F69640064635019A078BE67426D7CF4D2000004"
},
"operationType": "insert",
"clusterTime": Timestamp({ "t": 1666193824, "i": 1 }),
"collectionUUID": new UUID("ab1154ac-acd8-49a4-8c61-756d70d3b21f"),
"wallTime": ISODate("2022-10-19T15:37:04.604Z"),
"fullDocument": {
"_id": ObjectId("635019a078be67426d7cf4d2"'),
"name": "Giovanni Verga"
},
"ns": {
"db": "test",
"coll": "names"
},
"documentKey": {
"_id": ObjectId("635019a078be67426d7cf4d2")
}
}

Al utilizar el comando aggregate, la etapa de agregación $changeStream incluye un token de reanudación en el campo cursor.postBatchResumeToken:

{
"cursor": {
"firstBatch": [],
"postBatchResumeToken": {
"_data": "8263515EAC000000022B0429296E1404"
},
"id": Long("4309380460777152828"),
"ns": "test.names"
},
"ok": 1,
"$clusterTime": {
"clusterTime": Timestamp({ "t": 1666277036, "i": 1 }),
"signature": {
"hash": Binary(Buffer.from("0000000000000000000000000000000000000000", "hex"), 0),
"keyId": Long("0")
}
},
"operationTime": Timestamp({ "t": 1666277036, "i": 1 })
}

El comando getMore también incluye un token de reanudación en el campo cursor.postBatchResumeToken:

{
"cursor": {
"nextBatch": [],
"postBatchResumeToken": {
"_data": "8263515979000000022B0429296E1404"
},
"id": Long("7049907285270685005"),
"ns": "test.names"
},
"ok": 1,
"$clusterTime": {
"clusterTime": Timestamp( { "t": 1666275705, "i": 1 } ),
"signature": {
"hash": Binary(Buffer.from("0000000000000000000000000000000000000000", "hex"), 0),
"keyId": Long("0")
}
},
"operationTime": Timestamp({ "t": 1666275705, "i": 1 })
}

Los flujos de cambios pueden beneficiar las arquitecturas con sistemas empresariales confiables, ya que modifican los sistemas posteriores una vez que los cambios en los datos son duraderos. Por ejemplo, los flujos de cambio pueden ahorrar tiempo a los desarrolladores cuando se implementan servicios de extracción, transformación y carga (ETL), sincronización multiplataforma, funcionalidad de colaboración y servicios de notificación.

Para implementaciones que aplican autenticación en implementaciones autogestionadas y autorización:

  • Para abrir un flujo de cambios en una colección específica, las aplicaciones deben tener privilegios que otorguen acceso a las acciones changeStream y find en la colección correspondiente.

    { resource: { db: <dbname>, collection: <collection> }, actions: [ "find", "changeStream" ] }
  • Para abrir un flujo de cambios en una sola base de datos, las aplicaciones deben tener privilegios que otorguen acceso a las acciones changeStream y find en todas las colecciones que no seansystem en la base de datos.

    { resource: { db: <dbname>, collection: "" }, actions: [ "find", "changeStream" ] }
  • Para abrir un flujo de cambios en toda una implementación, las aplicaciones deben tener privilegios que otorguen acceso a las acciones changeStream y find en todas las colecciones que no seansystem para todas las bases de datos de la implementación.

    { resource: { db: "", collection: "" }, actions: [ "find", "changeStream" ] }

Los flujos de cambios solo notifican sobre los cambios de datos que han persistido en la mayoría de los nodos que contienen datos en el set de réplicas. Los flujos de cambios garantizan que las notificaciones se activen solo por cambios comprometidos por la mayoría que sean duraderos en escenarios de fallo.

Por ejemplo, piensa en un set de réplicas de 3 nodos con un cursor de flujo de cambios abierto contra el primario. Si un cliente emite una operación de inserción, el flujo de cambios solo notifica a la aplicación del cambio de datos una vez que esa inserción haya persistido en la mayoría de los nodos portadores de datos.

Si una operación está asociada con una transacción, el documento de evento de cambio incluye el txnNumber y el lsid.

Los flujos de cambio utilizan comparaciones binarias simple a menos que se proporcione una intercalación explícita.

A partir de MongoDB 5.3, durante la migración de rango, los eventos de flujo de cambios no se generan para las actualizaciones de documentos huérfanos.

A partir de MongoDB 6.0, puedes utilizar eventos de flujo de cambios para generar la versión de un documento antes y después de los cambios (las imágenes previas y posteriores del documento):

  • La imagen previa es el documento antes de reemplazarlo, actualizarlo o borrarlo. No existe una imagen previa para un documento insertado.

  • La imagen posterior es el documento tras insertarse, sustituirse o actualizarse. No hay imagen posterior para un documento borrado.

  • Activa changeStreamPreAndPostImages para una colección con db.createCollection(), create o collMod. Por ejemplo, cuando se utiliza el comando collMod:

    db.runCommand( {
    collMod: <collection>,
    changeStreamPreAndPostImages: { enabled: true }
    } )

Las imágenes previas y posteriores no están disponibles para un evento de flujo de cambios si las imágenes fueron:

  • No está habilitado en la colección en el momento de la operación de actualización o eliminación de un documento.

  • Eliminado después del tiempo de retención de imágenes previas y posteriores establecido en expireAfterSeconds.

    • El siguiente ejemplo configura expireAfterSeconds en 100 segundos en todo el clúster:

      use admin
      db.runCommand( {
      setClusterParameter:
      { changeStreamOptions: {
      preAndPostImages: { expireAfterSeconds: 100 }
      } }
      } )

      Nota

      El comando setClusterParameter no es compatible con los clústeres de MongoDB Atlas. Para obtener información sobre el soporte de Atlas para todos los comandos, consulta Comandos no compatibles en Atlas.

    • El siguiente ejemplo devuelve la configuración actual de changeStreamOptions, incluyendo expireAfterSeconds:

      db.adminCommand( { getClusterParameter: "changeStreamOptions" } )
    • Configurar expireAfterSeconds en off utiliza la política de retención por defecto: las imágenes previas y posteriores se conservan hasta que los eventos correspondientes del flujo de cambios se eliminan del oplog.

    • Si se elimina un evento de flujo de cambios del oplog, las imágenes previas y posteriores correspondientes también se eliminan independientemente del tiempo de retención de las imágenes previas y posteriores expireAfterSeconds.

Consideraciones adicionales:

  • Habilitar las imágenes previas y posteriores consume espacio de almacenamiento y aumenta el tiempo de procesamiento. Habilita solo las imágenes previas y de publicación si las necesitas.

  • Limita el tamaño del evento del flujo de cambios a menos de 16 mebibytes. Para limitar el tamaño del evento, puedes:

    • Limita el tamaño del documento a 8 megabytes. Puedes solicitar imágenes previas y posteriores simultáneamente en la salida del flujo de cambios si otros campos de eventos del flujo de cambios como updateDescription no son grandes.

    • Solicita solo imágenes posteriores en la salida del flujo de cambios para documentos de hasta 16 mebibytes si otros campos de eventos del flujo de cambios como updateDescription no son grandes.

    • Solicita solo imágenes previas en la salida del flujo de cambios para documentos de hasta 16 mebibytes si:

      • las actualizaciones de documento afectan solo a una pequeña fracción de la estructura o el contenido del documento, y

      • no cause un evento de cambio replace. Un evento replace siempre incluye la imagen de publicación.

  • Para realizar una solicitud de imagen previa, debes establecer fullDocumentBeforeChange en required o whenAvailable en db.collection.watch(). Para solicitar una imagen posterior, establece fullDocument mediante el mismo método.

  • Las imágenes previas se escriben en la colección config.system.preimages.

    • La colección config.system.preimages puede agrandarse. Para limitar el tamaño de la colección, puedes establecer el tiempo a expireAfterSeconds para las imágenes previas como se mostró antes.

    • Para supervisar el tamaño de config.system.preimages, conectar a un nodo de partición en un clúster particionado o a un nodo mongod en un set de réplicas. A continuación, ejecutar los siguientes comandos:

      use config
      db.system.preimages.totalSize()
      db.system.preimages.stats()

      Nota

      Debe tener la acción de privilegio collStats en la colección config.system.preimages para ejecutar estos comandos.

      Para ejecutar estos comandos en una implementación de MongoDB Atlas, debes tener el rol atlasAdmin.

    • Las imágenes previas se eliminan de forma asincrónica mediante un proceso en segundo plano.

Importante

Característica incompatible con versiones anteriores

A partir de MongoDB 6.0, si utilizas imágenes previas y posteriores de documentos para los flujos de cambios, debes deshabilitar changeStreamPreAndPostImages para cada colección mediante el comando collMod antes de poder volver a una versión anterior de MongoDB.

Tip

Para obtener ejemplos completos con la salida del flujo de cambios, consulta Flujos de cambio con imágenes previas y posteriores de documentos.

Si el parámetro initialSyncMethod para el clúster es fileCopyBased, entonces no hay impacto en los oyentes del flujo de cambios.

Si initialSyncMethod es logical y se abre un flujo de cambios en un nodo recién sincronizado y lee eventos de un punto en el tiempo anterior a la finalización de la sincronización inicial lógica, es posible que falten las imágenes previa y posterior.