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:testConfigure 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.




