Para agentes de IA: um índice de documentação está disponível em https://www.mongodb.com/pt-br/docs/llms.txt — as versões de markdown de todas as páginas estão disponíveis anexando .md a qualquer caminho de URL.
Menu Docs

Crie visualizações materializadas de streaming

Uma visualização materializada de streaming é uma coleção que um processador de fluxo de Atlas Stream Processing mantém atual. Para conseguir isso, um processador de fluxo lê cada alteração em uma coleção de origem, calcula o efeito dessa alteração e aplica o resultado à visualização. A visualização reflete o estado atual dos dados de origem sem uma atualização manual ou agendada.

As visualizações materializadas de streaming são adequadas para cargas de trabalho que leem um resultado computado com mais frequência do que os dados de origem são alterados, como dashboards, totais em execução e contagens de itens abertos. A visualização é atualizada à medida que os dados chegam, para que as leituras retornem os resultados atuais com baixa latência.

O MongoDB oferece suporte a duas abordagens para visualizações materializadas:

  • Uma visualização materializada sob demanda armazena o resultado de um pipeline de agregação que você executa manualmente ou em um agendamento.

  • Uma visualização materializada de streaming armazena o resultado de um processador de streaming do Atlas Stream Processing em execução contínua.

Ambas as abordagens armazenam resultados computados em disco e servem leituras diretamente da visualização. Eles diferem em como e quando a visualização é atualizada.

Característica
Visualizações materializadas sob demanda
Visualizações materializadas de streaming

Atualizar trigger

Manual ou agendado

Contínuo, orientado a mudanças

Latência

Minutos para dias

Subsegundo a baixo segundo

Atualização de dados

Snapshot de ponto

Sincronizado perpetuamente

Modelo de computação

Recalcula o resultado completo

Calcula o efeito incremental

Best fit

Relatório de lote, agregação periódica

Dashboards ao vivo, análise operacional

Manter uma visualização materializada de streaming é um padrão que você monta a partir de estágios de agregação padrão do Atlas Stream Processing. Um processador de fluxo que implementa esse padrão pode ter as seguintes características:

  • Lê uma fonte de fluxo de alterações. O estágio $source lê da coleção de origem com fullDocument e fullDocumentBeforeChange definidos como required para que o pipeline possa comparar o estado de cada documento antes e depois de uma alteração.

  • Calcula um delta assinado por evento. Um estágio $addFields pode usar uma expressão $switch para atribuir um valor positivo ou negativo a cada inserção, atualização ou exclusão com base em como a alteração afeta o resultado calculado.

  • Agrupa os resultados em uma janela. Como um processador de fluxo opera em um fluxo ilimitado, cada estágio $group deve ser executado dentro de um estágio de janela. A janela $group soma os deltas para cada chave sobre o intervalo da janela. O intervalo da janela também atua como um intervalo de atualização para a visualização, portanto, determina o quão atualizada a visualização está. O menor intervalo que você pode definir é um milissegundo.

  • Aplica os resultados de forma aditiva. O estágio $merge pode usar um pipeline whenMatched que adiciona o resultado de cada janela ao total em execução na visualização em vez de substituí-lo.

  • Termina em um estágio de sink. Um pipeline de processador de streaming deve terminar em um estágio de sink. Use $merge para gravar em uma coleção do Atlas.

Somente agregações incrementais se traduzem em uma visualização materializada de streaming. Agregações que exigem uma varredura de coleção completa não se qualificam.

Uma visualização materializada de streaming se comporta de forma diferente de uma agregação em lote de maneiras que afetam os consumidores downstream:

  • A visualização começa do zero. Por padrão, o processador não lê documentos pré-existentes, portanto, preencha a visualização antes de iniciar o processador ou ative initialSync no estágio $source.

  • Somente alterações de origem impulsionam as atualizações. Se o pipeline usar um estágio $lookup, alterações posteriores na coleção de referência não atualizarão os documentos que a exibição já gravou.

O tutorial a seguir cria uma visualização materializada de streaming que mantém uma contagem de vendas concluídas por método de compra. O processador de streaming lê o fluxo de alterações da coleção sample_supplies.sales e grava na coleção sample_supplies.sales_by_channel.

Observação

Um processador de fluxo também pode ler de um tópico do Apache Kafka e escrever em tabelas do Apache Geophyllum no AWS S3. Para saber mais, consulte Apache Kafka Corretor e $iceberg Estágio de agregação.

Antes de criar o processador de fluxo, você deve ter um espaço de trabalho de processamento de fluxo com uma conexão Atlas ao cluster que contém os dados de origem. Para adicionar uma conexão, consulte Gerenciar conexões.

Execute os comandos nesta seção em seu cluster. Para se conectar, consulte Conectar-se a um cluster via mongosh.

Conclua as seguintes etapas para preparar a coleção de origem e semear a visualização:

1

Este procedimento utiliza a coleção sales do conjunto de dados sample_supplies. Para aprender como carregar dados de amostra, consulte Importar dados de amostra para sua implantação do Atlas.

2

Adicione um campo status para que o processador de fluxo possa detectar vendas concluídas e devolvidas. Execute o seguinte comando no cluster:

db.sales.updateMany(
{ status: { $exists: false } },
{ $set: { status: "completed" } }
)
3

Habilite pré-imagens e imagens de publicação na coleção de origem para que o processador de streaming possa calcular o efeito de cada alteração. Execute o seguinte comando no cluster:

db.getSiblingDB("sample_supplies").runCommand({
collMod: "sales",
changeStreamPreAndPostImages: { enabled: true }
})
4

Execute a seguinte agregação de lote no cluster para preencher a visualização com as contagens atuais:

db.sales.aggregate([
{ $match: { status: "completed" } },
{ $group: { _id: "$purchaseMethod", active_count: { $sum: 1 } } },
{ $merge: {
into: "sales_by_channel",
whenMatched: "replace",
whenNotMatched: "insert"
} }
])

Para confirmar as contagens semeadas, query a visualização. A coleção sales_by_channel contém um documento por método de compra:

db.sales_by_channel.find()

Essa agregação inicial é, em si, uma visualização materializada sob demanda. Os dois tipos de visualização são complementares: você pode inicializar uma visualização materializada sob demanda com uma agregação em lote e, em seguida, iniciar um processador de streaming que mantém a mesma coleção atual. A visualização sob demanda se torna uma visualização materializada de streaming.

Observação

Se o seu pipeline não usar um estágio de janela, você poderá semear a visualização com initialSync em vez disso. Nesse caso, o processador de fluxo primeiro ingere todos os documentos existentes na coleção de origem como um evento de inserção e, em seguida, processa novos eventos de alteração. Para aprender sobre a opção $source initialSync, consulte MongoDB coleção Change Stream.

Selecione a IU do Atlas ou mongosh para criar o processador de fluxo.

Execute os comandos nesta seção em seu cluster. Para se conectar, consulte Conectar-se a um cluster via mongosh.

Depois de iniciar o processador, as alterações na coleção sales atualizam sales_by_channel em segundos. Para confirmar isso, insira uma nova venda online concluída:

db.sales.insertOne({
saleDate: new Date(),
purchaseMethod: "Online",
status: "completed",
items: [],
customer: {},
couponUsed: false
})

O processador incrementa a contagem Online e registra o limite da janela em 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 }
]

Quando um cliente retorna essa venda, a contagem de Online retorna ao seu 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 }
]

Esses exemplos percorrem o exemplo smv no repositório de exemplo do Atlas Stream Processing . Ele mantém uma queue_stats coleção que contém uma contagem contínua de tickets de suporte abertos, com um documento por nível de prioridade . O processador de stream lê o support_tickets change stream e atualiza cada contagem à medida que os tíquetes são abertos, resolvidos e excluídos.

Calcule um delta assinado por evento e adicione-o à contagem de cada prioridade.

O pipeline ajusta a contagem da prioridade de um ticket em +1 quando um ticket é aberto e em -1 quando um ticket é resolvido ou excluído. A agregação tem cinco estágios:

  1. O estágio $source lê o fluxo de alterações support_tickets com pré-imagens e imagens de publicação.

  2. O estágio $addFields usa um $switch em operationType para calcular _delta e extrair _priority como a chave do grupo.

  3. O estágio $match descarta eventos com um delta de zero.

  4. O estágio $tumblingWindow soma os deltas por prioridade em cada janela de um segundo.

  5. O estágio $merge adiciona o delta de cada janela ao total em execução, usando lastWindowStart como uma marca de água alta para evitar contagem dupla na reprodução.

[
{
"$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 na coleção queue_stats é semelhante ao seguinte:

{ _id: "P1", open_count: <num>, lastWindowStart: <timestamp> }

Distribua um único evento em um ajuste por chave de grupo afetada.

Quando um evento deve ajustar duas chaves de grupo, substitua o escalar _delta por um array _adjustments e distribua o array em um documento por ajuste. Modifique o pipeline de chave única da seguinte forma:

  1. Substitua o estágio $addFields para que cada ramificação $switch retorne um array _adjustments com um elemento por chave de grupo afetada. Uma escalada retorna dois elementos.

  2. Adicione um estágio $unwind e um estágio $set após $addFields. O estágio $unwind faz a divisão de cada evento em um documento por ajuste e descarta arrays vazios, o que substitui o estágio $match. [[ ## completed ## ]񎟟 O estágio $set promove os campos _adjustments para _priority e _delta de nível superior para que os estágios restantes funcionem inalterados.

Substitua os estágios 2 e 3 do pipeline de chave única pelo seguinte:

// 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 mais sobre os estágios e conceitos neste guia, consulte os seguintes recursos: