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

Compilar vistas materializadas de transmisión

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.

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

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 $source lee de la colección de origen con fullDocument y fullDocumentBeforeChange establecidos en required para 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 $addFields puede utilizar una expresión $switch para 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 $group debe ejecutarse dentro de una etapa de ventana. La ventana $group suma 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 $merge puede usar un pipeline whenMatched que 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 $merge para 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.

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 initialSync en 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ó.

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.

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:

1

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.

2

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
}
3

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, ... }
4

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.

1
  1. En la Interfaz de Usuario de Atlas, Go a la página Stream Processing de su Proyecto de Atlas.

  2. Haga clic en Manage en el panel de su espacio de trabajo de Stream Processing que contiene su conexión de Atlas al clúster de origen.

2
  1. Haga clic en Create stream processor.

  2. Selecciona el Visual Builder.

3
  1. En el campo Source, seleccione su conexión de Atlas al clúster de origen en la lista desplegable Connection.

  2. En el cuadro de texto JSON, configure la etapa $source para leer la colección sales con imágenes previas y de publicación:

{
"$source": {
"connectionName": "<connection-name>",
"db": "sample_supplies",
"coll": "sales",
"config": {
"fullDocument": "required",
"fullDocumentBeforeChange": "required"
}
}
}
4
  1. En el panel Start building your pipeline, haga clic en + Custom stage.

  2. En el cuadro de texto JSON, agregue una etapa $addFields que 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"
]
}
}
}
5
  1. Haga clic en + y, a continuación, seleccione Custom stage.

  2. En el cuadro de texto JSON, agregue una etapa $match que remueva los eventos con un delta de cero:

{
"$match": { "_delta": { "$ne": 0 } }
}
6
  1. Haga clic en + y, a continuación, seleccione Custom stage.

  2. En el cuadro de texto JSON, agregue una etapa $tumblingWindow que 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.

7
  1. En el campo Sink, seleccione su conexión de Atlas en la lista de descartar Connection.

  2. En el cuadro de texto JSON, configure la etapa $merge para agregar el resultado de cada ventana al total acumulado en sales_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

8
  1. En el campo Stream processor name, introduce sales_stats_sp.

  2. 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.

9

Haga clic en Create stream processor.

10

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.

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 }
]

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.

Calcule un delta firmado por evento y agréguelo al recuento de cada prioridad.

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:

  1. La etapa $source lee el flujo de cambios support_tickets con imágenes previas y de publicación.

  2. La etapa $addFields utiliza un $switch en operationType para calcular _delta y extraer _priority como clave de grupo.

  3. La etapa $match descarta los eventos con un delta de cero.

  4. La etapa $tumblingWindow suma los deltas por prioridad dentro de cada ventana de un segundo.

  5. La etapa $merge añade el delta de cada ventana al total acumulado, utilizando lastWindowStart como 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> }

Distribuya un solo evento en un ajuste por clave de grupo afectada.

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:

  1. Reemplace la etapa $addFields para que cada rama $switch devuelva un arreglo _adjustments con un elemento por clave de grupo afectada. Una escalada devuelve dos elementos.

  2. Agregue una etapa $unwind y una etapa $set después de $addFields. La etapa $unwind realiza 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 $set promueve los campos _adjustments a _priority y _delta de 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"
}
}

Para aprender más sobre las etapas y los conceptos de esta guía, consulte los siguientes recursos: