Una vista materializada de transmisión es una colección que un procesador de transmisión de Atlas Stream Processing mantiene actualizada. Para lograrlo, un procesador de transmisión lee cada cambio en una colección de origen, calcula el efecto de ese cambio y aplica el resultado a la vista. La vista refleja el estado actual de los datos de origen sin una actualización manual o programada.
Las vistas materializadas de transmisión se adaptan a las cargas de trabajo que leen un resultado calculado con más frecuencia de la que cambian los datos de origen, como tableros, totales acumulados y recuentos de elementos abiertos. La vista se actualiza a medida que llegan los datos, por lo que las lecturas devuelven los resultados actuales con baja latencia.
Comparación con vistas materializadas on-demand
MongoDB admite dos enfoques para las vistas materializadas:
Una vista materializada on-demand almacena el resultado de un pipeline de agregación que se ejecuta manualmente o en un cronograma.
Una vista materializada de transmisión almacena el resultado de un procesador de transmisión de Atlas Stream Processing que se ejecuta continuamente.
Ambos enfoques almacenan los resultados calculados en el disco y sirven las lecturas directamente desde la vista. Difieren en cómo y cuándo se actualiza la vista.
Característica | Vistas materializadas on-demand | Vistas materializadas de transmisión |
|---|---|---|
Actualizar activador | Manual o programado | Continuo, impulsado por el cambio |
latencia | Minutos a días | De subsegundo a bajo segundo |
Actualización de datos | Snapshot punto en el tiempo | Sincronizado permanentemente |
Modelo de cómputo | Vuelve a calcular el resultado completo | Calcula el efecto incremental |
Best fit | Reporte por agrupar, agregación periódica | Tableros en vivo, análisis operativos |
Características de un procesador de vistas materializadas de transmisión
El mantenimiento de una vista materializada de transmisión es un patrón que se ensambla a partir de las etapas de agregación estándar de Atlas Stream Processing. Un procesador de transmisión que implementa este patrón puede tener las siguientes características:
Lee una fuente de flujo de cambios. La etapa
$sourcelee de la colección de origen confullDocumentyfullDocumentBeforeChangeestablecidos enrequiredpara que el pipeline pueda comparar el estado de cada documento antes y después de un cambio.Calcula un delta firmado por evento. Una etapa
$addFieldspuede utilizar una expresión$switchpara asignar un valor positivo o negativo a cada inserción, actualización o eliminación en función de cómo el cambio afecta al resultado calculado.Agrupa los resultados en una ventana. Dado que un procesador de transmisión opera en una transmisión ilimitada, cada etapa
$groupdebe ejecutarse dentro de una etapa de ventana. La ventana$groupsuma los deltas de cada clave durante el intervalo de la ventana. El intervalo de la ventana también actúa como intervalo de actualización de la vista, por lo que determina la actualidad de la vista. El intervalo más pequeño que puede establecer es un milisegundo.Aplica los resultados de forma aditiva. La etapa
$mergepuede usar un pipelinewhenMatchedque agrega el resultado de cada ventana al total acumulado en la vista en lugar de reemplazarlo.Termina en una etapa de sumidero. Un pipeline de procesador de transmisión debe terminar en una etapa de sumidero. Utilice
$mergepara guardar en una colección de Atlas.
Solo las agregaciones incrementales se traducen en una vista materializada de transmisión. Las agregaciones que requieren un escaneo completo de la colección no califican.
Diferencias de comportamiento a considerar
Una vista materializada de transmisión se comporta de manera diferente a una agregación por agrupar de formas que afectan a los consumidores posteriores:
La vista comienza desde cero. Por defecto, el procesador no lee documentos preexistentes, así que siembre la vista antes de iniciar el procesador o habilite
initialSyncen la etapa$source.Solo los cambios de origen impulsan actualizar. Si el pipeline utiliza una etapa
$lookup, los cambios posteriores en la colección de referencia no actualizan los documentos que la vista ya guardó.
Crear una vista materializada de transmisión
El siguiente tutorial crea una vista materializada de transmisión que mantiene un recuento de las ventas completadas por método de compra. El procesador de transmisión lee el flujo de cambios de la colección sample_supplies.sales y guarda en la colección sample_supplies.sales_by_channel.
Nota
Un procesador de flujo también puede leer desde un tema de Apache Kafka y escribir en tablas de Apache Iceberg en AWS S3. Para obtener más información, consulte Apache Kafka Broker y $iceberg Aggregation Stage.
Requisitos previos
Antes de crear el procesador de transmisión, debe tener un espacio de trabajo de Stream Processing con una conexión de Atlas al clúster que contiene los datos de origen. Para agregar una conexión, consulte Gestionar conexiones.
Ejecute los comandos de esta sección en su clúster. Para conectarse, consulte Conectarse a un clúster mediante mongosh.
Complete los siguientes pasos para preparar la colección de origen y sembrar la vista:
Cargue el conjunto de datos de muestra.
Este procedimiento utiliza la colección sales del conjunto de datos sample_supplies. Para aprender cómo cargar datos de muestra, consulte Importar datos de muestra en su implementación de Atlas.
Habilite las imágenes previas y de publicación del flujo de cambios.
Habilite las imágenes previas y de publicación en la colección de origen para que el procesador de transmisión pueda calcular el efecto de cada cambio. Ejecute el siguiente comando en el clúster:
db.getSiblingDB("sample_supplies").runCommand({ collMod: "sales", changeStreamPreAndPostImages: { enabled: true } })
Inicialice la vista con el resultado actual.
Ejecute la siguiente agregación para agrupar en el clúster para rellenar la vista con los recuentos actuales:
db.sales.aggregate([ { $match: { status: "completed" } }, { $group: { _id: "$purchaseMethod", active_count: { $sum: 1 } } }, { $merge: { into: "sales_by_channel", whenMatched: "replace", whenNotMatched: "insert" } } ])
Para confirmar los recuentos iniciales, query la vista. La colección sales_by_channel contiene un documento por método de compra:
db.sales_by_channel.find()
Esta agregación inicial es en sí misma una vista materializada on-demand. Los dos tipos de vista son complementarios: puede inicializar una vista materializada on-demand con una agregación para agrupar y, a continuación, iniciar un procesador de transmisión que mantenga la misma colección actualizada. La vista on-demand se convierte en una vista materializada de transmisión.
Nota
Si su pipeline no utiliza una etapa de ventana, puede inicializar la vista con initialSync en su lugar. En este caso, el procesador de transmisión ingiere primero todos los documentos existentes en la colección de origen como un evento de inserción y, a continuación, procesa los nuevos eventos de cambio. Para aprender sobre la opción $source initialSync, consulte MongoDB colección Change Stream.
Procedimiento
Seleccione la interfaz de usuario de Atlas o mongosh para crear el procesador de transmisión.
Verificar que la vista se mantenga actualizada.
Ejecute los comandos de esta sección en su clúster. Para conectarse, consulte Conectarse a un clúster mediante mongosh.
Después de iniciar el procesador, los cambios en la colección sales actualizan sales_by_channel en segundos. Para confirmar esto, inserte una nueva venta en linea completada:
db.sales.insertOne({ saleDate: new Date(), purchaseMethod: "Online", status: "completed", items: [], customer: {}, couponUsed: false })
El procesador incrementa el recuento Online y registra el límite de la ventana en lastWindowStart:
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1586, lastWindowStart: ISODate('2026-07-23T15:18:01.000Z') }, { _id: 'In store', active_count: 2819 } ]
Cuando un cliente devuelve esa venta, el recuento de Online vuelve a su valor inicial:
db.sales.updateOne( { purchaseMethod: "Online", status: "completed" }, { $set: { status: "returned" } } )
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1585, lastWindowStart: ISODate('2026-07-23T15:18:13.000Z') }, { _id: 'In store', active_count: 2819 } ]
Ejemplos
Estos ejemplos explican el ejemplo smv del repositorio de ejemplos de Atlas Stream Processing. Este ejemplo mantiene una queue_stats colección que almacena un recuento continuo de los tickets de soporte abiertos, con un documento por nivel de prioridad. El procesador de support_tickets flujo lee el flujo de cambios y actualiza cada recuento a medida que se abren, resuelven y eliminan los tickets.
El pipeline ajusta el recuento de la prioridad de un ticket en +1 cuando se abre un ticket y en -1 cuando se resuelve o elimina un ticket. La agregación tiene cinco etapas:
La etapa
$sourcelee el flujo de cambiossupport_ticketscon imágenes previas y de publicación.La etapa
$addFieldsutiliza un$switchenoperationTypepara calcular_deltay extraer_prioritycomo clave de grupo.La etapa
$matchdescarta los eventos con un delta de cero.La etapa
$tumblingWindowsuma los deltas por prioridad dentro de cada ventana de un segundo.La etapa
$mergeañade el delta de cada ventana al total acumulado, utilizandolastWindowStartcomo marca de agua alta para evitar el doble recuento en la reproducción.
[ { "$source": { "connectionName": "<connection-name>", "db": "support", "coll": "support_tickets", "config": { "fullDocument": "required", "fullDocumentBeforeChange": "required" } } }, { "$addFields": { "_delta": { "$switch": { "branches": [ { "case": { "$and": [ { "$eq": ["$operationType", "insert"] }, { "$eq": ["$fullDocument.status", "open"] } ] }, "then": 1 }, { "case": { "$and": [ { "$eq": ["$operationType", "update"] }, { "$eq": ["$fullDocumentBeforeChange.status", "open"] }, { "$eq": ["$fullDocument.status", "resolved"] } ] }, "then": -1 }, { "case": { "$and": [ { "$eq": ["$operationType", "delete"] }, { "$eq": ["$fullDocumentBeforeChange.status", "open"] } ] }, "then": -1 } ], "default": 0 } }, "_priority": { "$ifNull": [ "$fullDocument.priority", "$fullDocumentBeforeChange.priority" ] } } }, { "$match": { "_delta": { "$ne": 0 } } }, { "$tumblingWindow": { "boundary": "processingTime", "interval": { "size": 1, "unit": "second" }, "pipeline": [ { "$group": { "_id": "$_priority", "open_count": { "$sum": "$_delta" }, "windowStart": { "$first": { "$meta": "stream.window.start" } } } } ] } }, { "$merge": { "into": { "connectionName": "<connection-name>", "db": "support", "coll": "queue_stats" }, "whenMatched": [ { "$set": { "open_count": { "$cond": [ { "$gt": ["$$new.windowStart", { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }] }, { "$add": ["$open_count", "$$new.open_count"] }, "$open_count" ] }, "lastWindowStart": { "$max": [ { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }, "$$new.windowStart" ] } } } ], "whenNotMatched": "insert" } } ]
Cada documento de la colección queue_stats se parece a lo siguiente:
{ _id: "P1", open_count: <num>, lastWindowStart: <timestamp> }
Cuando un evento debe ajustar dos claves de grupo, reemplace el escalar _delta con un arreglo _adjustments y distribuya el arreglo en un documento por ajuste. Modifique el pipeline de clave única de la siguiente manera:
Reemplace la etapa
$addFieldspara que cada rama$switchdevuelva un arreglo_adjustmentscon un elemento por clave de grupo afectada. Una escalada devuelve dos elementos.Agregue una etapa
$unwindy una etapa$setdespués de$addFields. La etapa$unwindrealiza la división de cada evento en un documento por ajuste y descarta los arreglos vacíos, lo que reemplaza la etapa$match. La etapa$setpromueve los campos_adjustmentsa_priorityy_deltade nivel superior para que las etapas restantes funcionen sin cambios.
Reemplace las etapas 2 y 3 del pipeline de clave única con lo siguiente:
// Stage 2 (replacement): Compute an _adjustments array. { $addFields: { _adjustments: { $switch: { branches: [ { case: { $and: [ { $eq: ["$operationType", "insert"] }, { $eq: ["$fullDocument.status", "open"] } ]}, then: [{ _priority: "$fullDocument.priority", _delta: 1 }] }, { case: { $and: [ { $eq: ["$operationType", "update"] }, { $eq: ["$fullDocumentBeforeChange.status", "open"] }, { $eq: ["$fullDocument.status", "resolved"] } ]}, then: [{ _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }] }, { case: { $and: [ { $eq: ["$operationType", "update"] }, { $eq: ["$fullDocument.status", "open"] }, { $ne: ["$fullDocument.priority", "$fullDocumentBeforeChange.priority"] } ]}, then: [ { _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }, { _priority: "$fullDocument.priority", _delta: 1 } ] }, { case: { $and: [ { $eq: ["$operationType", "delete"] }, { $eq: ["$fullDocumentBeforeChange.status", "open"] } ]}, then: [{ _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }] } ], default: [] } } } }, // Stage 3 (replacement): Fan out into one document per adjustment, // then lift the adjustment fields back to the top level. { $unwind: "$_adjustments" }, { $set: { _priority: "$_adjustments._priority", _delta: "$_adjustments._delta" } }
Información Adicional
Para aprender más sobre las etapas y los conceptos de esta guía, consulte los siguientes recursos:
Para aprender a crear, iniciar, detener y supervisar procesadores de transmisión, consulte Desarrollar y gestionar procesadores de transmisión.
Para aprender sobre las etapas de agregación que admite Atlas Stream Processing, consulta etapas del pipeline de agregación.
Para aprender sobre los orígenes que puede leer un procesador de transmisión, consulte
$sourceStage (Stream Processing).Para aprender sobre las etapas de ventana, consulte Windows del procesador de transmisión.
Para aprender más sobre las vistas materializadas on-demand, consulte Vistas materializadas On-demand.