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
- O conceito de fluxos de alteração (com base no oplog)
- db.collection.watch() Escuta
- Filtragem de pipeline ($match)
- fullDocument: 'updateLookup'
- Integração com o Mongoose
- Cenários típicos (notificações em tempo real, sincronização de dados)
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.
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:
// === 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:
- Os Change Streams devem ser executados em um conjunto de réplicas ou em um cluster fragmentado (dependendo do oplog).
updatePor padrão, o evento não retorna o documento completo; é necessário definirfullDocument: 'updateLookup'- 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 |
// === 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:
- A filtragem é realizada no lado do servidor, reduzindo o tráfego de rede desnecessário.
- O caminho do campo para
$matchutiliza a estrutura de eventos do oplog (comofullDocument.category), e não os campos do documento original. - As etapas que exigem estado global, como
$groupe$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:
whenAvailable(padrão): os eventos de inserção/substituição incluem o documento inteiro; os eventos de atualização incluem apenas os campos que foram alterados; os eventos de exclusão não incluem nenhum documentoupdateLookup: Todos os eventos incluem o documento completo (os eventos de atualização consultam adicionalmente o documento mais recente), mas há uma sobrecarga adicional na leitura e um leve atraso.required: É necessário retornar o documento completo; caso contrário, será gerado um erro.
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 |
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
// === 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:
updateLookupRealiza uma consulta adicional, o que aumenta a carga no banco de dados em cenários com atualizações frequentes.updateLookupretorna 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).- O evento
deletenão retorna um documento, mesmo queupdateLookupesteja 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.
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
// === 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:
- O
watch()do Mongoose retorna o Change Stream nativo do MongoDB, e a API é totalmente idêntica. - 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.
- 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.
// === 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:
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
// === 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.
// === 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
// 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á.
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 ⭐⭐)
// 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 leituraOrder 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
// 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
resumeAfterpara retomar de onde parou.
P: Como posso salvar o andamento do monitoramento? R: Use
resumeTokenpara armazená-lo externamente (por exemplo, no Redis) para que possa ser restaurado após a reinicialização.
📖 Resumo
- Fluxos de alterações: notificações de alterações em tempo real com base no oplog
- Método watch() para monitoramento + filtragem com $match
- operationType: inserir / atualizar / excluir / substituir / remover
- fullDocument: 'updateLookup' Obter o documento completo
- Cenários típicos: notificações em tempo real, CDC, invalidação de cache
- Número necessário de conjuntos de réplicas, limite de tamanho do oplog
📝 Exercícios
- Problema básico (⭐): Monitore todas as alterações na coleção
productse exiba-as no console. - Questão básica (⭐): Use $match para filtrar operações de inserção.
- Exercício avançado (⭐⭐): Fique atento a novos pedidos e envie notificações por e-mail (usando uma função de e-mail simulado).
- Exercício avançado (⭐⭐): Implementar a expiração do cache de produtos (sincronização com o Redis).
- Desafio (⭐⭐⭐): Concluir o sistema CDC (sincronização em tempo real entre o MongoDB e o Elasticsearch).