Os Consumer Groups no Kafka permitem distribuir partições de um tópico entre várias instâncias de uma aplicação Node.js. Cada partição é atribuída a apenas um membro do grupo por vez, enquanto grupos diferentes podem ler o mesmo tópico de forma independente.
Esse mecanismo oferece escalabilidade e tolerância a falhas, mas introduz conceitos importantes: offsets, heartbeats, session timeout, rebalance, ordenação por partição, concorrência e lag. Uma configuração inadequada pode provocar rebalances frequentes, mensagens reprocessadas ou consumidores removidos do grupo durante tarefas lentas.
Neste guia, você aprenderá a criar consumer groups com KafkaJS, entender atribuição de partições, processar mensagens com eachMessage e eachBatch, controlar offsets, pausar partições, medir lag e realizar shutdown sem perder trabalho.
O que é um consumer group?
A documentação de consumo do KafkaJS explica que um grupo coordena máquinas ou processos e distribui tópicos entre consumidores. Quando um membro falha, as partições são redistribuídas.
Os detalhes do protocolo e das configurações estão também na documentação oficial do Apache Kafka.
Partições definem paralelismo
Se um tópico possui quatro partições, um grupo pode ter no máximo quatro consumidores ativos processando esse tópico ao mesmo tempo. Um quinto membro permanece sem partição.
Tópico orders.events
Partição 0 → Consumer A
Partição 1 → Consumer B
Partição 2 → Consumer C
Partição 3 → Consumer D
Consumer E → ociosoAumentar réplicas sem aumentar partições não cria paralelismo adicional.
Criando o consumer
import { Kafka } from 'kafkajs';
const kafka = new Kafka({
clientId: 'billing-service',
brokers: process.env.KAFKA_BROKERS.split(',')
});
const consumer = kafka.consumer({
groupId: 'billing-order-events',
sessionTimeout: 30_000,
heartbeatInterval: 3_000,
rebalanceTimeout: 60_000
});O groupId representa uma aplicação lógica. Instâncias do mesmo serviço devem compartilhar o mesmo ID.
Grupos independentes
O serviço de faturamento e o serviço de notificações usam IDs diferentes:
billing-order-events
notifications-order-eventsCada grupo mantém seus próprios offsets e recebe todas as mensagens.
Conectando e assinando
await consumer.connect();
await consumer.subscribe({
topics: ['orders.events.v1'],
fromBeginning: false
});fromBeginning é usado somente quando o grupo não possui offset válido. Com false, começa no final; com true, no início disponível.
eachMessage
await consumer.run({
eachMessage: async ({ topic, partition, message, heartbeat, pause }) => {
const event = parseAndValidate(message.value);
await processEvent(event);
}
});eachMessage é a interface mais simples. KafkaJS resolve offsets e executa heartbeat conforme necessário, mas o handler não deve bloquear por mais que o session timeout.
Ordenação
Kafka garante ordem dentro de uma partição, não entre todas as partições. Para eventos do mesmo pedido, use a mesma chave:
await producer.send({
topic: 'orders.events.v1',
messages: [{
key: event.orderId,
value: JSON.stringify(event)
}]
});A chave consistente envia eventos relacionados para a mesma partição.
partitionsConsumedConcurrently
await consumer.run({
partitionsConsumedConcurrently: 3,
eachMessage: async ({ message }) => {
await processEvent(message);
}
});Mensagens da mesma partição continuam sequenciais. Partições diferentes podem ser processadas em paralelo.
Não configure concorrência maior que o número de partições atribuídas sem medir. Para trabalho CPU-bound, o limite também depende dos núcleos.
eachBatch
await consumer.run({
eachBatchAutoResolve: false,
eachBatch: async ({
batch,
resolveOffset,
heartbeat,
isRunning,
isStale
}) => {
for (const message of batch.messages) {
if (!isRunning() || isStale()) break;
await processEvent(message);
resolveOffset(message.offset);
await heartbeat();
}
}
});eachBatch oferece controle sobre offsets, heartbeat e shutdown. Use apenas quando o time compreende essas relações.
Offsets
O offset representa a posição de uma mensagem na partição. O grupo armazena o próximo ponto de leitura. Confirmar offset antes de concluir pode perder processamento; confirmar depois pode gerar reprocessamento.
Auto commit
await consumer.run({
autoCommitInterval: 5000,
autoCommitThreshold: 100,
eachMessage: async payload => {
await processMessage(payload.message);
}
});Offsets resolvidos são confirmados quando uma das condições é atingida. Commit frequente aumenta tráfego; commit raro aumenta reprocessamento após falha.
Commit manual
await consumer.run({
autoCommit: false,
eachMessage: async ({ topic, partition, message }) => {
await processMessage(message);
await consumer.commitOffsets([{
topic,
partition,
offset: String(BigInt(message.offset) + 1n)
}]);
}
});O offset confirmado normalmente é o próximo a ser lido, por isso soma-se um.
Offset junto com o banco
Para maior atomicidade, armazene o offset na mesma transação da alteração de negócio:
await db.transaction(async tx => {
await updateProjection(tx, event);
await saveOffset(tx, {
groupId,
topic,
partition,
nextOffset: String(BigInt(message.offset) + 1n)
});
});Ao reiniciar, use seek para começar no offset externo. Esse modelo exige implementação cuidadosa.
Inbox Pattern
Uma alternativa é registrar o ID do evento na mesma transação e descartar duplicatas. Veja Inbox Pattern no Node.js.
Heartbeats
Heartbeats informam ao coordenador que o consumidor continua ativo. Se o broker não recebe dentro do session timeout, remove o membro e inicia rebalance.
Em eachBatch, chame heartbeat() durante loops longos.
Session timeout
Um valor curto detecta falhas rapidamente, mas pode remover consumidores durante pausas de GC ou rede. Um valor longo reduz rebalances falsos, porém demora a redistribuir partições após falha.
Meça a aplicação e o ambiente; não copie valores sem contexto.
Rebalance
Um rebalance ocorre quando:
- um consumidor entra ou sai;
- um membro perde heartbeats;
- partições são adicionadas;
- a assinatura muda;
- o coordenador muda.
Durante o rebalance, o consumo pode pausar temporariamente.
Deploys e rebalances
Um rolling deploy inicia e encerra consumidores, causando redistribuição. Reduza impacto:
- faça shutdown gracioso;
- evite reinícios desnecessários;
- use readiness correta;
- mantenha handlers curtos;
- limite concorrência de deploy.
Shutdown gracioso
let shuttingDown = false;
async function shutdown() {
if (shuttingDown) return;
shuttingDown = true;
await consumer.stop();
await consumer.disconnect();
}
process.once('SIGTERM', shutdown);
process.once('SIGINT', shutdown);stop aguarda o handler atual conforme a API. Defina também um limite externo do orchestrator.
Veja Graceful Shutdown no Node.js.
Pause e resume
await consumer.run({
eachMessage: async ({ topic, partition, message, pause }) => {
try {
await callDependency(message);
} catch (error) {
if (error.statusCode === 429) {
const resume = pause();
setTimeout(resume, error.retryAfterMs);
}
throw error;
}
}
});Pausar uma partição permite que outras continuem. Isso evita parar todo o tópico por uma dependência lenta.
Retries
KafkaJS possui retry para operações do cliente, mas erro de negócio no handler precisa de estratégia própria. Opções:
- retry local com backoff;
- tópicos de retry;
- dead-letter topic;
- pausa de partição;
- marcação de falha no banco.
Consulte Retry com Backoff no Node.js.
Poison message
Uma mensagem inválida pode bloquear a partição se sempre falha e nunca avança o offset. Valide e, após política definida, publique em quarentena e resolva o offset.
Não pule silenciosamente sem registrar motivo e conteúdo seguro para investigação.
Lag
Lag é a diferença entre o fim da partição e o offset confirmado do grupo. Um lag crescente indica que o consumidor não acompanha a produção.
Monitore por tópico e partição, não apenas a soma. Uma única partição quente pode ficar atrasada enquanto as outras estão zeradas.
High watermark
Em eachBatch, batch.highWatermark ajuda a estimar lag. Converta offsets como inteiros grandes, pois podem ultrapassar o limite seguro de Number.
Partição quente
Se uma chave domina o tráfego, uma partição recebe muito mais mensagens. Soluções possíveis:
- revisar a chave;
- aumentar partições;
- usar sharding adicional;
- separar tipos de eventos;
- otimizar o consumidor.
Alterar chave pode mudar ordenação e distribuição. Planeje migração.
Número de partições
Aumentar partições é simples, mas mensagens com a mesma chave podem mudar de partição depois da alteração, dependendo do particionador. Isso afeta ordenação histórica.
fromBeginning
fromBeginning: true não força replay se o grupo já possui offsets. Para replay, use novo group ID, reset de offsets ou seek.
Seek
consumer.seek({
topic: 'orders.events.v1',
partition: 0,
offset: '12345'
});Mensagens de batches ativos podem se tornar stale. Em eachBatch, verifique isStale().
Expressões regulares
await consumer.subscribe({
topics: [/^tenant-.*-events$/]
});KafkaJS não adiciona automaticamente tópicos criados depois da assinatura regex em todos os cenários. Atualize metadados ou reinicie de forma controlada conforme a documentação.
Segurança
Use TLS e SASL quando necessário:
const kafka = new Kafka({
brokers,
ssl: true,
sasl: {
mechanism: 'scram-sha-512',
username: process.env.KAFKA_USERNAME,
password: process.env.KAFKA_PASSWORD
}
});Conceda ACL de leitura somente aos tópicos e grupos necessários.
Observabilidade
Monitore:
- lag por partição;
- rebalances;
- membros ativos;
- tempo de processamento;
- erros por tipo;
- commits;
- heartbeats;
- partições pausadas;
- throughput.
Veja Métricas Prometheus no Node.js.
Instrumentation events
KafkaJS expõe eventos de instrumentação. Use-os para métricas e logs, mas evite handlers pesados que prejudiquem o cliente.
Testes
Teste com broker real em container:
- dois consumidores e várias partições;
- entrada e saída de membro;
- handler lento;
- falha antes do commit;
- rebalance durante batch;
- pause e resume;
- shutdown;
- replay.
Erros comuns
- Mais consumidores que partições: réplicas ficam ociosas.
- Handler longo sem heartbeat: membro é removido.
- Commit antes do negócio: mensagem é perdida.
- Sem idempotência: reprocessamento duplica efeitos.
- Ignorar lag por partição: hotspot fica oculto.
- Group ID novo em cada deploy: todas as mensagens são relidas ou puladas.
- Rebalance frequente: throughput cai.
- Seek sem isStale: batch antigo continua processando.
Conclusão
Os Consumer Groups no Kafka distribuem partições entre instâncias e permitem escalar consumidores Node.js. O paralelismo real é limitado pelo número de partições, enquanto offsets controlam recuperação e replay.
Configure timeouts com base em métricas, mantenha handlers idempotentes, monitore lag e implemente shutdown gracioso. Com controle de offsets, heartbeats e rebalances, consumer groups oferecem escala sem sacrificar ordenação dentro de cada partição.



