MongoDB: Fluxos de alterações

O Change Streams é a base para as notificações de alterações em tempo real do MongoDB — CDC (Change Data Capture) — e para aplicativos orientados a eventos.

1. O que você vai aprender


100%
graph LR
    A[Applications 1<br/>Insert Order] -->|oplog| DB[(MongoDB<br/>oplog)]
    B[Applications 2<br/>Update Order] -->|oplog| DB
    C[Applications 3<br/>Delete Order] -->|oplog| DB

    DB -->|Change Stream| CS[db.collection.watch]

    CS -->|operationType=insert| N1[Email Notification<br/>WebSocket Push]
    CS -->|operationType=update| N2[Elasticsearch<br/>Synchronize]
    CS -->|operationType=delete| N3[Redis Cache<br/>Invalidation]

    style CS fill:#d4edda
    style N1 fill:#cce5ff

2. Noções básicas sobre Change Streams

Visão geral do conceito: O Change Streams é uma API de monitoramento de alterações em tempo real introduzida no MongoDB 3.6 e versões posteriores, implementada por meio do oplog do conjunto de réplicas. Ela permite que as aplicações detectem alterações nos dados em tempo real sem a necessidade de polling — servindo como base para o CDC (Change Data Capture) e para arquiteturas orientadas a eventos.

Princípios do Oplog: O Oplog (Registro de Operações) é uma coleção especial com limite de tamanho que registra todas as operações de gravação em um conjunto de réplicas. O primário acrescenta cada operação de gravação (inserção/atualização/exclusão) ao Oplog, e o secundário replica os dados acompanhando o Oplog. Os Change Streams são, essencialmente, um encapsulamento estruturado do Oplog — eles convertem entradas brutas do Oplog em eventos de alteração legíveis por humanos e oferecem recursos avançados, como filtragem e a funcionalidade de retomada a partir de um ponto de interrupção.

100%
graph TB
    subgraph "OpLog How It Works"
        W1[Write Operation] --> OP[(oplog<br/>Fixed Set)]
        OP --> S1[Secondary 1<br/>tailing oplog]
        OP --> S2[Secondary 2<br/>tailing oplog]
        OP --> CS[Change Stream<br/>Structured Events]
    end

    subgraph "Change Stream vs Polling"
        P1[Polling: Queries per second] --> C1[High latency<br/>High Load<br/>Duplicate Query]
        CS --> C2[Real-time Notifications<br/>Low load<br/>No duplicates]
    end

    style CS fill:#d4edda
    style C2 fill:#d4edda
    style C1 fill:#f8d7da

Fluxos de alterações x Polling: uma comparação:

Dimensão Pesquisa Fluxos de alterações
Atraso 1–60 segundos (dependendo do intervalo) < 100 ms (notificação em tempo real)
Carga do banco de dados Alta (varredura de índice em cada consulta) Baixa (recebimento passivo de eventos do oplog)
Integridade do evento Possível perda de eventos (várias alterações dentro de um intervalo) Sem perda (o oplog está ordenado)
Download da retomada Deve ser implementado manualmente resumeToken integrado
Uso de recursos Uso contínuo da conexão + CPU Conexões de longa duração, baixo uso da CPU

Alterar a estrutura do evento:

JAVASCRIPT
// === Monitor Set Changes ===
const changeStream = db.products.watch();

changeStream.on('change', (change) => {
  console.log('Change detected:', change);
  // {
  //   _id: { _data: '...' },           // resumeToken(For resuming downloads)
  //   operationType: 'insert',          // Operation Type
  //   fullDocument: { _id: ..., sku: ..., title: ..., ... },  // Complete Documentation
  //   ns: { db: 'shopdb', coll: 'products' },  // Namespace
  //   documentKey: { _id: ObjectId('...') }    // Document Primary Key
  // }
});
tipo de operação Significado documento completo
insert Inserir novo documento ✅ Documento concluído
update Atualização do documento ❌ Apenas alterações nos campos (requer updateLookup)
replace Substituição de documento ✅ Documento completo
delete Documento excluído ❌ Nenhum (O documento foi excluído)
drop Excluir conjunto ❌ Nenhum
rename Conjunto de renomeação ❌ Nenhum
invalidate Perda de efeito (por exemplo, conjunto excluído) ❌ Nenhuma

Análise dos pontos-chave:

  1. Os Change Streams devem ser executados em um conjunto de réplicas ou em um cluster fragmentado (dependendo do oplog).
  2. update Por padrão, o evento não retorna o documento completo; é necessário definir fullDocument: 'updateLookup'
  3. O campo _id é o resumeToken, usado para retomar downloads — salve-o e, após reiniciar, continue a reprodução de onde parou.

3. Filtragem em tubulação

Descrição do conceito: O Change Streams oferece suporte à filtragem no pipeline de agregação — ao aplicar etapas como $match e $project ao fluxo de eventos do oplog, apenas os eventos de interesse são encaminhados para o aplicativo. A filtragem é realizada no lado do servidor, reduzindo o tráfego de rede e a sobrecarga de processamento na camada de aplicativo.

Como funciona: watch() aceita uma matriz de pipelines de agregação como parâmetro; as etapas do pipeline são executadas no lado do servidor do MongoDB. Os eventos passam primeiro pelo pipeline para filtragem, e somente aqueles que são aprovados são enviados ao cliente. Ele suporta etapas como $match, $project, $addFields e $replaceRoot.

Política de filtragem:

Destino do filtro Etapa do pipeline Exemplo
Tipo de operação $match: {operationType} Apenas ouvir inserções
Campo do documento $match: {'fullDocument.field'} Apenas categorias específicas
Campo a ser alterado $match: {'updateDescription.updatedFields'} Apenas alteração de preço
Critérios de combinação $match: {$or: [...]} Inserção ou alteração de preço
JAVASCRIPT
// === Filter by Specific Actions ===
const changeStream = db.products.watch([
  { $match: { operationType: 'insert' } }
]);

// === Filter by Specific Fields ===
const changeStream = db.products.watch([
  { $match: { 'fullDocument.category': 'Electronics' } }
]);

// === Filter by Price Changes ===
const changeStream = db.products.watch([
  {
    $match: {
      $or: [
        { operationType: 'insert' },
        { operationType: 'update', 'updateDescription.updatedFields.price': { $exists: true } }
      ]
    }
  }
]);

Análise dos pontos-chave:

  1. A filtragem é realizada no lado do servidor, reduzindo o tráfego de rede desnecessário.
  2. O caminho do campo para $match utiliza a estrutura de eventos do oplog (como fullDocument.category), e não os campos do documento original.
  3. As etapas que exigem estado global, como $group e $limit, não são suportadas.

4. Configuração do fullDocument

Descrição do conceito: A opção fullDocument controla se os eventos do Change Stream incluem o documento completo. Por padrão, os eventos update retornam apenas os campos alterados (updateDescription) e não retornam o documento completo. Quando fullDocument: 'updateLookup' está definido, o MongoDB realiza uma consulta adicional para recuperar o conteúdo completo do documento atual.

Como funciona:

Pontos positivos e negativos do updateLookup:

Dimensão quando disponível atualizar consulta
Documentação completa Apenas inserir/substituir Todos os tipos de operação
Desempenho Teste de desempenho Sobrecarga adicional nas consultas (+10–20%)
Atualidade dos dados Hora da alteração Hora da consulta (pode estar sujeita a um pequeno atraso)
Casos de uso É necessário saber apenas quais campos foram alterados É necessário o documento completo para o processamento posterior
100%
sequenceDiagram
    participant App as Applications
    participant DB as MongoDB
    participant Doc as Document

    App->>DB: watch([], {fullDocument: 'updateLookup'})
    DB->>DB: oplog Generate update Event

    Note over DB: update By default, events contain only the changed fields.

    DB->>Doc: View the current complete document
    Doc-->>DB: Back to Latest Documents
    DB-->>App: Push Events + fullDocument

    Note over App: The incident includes the complete documentation
JAVASCRIPT
// === update Return to the full document ===
const changeStream = db.products.watch([], {
  fullDocument: 'updateLookup'
});

changeStream.on('change', (change) => {
  if (change.operationType === 'update') {
    console.log('Updated doc:', change.fullDocument);
    // Complete Documentation(After the default value changes)
  }
});

// === Return only the fields that have changed ===
const changeStream = db.products.watch([], {
  fullDocument: 'whenAvailable'  // Default
});

Análise dos pontos-chave:

  1. updateLookup Realiza uma consulta adicional, o que aumenta a carga no banco de dados em cenários com atualizações frequentes.
  2. updateLookup retorna o documento mais recente no momento da consulta, o que pode diferir ligeiramente do documento no momento da alteração (devido a outras modificações simultâneas).
  3. O evento delete não retorna um documento, mesmo que updateLookup esteja definido (o documento não existe mais)

5. Integração com o Mongoose

Explicação do conceito: O Mongoose 6+ oferece suporte nativo a Change Streams, que são retornados por meio de Model.watch(). Isso é totalmente consistente com a API do driver nativo do MongoDB, mas usá-lo diretamente na camada de modelo se alinha melhor às convenções de desenvolvimento do Mongoose.

Arquitetura orientada a eventos: Os fluxos de alterações são um componente essencial da arquitetura orientada a eventos — as alterações no banco de dados servem como fonte de eventos, gerando efeitos colaterais a jusante, como atualizações de cache, sincronização de buscas e envio de notificações.

100%
graph TB
    subgraph "Source of the Incident"
        DB[(MongoDB<br/>oplog)]
    end

    subgraph "Change Stream Bus"
        CS[Model.watch<br/>Change Flow]
    end

    subgraph "Event Consumer"
        N1[Email Notification]
        N2[Elasticsearch<br/>Search Synchronization]
        N3[Redis<br/>Cache Expiration]
        N4[WebSocket<br/>Real-time Notifications]
        N5[Audit Log<br/>Compliance Record]
    end

    DB --> CS
    CS --> N1
    CS --> N2
    CS --> N3
    CS --> N4
    CS --> N5

    style CS fill:#d4edda
JAVASCRIPT
// === mongoose Change Streams(mongoose 6+)===
const Product = mongoose.model('Product', productSchema);

// Monitoring Product Set Changes
const changeStream = Product.watch();

changeStream.on('change', (change) => {
  console.log(`${change.operationType}:`, change.fullDocument);
});

// Filter
const filteredStream = Product.watch([
  { $match: { operationType: { $in: ['insert', 'update'] } } }
]);

Análise dos pontos-chave:

  1. O watch() do Mongoose retorna o Change Stream nativo do MongoDB, e a API é totalmente idêntica.
  2. A conexão com o banco de dados deve ser mantida durante o monitoramento; caso a conexão seja perdida, o Change Stream é encerrado automaticamente.
  3. Em um ambiente de produção, recomenda-se encapsular o Change Stream Manager: reconexão automática + persistência do resumeToken

6. Cenários do mundo real

Visão geral do conceito: Os Change Streams têm três casos de uso clássicos em ambientes de produção: notificações em tempo real, sincronização de dados (CDC) e invalidação de cache. Cada um desses casos de uso é uma implementação típica de uma arquitetura orientada a eventos.

(1) Sistema de Notificação em Tempo Real

Cenário: A plataforma de comércio eletrônico ShopHub precisa enviar notificações por e-mail em tempo real e notificações push via WebSocket para o backend do comerciante sempre que um novo pedido for criado.

JAVASCRIPT
// === Monitor New Orders,Send a notification ===
const OrderStream = db.orders.watch([
  { $match: { operationType: 'insert' } }
]);

OrderStream.on('change', async (change) => {
  const order = change.fullDocument;

  // Send an Email
  await sendEmail(order.userId, 'Order confirmation', `Order ${order._id} received`);

  // Push Notifications
  await pushNotification(order.userId, {
    title: 'New Order',
    body: `Order total: $${order.total}`
  });

  // WebSocket Real-time Notifications
  io.emit('new_order', order);
});

(2) Sincronização de dados (CDC)

Cenário: O ShopHub precisa sincronizar os dados de produtos do MongoDB em tempo real com o Elasticsearch para oferecer suporte à pesquisa de texto completo. O Change Streams permite um pipeline de CDC (Control Data Change) com latência zero.

Arquitetura do CDC:

100%
graph LR
    A[MongoDB<br/>Product Data] -->|Change Stream| B[CDC Worker<br/>Node.jsProcess]
    B -->|insert/update| C[Elasticsearch<br/>Search Index]
    B -->|delete| C
    B -->|Change Log| D[(Redis<br/>resumeToken)]
    D -->|Restart and Recover| B

    style B fill:#d4edda
JAVASCRIPT
// === Monitoring MongoDB Change,Sync to Elasticsearch ===
const ProductStream = db.products.watch();

ProductStream.on('change', async (change) => {
  switch (change.operationType) {
    case 'insert':
    case 'update':
    case 'replace':
      await elasticsearch.index({
        index: 'products',
        id: change.documentKey._id.toString(),
        body: change.fullDocument
      });
      break;
    case 'delete':
      await elasticsearch.delete({
        index: 'products',
        id: change.documentKey._id.toString()
      });
      break;
  }
});

(3) Invalidação do cache

Cenário: O ShopHub usa o Redis para armazenar em cache as páginas de detalhes dos produtos. Quando os dados dos produtos são alterados, o cache é automaticamente esvaziado para evitar dados desatualizados.

JAVASCRIPT
// === Monitor Product Changes,Invalidate Redis cache ===
const ProductStream = db.products.watch();

ProductStream.on('change', async (change) => {
  const productId = change.documentKey._id.toString();
  await redis.del(`product:${productId}`);
  console.log(`Cache cleared for ${productId}`);
});

▶ Exemplo 1: Retomada do fluxo de alterações após uma interrupção

JAVASCRIPT
// Alice's TechCorp System: Change Stream Resume, Events Are Not Lost After a Restart
async function startResumableStream() {
  // 1. Retrieve last saved resumeToken from Redis
  let resumeToken = await redis.get('product_stream_token');
  let options = { fullDocument: 'updateLookup' };

  if (resumeToken) {
    options.resumeAfter = JSON.parse(resumeToken);
    console.log('Resuming from saved token');
  }

  // 2. Start Listening
  const changeStream = db.products.watch([], options);

  changeStream.on('change', async (change) => {
    // Processing Changes...
    console.log(`${change.operationType}: ${change.documentKey._id}`);

    // 3. Save after each event resumeToken
    await redis.set('product_stream_token', JSON.stringify(change._id));
  });

  changeStream.on('error', async (err) => {
    console.error('Stream error:', err.message);
    // 4. Delayed Reconnection After an Error
    setTimeout(startResumableStream, 5000);
  });
}

startResumableStream();

7. Limitações dos Change Streams

Nota conceitual: Os Change Streams dependem do oplog e dos conjuntos de réplicas e apresentam limitações claras. Compreender essas limitações é essencial para projetar sistemas confiáveis em tempo real.

Explicação detalhada das restrições:

Restrição Descrição Motivo Estratégia de mitigação
É necessário um conjunto de réplicas Não compatível com o modo autônomo Depende do oplog Use um conjunto de réplicas de nó único no ambiente de desenvolvimento
Limite de tamanho do oplog Padrão: 5% do espaço em disco O oplog é uma coleção com limite máximo Aumente o valor de oplogSize ou monitore a janela
Não é possível estender entre clusters Dentro de um único cluster O Oplog não se estende entre clusters Use o Kafka para conectar vários clusters
$where não é compatível Certos operadores O Oplog não registra detalhes da consulta Filtre usando o campo fullDocument
Ordem dos eventos Ordenados dentro de um único conjunto; não há ordem global entre conjuntos Oplog fragmentado por conjunto Classificação entre conjuntos realizada na camada de aplicação
Uso de memória Memória utilizada por conexão de observação O servidor mantém o estado do cursor Limite do número de conexões de observação

