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

Consumer Groups no Kafka

Atualizado em: 16 de setembro de 2026

Rack de servidores processando fluxos de dados no Node.js

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 → ocioso

Aumentar 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-events

Cada 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.

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