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ó.
Cree una vista materializada de transmisión en la interfaz de usuario de Atlas
El siguiente procedimiento 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 transmisión también puede leer de un tema de Apache Kafka y guardar en tablas de Apache Iceberg en AWS S3. Para aprender más, consulte Apache Kafka Broker y $iceberg agregación 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.
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.
Agregue un campo de estado a la colección de ventas.
Agregue un campo status para que el procesador de transmisión pueda detectar las ventas completadas y devueltas. Ejecute el siguiente comando en el clúster:
db.sales.updateMany( { status: { $exists: false } }, { $set: { status: "completed" } } )
{ acknowledged: true, insertedId: null, matchedCount: 5000, modifiedCount: 5000, upsertedCount: 0 }
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 } })
{ ok: 1, ... }
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()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1585 }, { _id: 'In store', active_count: 2819 } ]
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
Configure la fuente.
En el campo Source, seleccione su conexión de Atlas al clúster de origen en la lista desplegable Connection.
En el cuadro de texto JSON, configure la etapa
$sourcepara leer la colecciónsalescon imágenes previas y de publicación:
{ "$source": { "connectionName": "<connection-name>", "db": "sample_supplies", "coll": "sales", "config": { "fullDocument": "required", "fullDocumentBeforeChange": "required" } } }
Agregue una etapa que calcule el delta de cambio para cada evento.
En el panel Start building your pipeline, haga clic en + Custom stage.
En el cuadro de texto JSON, agregue una etapa
$addFieldsque asigne un delta firmado a cada inserción, actualización o eliminación y capture el método de compra:
{ "$addFields": { "_delta": { "$switch": { "branches": [ { "case": { "$and": [ { "$eq": ["$operationType", "insert"] }, { "$eq": ["$fullDocument.status", "completed"] } ] }, "then": 1 }, { "case": { "$and": [ { "$eq": ["$operationType", "update"] }, { "$eq": ["$fullDocumentBeforeChange.status", "completed"] }, { "$eq": ["$fullDocument.status", "returned"] } ] }, "then": -1 }, { "case": { "$and": [ { "$eq": ["$operationType", "delete"] }, { "$eq": ["$fullDocumentBeforeChange.status", "completed"] } ] }, "then": -1 } ], "default": 0 } }, "_channel": { "$ifNull": [ "$fullDocument.purchaseMethod", "$fullDocumentBeforeChange.purchaseMethod" ] } } }
Agregue una etapa que descarte los eventos sin efecto.
Haga clic en + y, a continuación, seleccione Custom stage.
En el cuadro de texto JSON, agregue una etapa
$matchque remueva los eventos con un delta de cero:
{ "$match": { "_delta": { "$ne": 0 } } }
Agregue una etapa de grupo con ventana.
Haga clic en + y, a continuación, seleccione Custom stage.
En el cuadro de texto JSON, agregue una etapa
$tumblingWindowque sume los deltas por método de compra en intervalos de un segundo:
{ "$tumblingWindow": { "boundary": "processingTime", "interval": { "size": 1, "unit": "second" }, "pipeline": [ { "$group": { "_id": "$_channel", "active_count": { "$sum": "$_delta" }, "windowStart": { "$first": { "$meta": "stream.window.start" } } } } ] } }
Importante
Un procesador de transmisión requiere que cada etapa $group se ejecute dentro de una etapa de ventana.
Configure el sumidero.
En el campo Sink, seleccione su conexión de Atlas en la lista de descartar Connection.
En el cuadro de texto JSON, configure la etapa
$mergepara agregar el resultado de cada ventana al total acumulado ensales_by_channel:
{ "$merge": { "into": { "connectionName": "<connection-name>", "db": "sample_supplies", "coll": "sales_by_channel" }, "whenMatched": [ { "$set": { "active_count": { "$cond": [ { "$gt": ["$$new.windowStart", { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }] }, { "$add": ["$active_count", "$$new.active_count"] }, "$active_count" ] }, "lastWindowStart": { "$max": [ { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }, "$$new.windowStart" ] } } } ], "whenNotMatched": "insert" } }
Nota
La marca de agua alta lastWindowStart evita que una ventana reproducida cuente dos veces
Introduzca los detalles del procesador.
En el campo Stream processor name, introduce
sales_stats_sp.Seleccione el nivel para el procesador de transmisión. Para elegir un nivel para su carga de trabajo, consulte la Guía de selección de nivel de Atlas Stream Processing.
Inicia el procesador de flujos.
En la pestaña Stream Processors, seleccione sales_stats_sp y haga clic en Start.
El procesador ahora mantiene sales_by_channel continuamente. Para aprender más sobre cómo iniciar, detener y realizar la supervisión de los procesadores de transmisión, consulte Desarrollar procesadores de transmisión.
Verificar que la vista se mantenga actualizada.
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 recorren el ejemplo de smv en el repositorio de ejemplos de Atlas Stream Processing. Mantiene una colección queue_stats que contiene un recuento continuo de tickets de soporte abiertos, con un documento por nivel de prioridad. El procesador de flujo lee el flujo de cambios support_tickets y actualiza cada recuento a medida que se abren, resuelven y borran 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 cómo crear, iniciar, detener y supervisar procesadores de transmisión, consulte Desarrollar 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.