Monitoramento da janela do oplog: O oplog tem um tamanho fixo, e as entradas mais antigas são sobrescritas. Se o Change Stream consumir dados mais lentamente do que a velocidade com que o oplog é sobrescrito, o resumeToken irá expirar, e a funcionalidade de retomada a partir do ponto de interrupção não funcionará.

100%
graph LR
    A[oplog Write Speed<br/>1000 ops/s] --> B[oplog Capacity<br/>5% Disk ≈ 50GB]
    B --> C[oplog Window<br/>about 72 hours]
    C --> D{Consumption Rate?}
    D -->|Keep up| E[✅ Normal]
    D -->|Behind > 72h| F[❌ resumeToken Failure<br/>A full resync is required]

    style E fill:#d4edda
    style F fill:#f8d7da
Item de monitoramento Comando Limite de alerta
Janela do oplog rs.printReplicationInfo() < 24 horas
uso do oplog db.oplog.rs.stats() > 80%
Número de conexões do Change Stream db.currentOp() > 100
Atraso no evento Monitoramento da camada de aplicação > 10 segundos

▶ Exemplo: Real-time Analytics Dashboard com Change Streams (Difficulty ⭐⭐)

JAVASCRIPT
// Scene: ShopHub admin dashboard showing live order statistics
const mongoose = require('mongoose');

async function startOrderAnalytics() {
  const Order = mongoose.model('Order');
  
  // Track daily order count and revenue in-memory (for dashboard)
  let stats = {
    date: new Date().toISOString().split('T')[0],
    orderCount: 0,
    totalRevenue: 0,
    lastOrderId: null
  };

  // Watch for insert operations only (new orders)
  const changeStream = Order.watch([
    { $match: { operationType: 'insert' } }
  ], { fullDocument: true });

  changeStream.on('change', async (change) => {
    const order = change.fullDocument;
    
    // Update in-memory stats
    stats.orderCount++;
    stats.totalRevenue += order.total || 0;
    stats.lastOrderId = order._id;

    // Broadcast to WebSocket clients (dashboard)
    broadcastToDashboard({
      event: 'new_order',
      data: {
        orderId: order._id,
        total: order.total,
        stats: stats
      }
    });

    // Log to console for demo
    console.log(`[LIVE] Order #${order.orderNumber} - $${order.total}`);
    console.log(`   Today: ${stats.orderCount} orders, $${stats.totalRevenue} revenue`);
  });

  changeStream.on('error', (err) => {
    console.error('Change Stream error:', err);
    // Reconnect logic would go here
  });

  console.log('Order analytics stream started...');
}

// Simulate WebSocket broadcast
function broadcastToDashboard(data) {
  // In production: io.emit('order_update', data)
}

startOrderAnalytics();

Saída:

TEXT 📖 Somente leitura
Order analytics stream started...
[LIVE] Order #ORD-2026-001 - $599
   Today: 1 orders, $599 revenue
[LIVE] Order #ORD-2026-002 - $1299
   Today: 2 orders, $1898 revenue

▶ Exemplo: Guia prático para notificações de pedidos em tempo real no Change Streams

JAVASCRIPT
// Scene: Monitor New Orders, Send emails in real time + WebSocket Push + Sync to Elasticsearch

// 1. Start Listening (in Node.js)
const { MongoClient } = require('mongodb');

