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

Kafka com Node.js: Guia Prático

Atualizado em: 19 de agosto de 2026

Rack de servidores processando fluxos de dados no Node.js

Kafka com Node.js é uma opção para sistemas que precisam distribuir eventos em grande volume, manter histórico e permitir que diferentes consumidores processem o mesmo fluxo de forma independente. Ao contrário de uma fila tradicional em que uma mensagem costuma desaparecer depois do consumo, o Kafka grava registros em tópicos por um período configurado e controla a posição de cada grupo por meio de offsets.

Essa arquitetura é útil para integração entre serviços, auditoria, pipelines de dados, notificações e processamento assíncrono. Ela também exige decisões sobre partições, chaves, ordenação, retenção, duplicidade e reprocessamento. Para revisar a plataforma, consulte o que é Node.js e como criar uma API com Node.js.

Conceitos fundamentais

Um tópico é dividido em partições. Cada partição mantém uma sequência ordenada de registros. Produtores escolhem a partição diretamente ou por uma chave. Consumidores pertencem a grupos; dentro do mesmo grupo, cada partição é atribuída a apenas um consumidor por vez. Grupos diferentes podem ler o mesmo evento para finalidades distintas.

A ordem é garantida apenas dentro de uma partição. Se eventos do mesmo pedido precisam manter sequência, use o identificador do pedido como chave. Isso direciona registros com a mesma chave para a mesma partição, desde que a estratégia de particionamento permaneça consistente.

Instalação

npm init -y
npm install kafkajs
npm install --save-dev node:test

Configure brokers, clientId, credenciais e TLS por ambiente. Não codifique senhas no projeto e não desative a validação de certificados para resolver problemas de desenvolvimento.

Criando o cliente e o produtor

import { Kafka, logLevel } from 'kafkajs';

const kafka = new Kafka({
  clientId: 'orders-api',
  brokers: process.env.KAFKA_BROKERS.split(','),
  ssl: true,
  sasl: {
    mechanism: 'scram-sha-512',
    username: process.env.KAFKA_USERNAME,
    password: process.env.KAFKA_PASSWORD
  },
  logLevel: logLevel.INFO
});

const producer = kafka.producer({
  allowAutoTopicCreation: false,
  idempotent: true,
  maxInFlightRequests: 5
});

await producer.connect();

Desative criação automática de tópicos em produção para evitar nomes digitados incorretamente e configurações inadequadas. Um produtor idempotente reduz duplicações causadas por retries no protocolo, mas não elimina duplicidade no fluxo completo da aplicação.

Publicando um evento

export async function publishOrderCreated(order) {
  const event = {
    eventId: crypto.randomUUID(),
    eventType: 'order.created',
    version: 1,
    occurredAt: new Date().toISOString(),
    data: {
      orderId: order.id,
      customerId: order.customerId,
      total: order.total
    }
  };

  await producer.send({
    topic: 'orders.events.v1',
    acks: -1,
    messages: [{
      key: order.id,
      value: JSON.stringify(event),
      headers: {
        'content-type': 'application/json',
        'correlation-id': order.correlationId
      }
    }]
  });
}

acks: -1 solicita confirmação de todas as réplicas sincronizadas exigidas pela configuração do tópico. O resultado significa que o cluster aceitou o registro, não que algum consumidor concluiu o processamento.

Consumindo em grupo

const consumer = kafka.consumer({
  groupId: 'billing-service-v1',
  sessionTimeout: 30_000,
  heartbeatInterval: 3_000
});

await consumer.connect();
await consumer.subscribe({
  topic: 'orders.events.v1',
  fromBeginning: false
});

await consumer.run({
  partitionsConsumedConcurrently: 3,
  eachMessage: async ({ topic, partition, message }) => {
    const event = JSON.parse(message.value.toString('utf8'));
    validateOrderEvent(event);
    await processEventOnce(event);
  }
});

O método eachMessage simplifica commits e heartbeats para muitos casos. Trabalhos muito longos precisam respeitar timeouts ou usar eachBatch com controle explícito. Não bloqueie o event loop com CPU pesada; use workers ou serviços especializados.

Offsets e garantias

Um offset identifica a posição na partição. Normalmente o consumidor confirma o offset depois do handler. Se o processo falhar antes do commit, o evento será lido novamente. Se confirmar antes de concluir o efeito, pode perder trabalho. Por isso, o padrão mais comum é at-least-once com consumidores idempotentes.

Exactly-once é uma propriedade do fluxo completo, não apenas uma opção da biblioteca. Transações Kafka podem coordenar leitura e escrita dentro do Kafka, mas não tornam automaticamente uma atualização em banco externo exatamente uma vez.

Idempotência do consumidor

async function processEventOnce(event) {
  await database.transaction(async tx => {
    const inserted = await tx.insertProcessedEvent({
      eventId: event.eventId,
      consumer: 'billing-service-v1'
    });

    if (!inserted) return;

    await tx.createInvoice({
      orderId: event.data.orderId,
      total: event.data.total
    });
  });
}

Use uma restrição única para eventId e consumidor. A gravação do marcador e do efeito deve ocorrer na mesma transação. Essa abordagem suporta reentregas, reinícios e reprocessamentos controlados.

Partições e escalabilidade

O número de consumidores ativos em um grupo não pode ultrapassar o número de partições com ganho de paralelismo. Aumentar partições melhora capacidade, mas altera distribuição de chaves e aumenta custo operacional. Planeje crescimento antes de criar tópicos críticos.

  • use chaves quando a ordem por entidade importa;
  • evite uma chave única que concentre todo o tráfego;
  • meça skew entre partições;
  • dimensione retenção pelo volume real;
  • documente a estratégia de particionamento;
  • teste rebalances durante picos.

Contratos e evolução

Mensagens são APIs duradouras. Inclua versão, identificador, horário e tipo. Faça mudanças compatíveis, adicionando campos opcionais antes de remover ou alterar significados. Valide payloads em runtime; veja Zod no TypeScript.

Para ambientes com muitas equipes, considere um schema registry. JSON é simples, mas Avro, Protobuf ou JSON Schema podem melhorar governança e compatibilidade.

Retries e dead-letter

O Kafka não possui uma dead-letter queue automática universal. Crie tópicos de retry com atrasos definidos pela aplicação ou infraestrutura. Após o limite, publique em um tópico de dead-letter contendo evento original, erro sanitizado, tentativa e horário.

Não faça loop imediato na mesma mensagem. Uma falha permanente pode bloquear a partição e impedir eventos posteriores da mesma chave. Classifique erros e defina quando preservar ordem é mais importante que continuar o fluxo.

Segurança

  • use TLS entre clientes e brokers;
  • ative SASL com mecanismo suportado;
  • aplique ACL por tópico e grupo;
  • não coloque segredos ou dados pessoais desnecessários no evento;
  • limite tamanho das mensagens;
  • proteja ferramentas administrativas;
  • rotacione credenciais sem interromper todos os consumidores;
  • audite alterações de tópicos e ACLs.

Observabilidade

Meça taxa de produção, erros, latência, tamanho dos batches, consumer lag, rebalances, commits, tempo de processamento e distribuição por partição. O lag mostra a distância entre o último registro disponível e o offset confirmado pelo grupo.

Logs devem incluir topic, partition, offset, key, eventId, groupId e correlationId. Não use valores únicos como labels de métrica. Integre traces com cuidado para não aumentar demais o tamanho dos headers.

Graceful shutdown

Ao encerrar, pare de aceitar novo trabalho, conclua handlers em andamento e desconecte consumidor e produtor. A desconexão limpa a participação no grupo e reduz o tempo de rebalance.

async function shutdown() {
  await Promise.allSettled([
    consumer.disconnect(),
    producer.disconnect()
  ]);
}

process.once('SIGTERM', shutdown);
process.once('SIGINT', shutdown);

Consulte graceful shutdown no Node.js para organizar outros recursos do processo.

Testes

Teste serialização, chave, validação, duplicidade, retry, dead-letter e reprocessamento. Em integração, use tópicos isolados e aguarde condições com timeout em vez de sleeps fixos. Simule queda depois do efeito e antes do commit para comprovar idempotência.

O guia de testes unitários com Jest ajuda a manter regras separadas do cliente Kafka.

Erros comuns

  • Confiar em ordem global: a ordem existe apenas por partição;
  • Ignorar consumer lag: o sistema parece saudável enquanto acumula atraso;
  • Sem chave adequada: eventos relacionados podem chegar fora de sequência;
  • Assumir entrega única: reprocessamentos geram efeitos duplicados;
  • Criar tópicos automaticamente: configurações erradas entram em produção;
  • Executar trabalho longo sem heartbeat: o grupo inicia rebalances repetidos.

Checklist

  • tópicos e partições foram planejados;
  • chaves preservam a ordem necessária;
  • producer usa confirmação apropriada;
  • consumidores são idempotentes;
  • contratos possuem versão;
  • retry e dead-letter estão definidos;
  • lag e rebalances são monitorados;
  • shutdown foi testado.

Referências oficiais

Conclusão

Kafka com Node.js atende fluxos em que eventos precisam ser retidos, distribuídos e reprocessados. O desenho correto depende de escolher chaves e partições, entender offsets, aceitar reentregas e tornar consumidores idempotentes.

Comece com poucos tópicos e contratos claros. Monitore lag, teste rebalances e trate falhas como parte normal do sistema. Dessa forma, o Kafka deixa de ser apenas um canal rápido e se torna uma plataforma previsível para integração orientada a eventos.

Os 10 Melhores Cursos de Programação de 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