Usar Redis Streams no Node.js permite implementar filas, logs de eventos e processamento distribuído com retenção, IDs ordenados e consumer groups. Diferentemente de Pub/Sub, as mensagens permanecem armazenadas até serem removidas, possibilitando leitura posterior, recuperação de pendências e múltiplos grupos independentes.
Streams são úteis para tarefas assíncronas, integração entre serviços, processamento de eventos e pipelines que precisam de baixa latência. Porém, a aplicação precisa controlar retenção, mensagens pendentes, retries, idempotência e observabilidade.
Neste guia, você aprenderá XADD, XREAD, consumer groups, XREADGROUP, XACK, pending entries, claim, trimming, dead-letter, Node Redis, shutdown e testes.
O que são Redis Streams?
Redis Streams é uma estrutura de dados append-only com entradas identificadas por IDs. A documentação oficial de Redis Streams apresenta os comandos XADD, XRANGE, XREAD, consumer groups e gerenciamento de mensagens pendentes. Para o cliente JavaScript, consulte a documentação oficial do Node Redis.
Para filas especializadas em jobs, consulte BullMQ no Node.js. Para eventos transacionais, veja Outbox Pattern no Node.js.
Instalando o cliente
npm install redisConexão
import { createClient } from 'redis';
const redis = createClient({
url: process.env.REDIS_URL,
socket: {
connectTimeout: 5000,
reconnectStrategy(retries) {
return Math.min(retries * 100, 3000);
}
}
});
redis.on('error', error => {
logger.error({ err: error }, 'Erro no Redis');
});
await redis.connect();Não registre a URL completa, pois pode conter credenciais.
Adicionando uma entrada
const id = await redis.xAdd(
'orders:events',
'*',
{
type: 'order.created',
eventId: crypto.randomUUID(),
orderId,
payload: JSON.stringify(payload),
version: '1'
}
);O ID * pede ao Redis para gerar um valor baseado em tempo e sequência.
Formato da mensagem
Inclua campos estáveis:
- eventId;
- type;
- version;
- aggregateId;
- occurredAt;
- correlationId;
- payload serializado.
Não inclua objetos gigantes ou dados sensíveis sem necessidade.
Lendo histórico com XRANGE
const messages = await redis.xRange(
'orders:events',
'-',
'+',
{ COUNT: 100 }
);XRANGE é útil para inspeção, replay controlado e paginação por ID. Não carregue o stream inteiro em memória.
Leitura bloqueante
const result = await redis.xRead(
{
key: 'orders:events',
id: lastId
},
{
BLOCK: 5000,
COUNT: 50
}
);XREAD espera novas mensagens. O consumidor precisa armazenar o último ID confirmado.
Consumer groups
Consumer groups distribuem mensagens entre workers e mantêm uma lista de pendências.
try {
await redis.xGroupCreate(
'orders:events',
'order-workers',
'0',
{ MKSTREAM: true }
);
} catch (error) {
if (!String(error).includes('BUSYGROUP')) {
throw error;
}
}O ID 0 permite processar histórico. Use $ para começar apenas com novas entradas.
Lendo com XREADGROUP
const result = await redis.xReadGroup(
'order-workers',
consumerName,
{
key: 'orders:events',
id: '>'
},
{
BLOCK: 5000,
COUNT: 20
}
);O ID > solicita mensagens nunca entregues ao grupo.
Processamento
for (const stream of result ?? []) {
for (const message of stream.messages) {
try {
await handleMessage(message.message);
await redis.xAck(
'orders:events',
'order-workers',
message.id
);
} catch (error) {
logger.warn({
err: error,
messageId: message.id
}, 'Falha ao processar mensagem');
}
}
}Confirme apenas depois que o efeito foi concluído.
Pending Entries List
Mensagens entregues e não confirmadas ficam na PEL. Consulte resumo:
const pending = await redis.xPending(
'orders:events',
'order-workers'
);Monitore quantidade, consumidor, menor e maior ID.
Inspecionando pendências
const items = await redis.xPendingRange(
'orders:events',
'order-workers',
'-',
'+',
100
);Os resultados incluem tempo ocioso e quantidade de entregas.
Recuperando mensagens abandonadas
Um worker pode morrer antes do ACK. Use XAUTOCLAIM:
const claimed = await redis.xAutoClaim(
'orders:events',
'order-workers',
consumerName,
60000,
'0-0',
{ COUNT: 50 }
);O valor 60000 representa tempo ocioso mínimo em milissegundos. Escolha acima da duração normal da tarefa.
Evite claim prematuro
Se o processamento pode levar dois minutos, reclamar após trinta segundos gera execução concorrente. Meça a distribuição e defina heartbeat ou tarefas menores.
Retries
Redis Streams registra quantidade de entregas, mas não implementa automaticamente uma política de retry. A aplicação precisa decidir:
- falha transitória: tentar novamente;
- mensagem inválida: dead-letter;
- dependência indisponível: aguardar backoff;
- bug: alertar e bloquear promoção.
Backoff
Não reprocese imediatamente em loop. Uma estratégia é mover a mensagem para um sorted set com horário futuro ou usar outro stream de retry.
await redis.zAdd('orders:retry', {
score: Date.now() + delayMs,
value: JSON.stringify({ streamId, message })
});Um scheduler devolve itens vencidos ao stream.
Dead-letter stream
await redis.xAdd(
'orders:dead-letter',
'*',
{
originalId: message.id,
reason: errorCode,
deliveries: String(deliveries),
payload: JSON.stringify(message.message)
}
);
await redis.xAck(
'orders:events',
'order-workers',
message.id
);Proteja payloads sensíveis e defina retenção.
Idempotência
Entrega pode ocorrer mais de uma vez. O consumidor precisa impedir efeitos duplicados.
const key = `processed:${eventId}`;
const acquired = await redis.set(
key,
'1',
{ NX: true, EX: 86400 }
);
if (!acquired) {
return;
}Para efeitos críticos, deduplicação no mesmo banco da alteração é mais forte. Consulte Idempotência em APIs Node.js.
Problema da deduplicação separada
Se marcar no Redis e a atualização PostgreSQL falhar, uma nova entrega pode ser ignorada. Use transação no banco de destino com uma tabela de eventos processados quando a consistência for essencial.
Retenção
Sem trimming, o stream cresce indefinidamente.
await redis.xAdd(
'orders:events',
'*',
message,
{
TRIM: {
strategy: 'MAXLEN',
strategyModifier: '~',
threshold: 1000000
}
}
);O modificador aproximado reduz custo.
MAXLEN versus MINID
MAXLEN limita quantidade aproximada. MINID remove IDs anteriores a um limite temporal. Escolha conforme retenção por volume ou idade.
Cuidado com mensagens pendentes
Trimming pode remover o conteúdo de uma entrada ainda referenciada na PEL. Monitore lag dos grupos e mantenha retenção maior que o pior atraso operacional.
Múltiplos grupos
Cada grupo recebe todas as mensagens de forma independente. Use grupos diferentes para:
- faturamento;
- notificações;
- analytics;
- projeções;
- auditoria.
Dentro de um grupo, consumidores dividem o trabalho.
Ordenação
O stream preserva ordem por ID, mas workers concorrentes podem concluir fora de ordem. Se a ordem por agregado for crítica:
- particione por chave;
- use lock por agregado;
- inclua aggregateVersion;
- detecte lacunas;
- processe sequencialmente por chave.
Payload versionado
{
"type": "order.created",
"version": 2,
"eventId": "...",
"data": { ... }
}Consumidores devem tolerar campos adicionais.
Outbox para Redis Streams
Não grave no PostgreSQL e XADD separadamente sem proteção. Salve um evento na outbox e use um relay para publicar no Redis Stream.
Consulte Domain Events no Node.js.
Streams versus Pub/Sub
- Streams: persistência, histórico, groups, ACK e pendências.
- Pub/Sub: mensagens efêmeras para assinantes conectados.
Use Pub/Sub para invalidação ou sinalização em tempo real que pode ser perdida; use Streams quando processamento precisa ser recuperável.
Streams versus BullMQ
BullMQ fornece delay, retry, prioridades, repeatable jobs e abstrações de worker. Redis Streams oferece primitivos mais gerais e exige implementar políticas.
Streams versus Kafka
Kafka é projetado para logs distribuídos de alto throughput, retenção longa e particionamento. Redis Streams pode ser mais simples para volumes moderados e infraestrutura já baseada em Redis.
Consulte Kafka com Node.js.
Concorrência
Defina quantidade de consumidores com base em:
- CPU;
- latência de dependências;
- pool de banco;
- memória;
- limites externos;
- lag aceitável.
Mais workers não resolvem dependência saturada.
Backpressure
Use COUNT limitado e processe em lotes pequenos. Não leia milhares de mensagens se há poucas conexões disponíveis no banco.
Shutdown
Ao receber SIGTERM:
- pare novos loops de leitura;
- cancele XREADGROUP bloqueante;
- aguarde handlers atuais;
- não faça ACK de tarefas incompletas;
- feche clientes Redis.
await redis.quit();Consulte Graceful Shutdown no Node.js.
Cliente separado para bloqueio
Comandos bloqueantes podem usar uma conexão duplicada:
const consumer = redis.duplicate();
await consumer.connect();Assim, comandos administrativos não ficam presos atrás do XREADGROUP.
Observabilidade
Monitore:
- comprimento do stream;
- lag por grupo;
- mensagens pendentes;
- idade da pendência mais antiga;
- entregas por mensagem;
- taxa de ACK;
- retries;
- dead letters;
- duração do handler.
Logs
logger.info({
stream: 'orders:events',
group: 'order-workers',
messageId,
eventId,
eventType
}, 'Mensagem processada');Não registre payload completo por padrão.
Segurança
- use TLS;
- ACL com comandos necessários;
- rede privada;
- credenciais por Secret;
- limites de memória;
- política de eviction adequada;
- backup quando exigido.
Eviction de chaves pode comprometer streams. Configure Redis conforme criticidade.
Testes de integração
Use Redis real em container e teste:
- XADD e leitura;
- criação idempotente do grupo;
- ACK;
- mensagem pendente;
- XAUTOCLAIM;
- duplicação;
- dead-letter;
- trimming;
- shutdown.
Teste de crash
Leia uma mensagem sem ACK, encerre o consumidor e confirme que outro worker a reclama após o idle mínimo.
Teste de idempotência
Entregue o mesmo eventId duas vezes e confirme apenas um efeito no banco.
Erros comuns
- ACK antes do efeito: mensagem é perdida.
- Sem trimming: memória cresce.
- Claim cedo: processamento duplica.
- Sem idempotência: efeitos repetem.
- Loop sem backoff: falha consome CPU.
- Um cliente bloqueado para tudo: comandos atrasam.
- Payload sem versão: consumidores quebram.
- Retenção menor que lag: conteúdo desaparece.
Boas práticas
- Use eventId único.
- Crie consumer groups explicitamente.
- Faça ACK após sucesso.
- Monitore PEL.
- Recupere mensagens abandonadas.
- Implemente retries com backoff.
- Use dead-letter.
- Controle retenção.
- Torne consumidores idempotentes.
- Teste falhas reais.
Conclusão
Usar Redis Streams no Node.js oferece mensagens persistentes, consumer groups e recuperação de pendências com baixa latência. XADD, XREADGROUP e XACK formam a base do processamento distribuído.
Para produção, é essencial controlar trimming, retries, claim, idempotência e shutdown. Com outbox, métricas de lag e testes de crash, Redis Streams pode sustentar filas e pipelines confiáveis sem esconder os mecanismos de entrega.



