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 procedimento 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 Apache Kafka e gravar em tabelas Apache Iceberg no Amazon Web Services S3. Para aprender mais, consulte Apache Kafka Broker e $iceberg agregação Stage.

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.

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" } }
)
{
acknowledged: true,
insertedId: null,
matchedCount: 5000,
modifiedCount: 5000,
upsertedCount: 0
}
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 }
})
{ ok: 1, ... }
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()
[
{ _id: 'Phone', active_count: 596 },
{ _id: 'Online', active_count: 1585 },
{ _id: 'In store', active_count: 2819 }
]

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.

1
  1. Na IU do Atlas, Go para a página Stream Processing do seu projeto Atlas.

  2. Clique em Manage no painel do seu espaço de trabalho de Stream Processing que contém sua conexão Atlas com o cluster de origem.

2
  1. Clique em Create stream processor.

  2. Selecione o Visual Builder.

3
  1. No campo Source, selecione sua conexão Atlas com o cluster de origem na lista suspensa Connection.

  2. Na caixa de texto JSON, configure o estágio $source para ler a coleção sales com pré- e publicação-imagens:

{
"$source": {
"connectionName": "<connection-name>",
"db": "sample_supplies",
"coll": "sales",
"config": {
"fullDocument": "required",
"fullDocumentBeforeChange": "required"
}
}
}
4
  1. No painel Start building your pipeline, clique em + Custom stage.

  2. Na caixa de texto JSON, adicione um estágio $addFields que atribui um delta assinado a cada inserção, atualizar ou excluir e captura o 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. Clique em + e selecione Custom stage.

  2. Na caixa de texto JSON, adicione um estágio $match que remove eventos com um delta de zero:

{
"$match": { "_delta": { "$ne": 0 } }
}
6
  1. Clique em + e selecione Custom stage.

  2. Na caixa de texto JSON, adicione um estágio $tumblingWindow que soma os deltas por método de compra em intervalos de um segundo:

{
"$tumblingWindow": {
"boundary": "processingTime",
"interval": { "size": 1, "unit": "second" },
"pipeline": [
{
"$group": {
"_id": "$_channel",
"active_count": { "$sum": "$_delta" },
"windowStart": {
"$first": { "$meta": "stream.window.start" }
}
}
}
]
}
}

Importante

Um processador de fluxo exige que cada estágio $group seja executado dentro de um estágio de janela.

7
  1. No campo Sink, selecione sua conexão do Atlas na lista suspensa Connection.

  2. Na caixa de texto JSON, configure o estágio $merge para adicionar o resultado de cada janela ao total em execução em 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"
}
}

Observação

A marca d'água alta lastWindowStart impede que uma janela reproduzida conte duas vezes

8
  1. No campo Stream processor name, digite sales_stats_sp.

  2. Selecione o tier para o processador de fluxo. Para escolher um tier para sua carga de trabalho, consulte Guia de seleção do tier de Atlas Stream Processing.

9

Clique em Create stream processor.

10

Na aba Stream Processors, selecione sales_stats_sp e clique em Start.

O processador agora mantém sales_by_channel continuamente. Para aprender mais sobre como iniciar, parar e fazer monitoramento de processadores de fluxo, consulte Desenvolver processadores de fluxo.

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 de smv no repositório de exemplos do Atlas Stream Processing. Ele mantém uma coleção queue_stats que contém uma contagem contínua de tickets de suporte abertos, com um documento por nível de prioridade. O processador de streaming lê o fluxo de alterações support_tickets e atualiza cada contagem à medida que os tickets 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: