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.
Criar uma visualização materializada de streaming
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.
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.
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:
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.
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 } })
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()
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
Selecione a IU do Atlas ou mongosh para criar o processador de fluxo.
Verifique se a visualização permanece atual
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 } ]
Exemplos
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.
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 saber como criar, iniciar, parar e monitorar processadores de fluxo, consulte Desenvolver e gerenciar 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.