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.
Comparação com visualizações materializadas on-demand
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 |
Características de um processador de visualizações materializadas de streaming
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
$sourcelê da coleção de origem comfullDocumentefullDocumentBeforeChangedefinidos comorequiredpara 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
$addFieldspode usar uma expressão$switchpara 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
$groupdeve ser executado dentro de um estágio de janela. A janela$groupsoma 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
$mergepode usar um pipelinewhenMatchedque 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
$mergepara 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.
Diferenças de comportamento a serem consideradas
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
initialSyncno 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.
Crie uma visualização materializada de streaming na IU do Atlas
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.
Pré-requisitos
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:
Carregue o conjunto de dados de amostra.
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.
Adicione um campo de status à coleção de vendas.
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 }
Ative a pré e publicação-imagens do fluxo de alterações.
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, ... }
Preencha a visualização com o resultado atual.
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.
Procedimento
Configurar a origem.
No campo Source, selecione sua conexão Atlas com o cluster de origem na lista suspensa Connection.
Na caixa de texto JSON, configure o estágio
$sourcepara ler a coleçãosalescom pré- e publicação-imagens:
{ "$source": { "connectionName": "<connection-name>", "db": "sample_supplies", "coll": "sales", "config": { "fullDocument": "required", "fullDocumentBeforeChange": "required" } } }
Adicione um estágio que calcule o delta de alterações para cada evento.
No painel Start building your pipeline, clique em + Custom stage.
Na caixa de texto JSON, adicione um estágio
$addFieldsque 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" ] } } }
Adicione um estágio que descarte eventos sem efeito.
Clique em + e selecione Custom stage.
Na caixa de texto JSON, adicione um estágio
$matchque remove eventos com um delta de zero:
{ "$match": { "_delta": { "$ne": 0 } } }
Adicione um estágio de grupo de janela.
Clique em + e selecione Custom stage.
Na caixa de texto JSON, adicione um estágio
$tumblingWindowque 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.
Configurar o coletor.
No campo Sink, selecione sua conexão do Atlas na lista suspensa Connection.
Na caixa de texto JSON, configure o estágio
$mergepara adicionar o resultado de cada janela ao total em execução emsales_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
Insira os detalhes do processador.
No campo Stream processor name, digite
sales_stats_sp.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.
Inicie o processador de fluxo.
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.
Verifique se a visualização permanece atual
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 } ]
Exemplos
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.
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:
O estágio
$sourcelê o fluxo de alteraçõessupport_ticketscom pré-imagens e imagens de publicação.O estágio
$addFieldsusa um$switchemoperationTypepara calcular_deltae extrair_prioritycomo a chave do grupo.O estágio
$matchdescarta eventos com um delta de zero.O estágio
$tumblingWindowsoma os deltas por prioridade em cada janela de um segundo.O estágio
$mergeadiciona o delta de cada janela ao total em execução, usandolastWindowStartcomo 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> }
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:
Substitua o estágio
$addFieldspara que cada ramificação$switchretorne um array_adjustmentscom um elemento por chave de grupo afetada. Uma escalada retorna dois elementos.Adicione um estágio
$unwinde um estágio$setapós$addFields. O estágio$unwindfaz 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$setpromove os campos_adjustmentspara_prioritye_deltade 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" } }
Informações adicionais
Para aprender mais sobre os estágios e conceitos neste guia, consulte os seguintes recursos:
Para aprender como criar, iniciar, parar e monitorar processadores de fluxo, consulte Desenvolver processadores de fluxo.
Para aprender sobre os estágios de agregação que o Atlas Stream Processing suporta, consulte fases do pipeline de agregação.
Para aprender sobre as fontes que um processador de fluxo pode ler, consulte Estágio
$source(Stream Processing).Para aprender sobre os estágios da janela, consulte Windows do processador de streaming.
Para saber mais sobre visualizações materializadas sob demanda, consulte Visualizações materializadas sob demanda.