Home > Blog > Desenvolvimento Web
Desenvolvimento Web
JavaScript
Programação

Redis Streams no Node.js

Atualizado em: 4 de setembro de 2026

Rack de servidores processando fluxos de dados no Node.js

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 redis

Conexã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:

  1. pare novos loops de leitura;
  2. cancele XREADGROUP bloqueante;
  3. aguarde handlers atuais;
  4. não faça ACK de tarefas incompletas;
  5. 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.

10 melhores cursos de programação em 2026

Descubra os melhores cursos de programação. Aprenda a escolher o curso ideal para iniciar ou avançar na carreira de desenvolvedor

POSTS RELACIONADOS

Ver todos

Seta para a direita