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.
Disponibilidad
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 demajoritypuede 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.
Compatibilidad con Stable API
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.
Conectar
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.
Supervisa una colección, base de datos o implementació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 Los ejemplos de esta página incluyen ejemplos de driver de MongoDB y de |
Una base de datos | Puede abrir un cursor de flujo de cambios para una sola base de datos (excluyendo las bases de datos Para el método del controlador de MongoDB, consulte la documentación de su controlador. Véase también el |
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 Para el método del controlador de MongoDB, consulte la documentación de su controlador. Véase también el |
Consideraciones sobre el rendimiento de los flujos de cambios
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.
Consideraciones sobre clústeres fragmentados
Cuando se abre un flujo de cambios en un clúster fragmentado:
mongoscrea 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
mongosrecibe los resultados del flujo de cambios, los ordena y filtra. Si es necesario,mongostambién realiza una búsqueda defullDocument.
Para obtener el mejor rendimiento, limita el uso de queries $lookup en los flujos de cambios.
Abre un flujo 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. |
Modifica la salida del flujo de cambios
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.
Revisa un documento completo para operaciones de actualización
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.
Reanudar un 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.
resumeAfter para Change Streams
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
resumeAfterpara 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.
startAfter para Change Streams
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.
Tokens de reanudación
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 | |
La etapa de agregación Este campo solo aparece cuando se utiliza el comando | |
El comando |
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:
Tokens de reanudación de eventos de cambios
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") } }
Reanudar tokens desde aggregate
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 }) }
Reanudar tokens desde getMore
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 }) }
Casos de uso
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.
Control de acceso
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
changeStreamyfinden 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
changeStreamyfinden todas las colecciones que no seansystemen 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
changeStreamyfinden todas las colecciones que no seansystempara todas las bases de datos de la implementación.{ resource: { db: "", collection: "" }, actions: [ "find", "changeStream" ] }
Notificación de evento
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.
Intercalación
Los flujos de cambio utilizan comparaciones binarias simple a menos que se proporcione una intercalación explícita.
Change Streams y documentos huérfanos
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.
Change Streams con imágenes previas y posteriores de los documentos
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
changeStreamPreAndPostImagespara una colección condb.createCollection(),createocollMod. Por ejemplo, cuando se utiliza el comandocollMod: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
expireAfterSecondsen100segundos en todo el clúster:use admin db.runCommand( { setClusterParameter: { changeStreamOptions: { preAndPostImages: { expireAfterSeconds: 100 } } } } ) Nota
El comando
setClusterParameterno 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, incluyendoexpireAfterSeconds:db.adminCommand( { getClusterParameter: "changeStreamOptions" } ) Configurar
expireAfterSecondsenoffutiliza 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
updateDescriptionno 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
updateDescriptionno 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 eventoreplacesiempre incluye la imagen de publicación.
Para realizar una solicitud de imagen previa, debes establecer
fullDocumentBeforeChangeenrequiredowhenAvailableendb.collection.watch(). Para solicitar una imagen posterior, establecefullDocumentmediante el mismo método.Las imágenes previas se escriben en la colección
config.system.preimages.La colección
config.system.preimagespuede agrandarse. Para limitar el tamaño de la colección, puedes establecer el tiempo aexpireAfterSecondspara 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 nodomongoden 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
collStatsen la colecciónconfig.system.preimagespara 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 los eventos y resultados del flujo de cambios, consulta Eventos de cambio.
Para buscar cambios en una colección, consulta
db.collection.watch().Para obtener ejemplos completos con la salida del flujo de cambios, consulta Flujos de cambio con imágenes previas y posteriores de documentos.
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.