async function startOrderListener() {
  const client = new MongoClient('mongodb://localhost:27017/?replicaSet=rs0');
  await client.connect();
  const orders = client.db('shopdb').collection('orders');

  // Monitor all changes to the order collection
  const changeStream = orders.watch([
    {
      $match: {
        operationType: 'insert',  // Listen-only insertion
        'fullDocument.status': 'paid'  // Process only paid orders
      }
    }
  ]);

  changeStream.on('change', async (change) => {
    const order = change.fullDocument;
    console.log(`New Orders: ${order._id}, Amount: $${order.total}`);

    // 1. Send an email notification
    await sendEmail(order.userId, {
      subject: 'Order Confirmation',
      body: `Your Order ${order._id} Submitted,Total Amount $${order.total}`
    });

    // 2. WebSocket Real-time push notifications to the merchant's backend
    io.to('merchant-dashboard').emit('new_order', {
      orderId: order._id,
      total: order.total,
      items: order.items,
      timestamp: order.createdAt
    });

    // 3. Sync to Elasticsearch For Search
    await elasticsearch.index({
      index: 'orders',
      id: order._id.toString(),
      body: order
    });
  });

  // 4. Monitor Product Changes,Synchronized Update Elasticsearch
  const productStream = client.db('shopdb').collection('products').watch();
  productStream.on('change', async (change) => {
    if (change.operationType === 'delete') {
      await elasticsearch.delete({
        index: 'products',
        id: change.documentKey._id.toString()
      });
    } else if (change.fullDocument) {
      await elasticsearch.index({
        index: 'products',
        id: change.documentKey._id.toString(),
        body: change.fullDocument
      });
    }
  });

  console.log('Change Streams listening...');
}

// 5. Resume Download(Restore progress tracking after the app restarts)
async function resumeAfterRestart() {
  const resumeToken = await redis.get('change_stream_resume_token');
  const changeStream = orders.watch([], {
    resumeAfter: JSON.parse(resumeToken),
    fullDocument: 'updateLookup'
  });

  changeStream.on('change', async (change) => {
    // Processing Changes...
    // Save resume token
    await redis.set('change_stream_resume_token', JSON.stringify(change._id));
  });
}

startOrderListener().catch(console.error);

// 2. Manual Testing in mongosh
// Trigger Event:Insert a New Order
db.orders.insertOne({
  userId: 'user_001',
  items: [{ sku: 'PHONE-001', qty: 1, price: 599 }],
  total: 599,
  status: 'paid',
  createdAt: new Date()
});
// The app receives a notification immediately:Send an Email + WebSocket Push + ES Index

Saída: Acionado em menos de 100 ms após a inserção de um pedido, concluindo automaticamente todos os efeitos colaterais — como e-mails, notificações push e sincronização de pesquisa — sem a necessidade de polling.

❓ Perguntas Frequentes

P: O Change Streams funciona em tempo real? R: Quase em tempo real. O Change Streams é acionado imediatamente após a gravação de uma entrada no oplog, com uma latência inferior a 100 ms.

P: O Change Streams perde eventos? R: Não (a menos que o oplog seja sobrescrito). Você pode usar resumeAfter para retomar de onde parou.

P: Como posso salvar o andamento do monitoramento? R: Use resumeToken para armazená-lo externamente (por exemplo, no Redis) para que possa ser restaurado após a reinicialização.


📖 Resumo


📝 Exercícios

  1. Problema básico (⭐): Monitore todas as alterações na coleção products e exiba-as no console.
  2. Questão básica (⭐): Use $match para filtrar operações de inserção.
  3. Exercício avançado (⭐⭐): Fique atento a novos pedidos e envie notificações por e-mail (usando uma função de e-mail simulado).
  4. Exercício avançado (⭐⭐): Implementar a expiração do cache de produtos (sincronização com o Redis).
  5. Desafio (⭐⭐⭐): Concluir o sistema CDC (sincronização em tempo real entre o MongoDB e o Elasticsearch).
Web-Tutorial.com

Equipe Técnica Web-Tutorial

Uma plataforma de tutoriais mantida por diversos desenvolvedores. Cada tutorial é escrito e revisado por profissionais da área correspondente. Trabalhamos para manter nosso conteúdo preciso e confiável — se encontrar algum problema, avise-nos.

100%