Event Streaming

Apache Kafka

Kafka é um commit log distribuído e append-only, não uma fila. Os registros são gravados uma vez e lidos várias vezes, o que lhe confere garantias de replay, fan-out e ordenação que um broker clássico não consegue igualar.

advanced15 min readUpdated 16 de set. de 2026
producer.ts
ts
// producer.ts
import { Kafka } from "kafkajs";

const kafka = new Kafka({
  clientId: "orders",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer();

await producer.connect();
await producer.send({
  topic: "orders.created",
  messages: [
    // The key decides the partition, so all events for one
    // customer stay in order.
    { key: order.customerId, value: JSON.stringify(order) },
  ],
});
await producer.disconnect();
Lançado
2011
Origem
LinkedIn
Modelo
Append-only commit log
Ordenação
Por partição
Entrega padrão
At-least-once
Retenção
Baseada em tempo ou tamanho

Por que importa

Por que as equipes escolhem o Kafka

Um log durável e reprodutível

Os registros são anexados e nunca removidos ao serem lidos. Qualquer consumer pode retroceder e reprocessar o histórico, algo que uma fila clássica simplesmente não pode oferecer.

Partições para escala e ordem

Um tópico é dividido em partições que se espalham entre os brokers. O paralelismo vem das partições, e a ordenação é garantida dentro de cada uma delas.

Múltiplos leitores independentes

Os consumer groups rastreiam seus próprios offsets, permitindo que um job de faturamento, um indexador de busca e um pipeline de analytics leiam os mesmos eventos sem interferirem entre si.

O panorama completo

As três ideias por trás do Kafka

Um tópico é um log append-only, producers gravam registros com chaves nele, e consumers rastreiam sua própria posição com offsets.

Topic

Append

Um log nomeado e particionado. Producers anexam registros ao final e cada registro mantém seu offset até que a retenção expire.

Consumer group

Share

Um grupo de consumers divide as partições entre si e rastreia seus próprios offsets commitados, mantendo a escala e o progresso independentes.

Offset

Track

A posição de um grupo em uma partição. Commitá-lo é o reconhecimento que decide entre at-least-once versus at-most-once.

HTML5 de uma olhada

Blocos de construção do Kafka

Topics

Logs nomeados e append-only que mantêm registros até que a retenção expire.

Partitions

A unidade de paralelismo e o limite da ordenação.

Producers

Serializam um registro, escolhem uma chave e o anexam a uma partição.

Consumer groups

Dividem partições entre membros e commitam offsets.

Offsets

A posição de um grupo em uma partição e seu ponto de reconhecimento.

Retention

Exclusão por tempo ou tamanho, ou compactação para manter o valor mais recente por chave.

Fluxo

A jornada de um registro

Um registro é gravado uma vez e lido várias vezes. Nada é removido quando um consumer o lê, o que é exatamente o que torna o replay possível.

  1. 1

    Produce

    O producer serializa um registro e escolhe uma partição através do hash de sua chave, garantindo que a mesma chave sempre caia no mesmo lugar.

  2. 2

    Append

    O broker líder da partição anexa o registro ao final de seu log e retorna o novo offset para o producer.

  3. 3

    Consume

    Cada consumer em um grupo detém um subconjunto das partições e lê seus registros na ordem do offset.

  4. 4

    Commit the offset

    Após o processamento, o consumer commita sua posição para que o grupo possa retomar dali após um reinício.

  5. 5

    Retain or replay

    O registro permanece no log até que a retenção expire, permitindo que outro grupo ou um job posterior retroceda e o leia novamente.

Uma breve historia

Dos logs do LinkedIn ao padrão de streaming

  1. 2011

    Kafka torna-se open-source

    O LinkedIn lança o Kafka como um commit log distribuído para seus fluxos de atividade.

    11
  2. 2012

    Incubação na Apache

    Kafka torna-se um projeto Apache e a adoção se espalha muito além do LinkedIn.

    12
  3. 2016

    Kafka Streams

    A versão 0.10 traz uma biblioteca de processamento de stream, seguida posteriormente pelo ksqlDB.

    16
  4. 2017

    Semântica Exactly-once

    A versão 0.11 adiciona producers idempotentes e transações para pipelines exactly-once.

    17
  5. 2021

    Início do KRaft

    O KIP-500 começa a substituir o ZooKeeper por um quorum de metadados integrado aos brokers.

    21
  6. 2025

    Kafka 4.0 remove o ZooKeeper

    Novos clusters rodam apenas em modo KRaft, encerrando a era do ZooKeeper.

    25

O guia completo

Apache Kafka: Tudo que voce precisa saber

O que o Kafka realmente é

O Apache Kafka é um distributed, append-only commit log. Essa única frase explica quase todo o resto. Os registros são anexados ao final de um log, cada um recebe um número crescentemente monotônico chamado de offset, e nada é jamais alterado no local. Os leitores não removem os registros quando os consomem; eles simplesmente movem um cursor para frente.

Esta é a diferença fundamental em relação a uma fila de mensagens clássica. Em uma fila, um consumidor pega uma mensagem e a mensagem desaparece. No Kafka, um consumidor lê um registro, lembra-se do seu offset, e o registro permanece onde está enquanto a política de retenção do topic permitir. Dez consumidores diferentes — e dez aplicações diferentes — podem ler o mesmo registro de forma independente, no seu próprio ritmo, sem precisar coordenar entre si.

O Kafka foi criado no LinkedIn para lidar com fluxos de atividade: visualizações de página, cliques, linhas de log, tudo isso a milhões de eventos por segundo. Ele foi aberto como open-source em 2011 e tornou-se um projeto Apache em 2012. Hoje, ele é a espinha dorsal padrão para event streaming: change data capture, pipelines de métricas, event sourcing, agregação de logs e stream processing.

Se você vem do RabbitMQ ou Redis, a mudança de mentalidade é parar de pensar em “uma mensagem a ser entregue” e começar a pensar em “um fato que foi registrado”.

Tópicos, partições e ordenação

Um tópico é um log nomeado e durável. Producers escrevem nele e consumers leem dele. Os tópicos são divididos em partições, e a partição é a unidade tanto de paralelismo quanto de ordenação.

  • Registros dentro de uma única partição são estritamente ordenados por offset.
  • Não há garantia de ordenação entre partições.
  • Mais partições permitem mais consumers paralelos, ao custo de mais arquivos, mais replicação e rebalances mais lentos.

Essa é a garantia mais importante do Kafka: a ordenação é por partição. Se dois eventos devem ser processados em ordem, eles devem cair na mesma partição. Caso contrário, o Kafka pode entregá-los a consumers diferentes que rodam simultaneamente, e a ordem será perdida.

As partições também determinam o paralelismo máximo de consumers em um grupo: um grupo pode ter, no máximo, um consumer por partição lendo-a ativamente. Seis partições significam, no máximo, seis consumers úteis; um sétimo ficaria ocioso.

# three partitions, replicated across three brokers
kafka-topics.sh --create \
  --topic orders.created \
  --partitions 6 \
  --replication-factor 3 \
  --bootstrap-server localhost:9092

Adicionar partições posteriormente é possível, mas isso altera o mapeamento de chave para partição para as chaves existentes. Assim, eventos de uma mesma entidade podem acabar divididos entre duas partições e perder sua ordem relativa. Defina a quantidade de partições no momento da criação do tópico, deixando uma margem para crescimento.

Producers, chaves e particionamento

Um producer serializa um registro e decide a qual partição ele pertence. O partitioner padrão gera um hash da chave (key) do registro e a mapeia para uma partição. A mesma chave sempre vai para a mesma partição, e é assim que você preserva a ordem para uma única entidade enquanto distribui entidades diferentes entre as partições.

await producer.send({
  topic: "orders.created",
  messages: [
    { key: order.customerId, value: JSON.stringify(order) },
  ],
});

A escolha da chave é uma decisão de design, não um detalhe. Use a chave por customerId e todos os eventos de um cliente serão ordenados juntos. Use a chave por orderId e você terá a distribuição máxima, mas sem ordenação entre eventos. Registros com chave nula são distribuídos para balanceamento — o sticky partitioner do Kafka preenche uma partição antes de passar para a próxima — mas eles não possuem nenhuma garantia de ordenação.

O producer também controla a durabilidade e o throughput através de batching e acknowledgements. acks decide quantas réplicas devem confirmar uma gravação, linger.ms e batch.size decidem quanto tempo ele espera para preencher um batch, e compression.type decide quanto de CPU será trocado por rede e disco. Mais detalhes sobre isso abaixo.

Brokers, replicação e o líder

Um broker é um servidor Kafka. Um cluster é composto por vários brokers trabalhando juntos. Cada partição possui um broker leader (líder) e zero ou mais followers (seguidores). Producers e consumers comunicam-se com o líder; os followers replicam o log.

A replicação é o que torna o Kafka durável. Se um líder falhar, uma das réplicas sincronizadas (ISR) é promovida e o cluster continua operando. A configuração acks define quanto tempo o producer aguarda:

  • acks=0 — “dispare e esqueça”; a opção mais rápida e menos segura.
  • acks=1 — o líder escreveu o dado; será perdido se o líder falhar antes da replicação.
  • acks=all — todas as réplicas sincronizadas confirmaram o recebimento; a opção mais segura.

Combine acks=all com min.insync.replicas=2 e um fator de replicação de três para um ambiente de produção: uma escrita só é confirmada quando pelo menos duas réplicas a possuem, portanto, a perda de um broker não resulta em perda de dados.

const producer = kafka.producer({
  idempotent: true,
  maxInFlightRequests: 1,
  transactionalId: "orders-producer",
});

O producer idempotente adiciona um número de sequência a cada batch para que o broker possa descartar duplicatas, eliminando as duplicatas acidentais que as tentativas de reenvio (retries) poderiam introduzir. Um cluster também possui um controller que gerencia a liderança das partições e os metadados. Versões antigas do Kafka utilizavam ZooKeeper para isso; o Kafka moderno utiliza KRaft, onde os brokers formam seu próprio quorum de metadados e o ZooKeeper é totalmente removido.

Grupos de consumidores e atribuição de partições

Os consumidores pertencem a um consumer group, identificado por groupId. O Kafka atribui cada partição a exatamente um consumidor dentro do grupo. Isso proporciona duas coisas simultaneamente: escalonamento horizontal, pois as partições são compartilhadas, e balanceamento de carga, pois nenhum consumidor em um grupo processa a mesma partição que outro.

const consumer = kafka.consumer({ groupId: "billing" });

await consumer.connect();
await consumer.subscribe({ topic: "orders.created", fromBeginning: true });

await consumer.run({
  eachMessage: async ({ partition, message }) => {
    const order = JSON.parse(message.value!.toString());
    console.log(`p${partition} @ ${message.offset}`, order.id);
  },
});

Quando consumidores entram ou saem, o Kafka realiza o rebalance: ele revoga as atribuições e distribui novas. O rebalance pausa o consumo, por isso é uma operação custosa; protocolos de rebalanceamento cooperativo reduzem a interrupção movendo apenas as partições que realmente precisam ser movidas. Um consumidor que para de enviar heartbeats dentro de session.timeout.ms é considerado morto e suas partições são reatribuídas.

O grupo também é a forma como o Kafka lembra o progresso. Cada grupo possui seus próprios offsets commitados, armazenados no tópico interno __consumer_offsets, permitindo que dois grupos lendo o mesmo tópico estejam em posições completamente diferentes sem qualquer coordenação.

Rebalanceamento e o loop de poll

Por baixo dos panos, um consumer é um poll loop. Ele busca registros, os entrega ao seu handler e, em seguida, busca novamente. O broker monitora a atividade separadamente através de heartbeats, portanto, um handler lento não parece “morto” imediatamente, mas existe um limite rígido: se um único poll demorar mais do que max.poll.interval.ms, o broker assume que o consumer travou e dispara um rebalanceamento.

const consumer = kafka.consumer({
  groupId: "billing",
  sessionTimeout: 45_000,
  heartbeatInterval: 3_000,
  rebalanceTimeout: 60_000,
});

Um rebalanceamento é disparado sempre que um consumer entra, sai ou é removido, sempre que partições ou tópicos mudam, e sempre que uma assinatura é alterada. Durante o rebalanceamento, o Kafka revoga as atribuições e pausa o consumo, portanto, rebalanceamentos frequentes destroem o throughput. Três hábitos ajudam a mantê-los raros:

  • Limite o trabalho por registro. Um handler que ocasionalmente executa por minutos eventualmente ultrapassará max.poll.interval.ms.
  • Use rebalanceamento cooperativo. O protocolo cooperativo incremental move apenas as partições que realmente precisam ser movidas, em vez de parar todo o grupo.
  • Use static membership. Definir um group.instance.id estável permite que um consumer reiniciado recupere suas partições antigas sem disparar um rebalanceamento.

O Kafka também oferece consumer.pause() e consumer.resume() para aplicar backpressure quando uma dependência downstream está com problemas. Pausar é melhor do que bloquear o poll loop: o consumer continua enviando heartbeats, permanece no grupo e simplesmente para de buscar dados até que você esteja pronto.

Offsets e garantias de entrega

Um offset é a posição de um consumer group em uma partição. Quando um consumer processa um registro, ele pode dar commit no offset, o que informa ao Kafka que “este grupo terminou tudo até aqui”. Se você fizer o commit cedo demais, uma falha resultará em perda de trabalho; se fizer tarde demais, uma falha causará o reprocessamento do trabalho.

  • At-most-once (no máximo uma vez) — commit antes do processamento. Uma falha faz com que o registro seja perdido.
  • At-least-once (pelo menos uma vez) — processamento e, depois, commit. Uma falha faz com que o registro seja reprocessado. Este é o padrão mais comum.
  • Exactly-once (exatamente uma vez) — transações e produtores idempotentes tornam o ciclo de leitura-processamento-escrita atômico dentro do Kafka.

O auto-commit acontece via timer em segundo plano. É conveniente, mas perigoso, pois pode realizar o commit de offsets de registros que ainda estão sendo processados. O commit manual após a conclusão do trabalho é a escolha padrão mais segura.

await consumer.run({
  autoCommit: false,
  eachMessage: async ({ topic, partition, message }) => {
    await chargeOrder(JSON.parse(message.value!.toString()));
    await consumer.commitOffsets([
      { topic, partition, offset: String(Number(message.offset) + 1) },
    ]);
  },
});

Como o at-least-once é a norma, os handlers devem ser idempotentes. O guia de Batch Processing aborda o mesmo contrato para filas de jobs: assuma que o trabalho pode ser executado duas vezes e faça com que a segunda execução seja um no-op. Uma chave de deduplicação derivada do trabalho — e não do offset do Kafka — torna o handler seguro mesmo que o mesmo evento lógico seja produzido duas vezes.

O exactly-once é real, mas mais limitado do que parece. Ele funciona para pipelines de Kafka para Kafka com produtores transacionais e isolation.level=read_committed, mas no momento em que você escreve em um banco de dados externo, você volta a precisar de escritas idempotentes nele.

Retenção, replay e compactação de log

O Kafka mantém os registros de acordo com uma política de retenção, e não por confirmação de entrega. retention.ms (padrão de sete dias) e retention.bytes limitam cada partição; quando qualquer um desses limites é excedido, os segmentos de log antigos são excluídos. Como nada é removido durante a leitura, um consumidor pode fazer o replay do histórico buscando um offset anterior ou resetando o grupo.

# rewind a group to the beginning of a topic
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group billing --topic orders.created \
  --reset-offsets --to-earliest --execute

O replay é o superpoder do Kafka. Um novo serviço pode ser iniciado lendo todo o histórico de um tópico. Uma correção de bug pode ser implantada e o último dia reprocessado. Um pipeline de analytics pode ser reconstruído a partir do log de eventos brutos. Uma fila não consegue fazer nada disso, pois os dados já teriam sumido.

Para estado com chave (keyed state), o outro modo de retenção é a compactação de log (log compaction). Com cleanup.policy=compact, o Kafka mantém pelo menos o valor mais recente para cada chave e descarta os valores mais antigos. Um valor null é um tombstone que deleta a chave. O log torna-se um changelog capaz de reconstruir uma tabela — exatamente o que o Kafka Streams usa para seus state stores e no que pipelines de change-data-capture se baseiam.

# keep the latest value per key instead of deleting by age
kafka-configs.sh --bootstrap-server localhost:9092 \
  --alter --entity-type topics --entity-name user.profiles \
  --add-config cleanup.policy=compact,min.cleanable.dirty.ratio=0.1

Schemas e o Schema Registry

Como o log sobrevive a qualquer aplicação individual, o formato de um registro torna-se um contrato entre equipes. Um produtor e um consumidor estão acoplados pelos bytes que trocam, e esse acoplamento persiste entre os deploys. Um schema registry transforma isso em um problema de compatibilidade gerenciada, em vez de uma surpresa.

O registry armazena schemas versionados — geralmente Avro, Protobuf ou JSON Schema — e atribui a cada um um id numérico. Os produtores registram um schema e escrevem o id junto ao payload; os consumidores buscam o schema pelo id e fazem a desserialização. Quando um schema muda, o registry impõe um modo de compatibilidade, como backward ou forward, rejeitando qualquer alteração que quebraria os leitores existentes.

import { SchemaRegistry } from "@kafkajs/confluent-schema-registry";

const registry = new SchemaRegistry({ host: "http://localhost:8081" });

const encoded = await registry.encode(schemaId, {
  orderId: order.id,
  totalCents: order.totalCents,
});

await producer.send({
  topic: "orders.created",
  messages: [{ key: order.customerId, value: encoded }],
});

Mesmo que você não utilize um registry, trate os payloads como uma API: adicione campos em vez de renomeá-los, defina valores padrão para novos campos e versione quando o significado mudar. Um consumidor executando a versão anterior deve ser capaz de ler um registro escrito pela versão seguinte.

Kafka Connect e Kafka Streams

Duas partes da plataforma evitam que você tenha que escrever a mesma “infraestrutura básica” duas vezes.

Kafka Connect é um framework para mover dados para dentro e para fora do Kafka via configuração, em vez de código. Conectores de origem (source connectors) extraem dados de bancos de dados, armazenamento de objetos ou APIs SaaS; conectores de destino (sink connectors) enviam dados para warehouses, índices de busca ou outro banco de dados. O Debezium, por exemplo, transforma um write-ahead log do PostgreSQL em um fluxo de eventos de alteração. O Connect roda como um cluster, rastreia offsets e tenta recuperar falhas, sendo a resposta padrão para “como colocar dados no Kafka” e “como tirar dados do Kafka”.

Kafka Streams é uma biblioteca cliente para processar dados no Kafka. Ela oferece uma DSL de stream — map, filter, groupByKey, join, window — sobre as abstrações de KStream e KTable, com stores de estado suportados por tópicos compactados. Ela roda dentro da sua aplicação, escala através da adição de instâncias e gerencia a tolerância a falhas por meio de tópicos de changelog. Se você prefere SQL, o ksqlDB oferece uma camada de consulta sobre esses mesmos conceitos.

const stream = builder.stream("orders.created");

stream
  .filter((key, order) => order.totalCents > 10_000)
  .groupBy((key) => order.customerId)
  .windowedBy(tumblingWindow({ size: 60 * 60 * 1000 }))
  .count()
  .toStream()
  .to("customer.hourly_orders");

Ambos valem a pena ser conhecidos, mesmo que você comece com produtores e consumidores simples, pois eles definem como é “o jeito Kafka” de trabalhar em escala.

Um exemplo minimalista de ponta a ponta

Ajuda visualizar tudo em um único arquivo: conectar um producer e um consumer, fazer a inscrição, executar e encerrar a aplicação de forma limpa.

import { Kafka } from "kafkajs";

const kafka = new Kafka({
  clientId: "orders-app",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer();
const consumer = kafka.consumer({ groupId: "orders-app" });

async function main() {
  await producer.connect();
  await producer.send({
    topic: "orders.created",
    messages: [{ key: "customer_1", value: JSON.stringify({ id: "order_1" }) }],
  });

  await consumer.connect();
  await consumer.subscribe({ topic: "orders.created", fromBeginning: true });
  await consumer.run({
    eachMessage: async ({ partition, message }) => {
      const order = JSON.parse(message.value!.toString());
      console.log(`p${partition} @ ${message.offset}`, order.id);
    },
  });
}

async function shutdown() {
  await consumer.disconnect();
  await producer.disconnect();
  process.exit(0);
}

process.on("SIGTERM", shutdown);
process.on("SIGINT", shutdown);

main().catch((err) => {
  console.error(err);
  process.exit(1);
});

Três pontos nesse arquivo são fundamentais em produção. O producer é conectado apenas uma vez e reutilizado, em vez de ser criado para cada mensagem. O consumer realiza a inscrição antes de iniciar a execução. E o processo trata SIGTERM, pois um encerramento abrupto durante um rebalance ou um commit deixa o grupo em um estado pior do que um fechamento gracioso.

Event sourcing e o outbox pattern

O Kafka é frequentemente descrito como a espinha dorsal do event sourcing, onde o log é a fonte da verdade e o estado atual é uma projeção de eventos. Em vez de armazenar apenas a linha mais recente, você armazena a sequência de fatos — order.created, order.paid, order.shipped — e reconstrói qualquer visualização ao reproduzi-los. Tópicos compactados transformam a projeção em uma tabela; a retenção a torna auditável.

A parte difícil nunca é escrever o evento; é escrever o evento e a linha do banco de dados de forma atômica. Um processo pode travar após commitar a linha e antes de publicar, ou após publicar e antes de commitar. Publicar dentro de uma transação de banco de dados é impossível, então a resposta padrão é o outbox pattern: escreva o evento em uma tabela outbox na mesma transação da alteração de estado e, em seguida, deixe que um relay ou conector CDC publique essas linhas no Kafka.

BEGIN;
INSERT INTO orders (id, status, total_cents)
VALUES ($1, 'created', $2);

INSERT INTO outbox (id, topic, payload)
VALUES ($1, 'orders.created', $2);
COMMIT;

O Debezium ou um pequeno relay de polling então monitora o outbox e produz para o tópico, deletando as linhas assim que forem publicadas. Como a escrita no outbox compartilha a transação, o evento existe exatamente quando o estado existe. Os consumidores ainda devem ser idempotentes, pois o relay pode publicar uma linha duas vezes após uma falha.

Desenvolvimento local e testes

Você não precisa de um cluster Kafka completo para desenvolver. Um broker de nó único no Docker Compose, ou o Redpanda em modo de compatibilidade, fornece tópicos, consumer groups e offsets diretamente no seu laptop.

services:
  kafka:
    image: redpandadata/redpanda:latest
    command: >
      redpanda start --overprovisioned --smp 1
      --kafka-addr PLAINTEXT://0.0.0.0:9092
      --advertise-kafka-addr PLAINTEXT://localhost:9092
    ports:
      - "9092:9092"

Para testes, o padrão é o mesmo de qualquer outra dependência de integração: inicie o broker em um container, crie os tópicos que o teste necessita e use um group id único por execução de teste para que os offsets commitados nunca vazem entre as execuções. Resete os offsets explicitamente quando um teste precisar ler o histórico e faça as asserções nos registros que seu consumer recebeu, em vez de baseá-las no tempo (timing).

const groupId = `test-${crypto.randomUUID()}`;
const consumer = kafka.consumer({ groupId });

Mantenha os handlers do consumer puros e enxutos — faça o parse, valide e chame um serviço — para que a maior parte da lógica possa ser testada via testes unitários sem a necessidade de um broker. Reserve os testes de integração para a fiação (wiring): um registro produzido chega ao handler correto e o offset é commitado em seguida?

Throughput, latência e batching

O Kafka é rápido porque utiliza batching e grava dados sequencialmente. Ambos são ajustáveis, e esse ajuste funciona como um botão de controle entre latência e throughput.

  • linger.ms — quanto tempo um producer espera para acumular um batch. Valores mais altos significam batches maiores e mais throughput, ao custo de maior latência.
  • batch.size — o limite máximo de bytes por batch de partição.
  • compression.typesnappy, lz4 ou zstd. A compressão reduz o uso de rede e disco; o zstd geralmente vence em termos de taxa de compressão.
  • fetch.min.bytes e fetch.max.wait.ms — quanto tempo os consumers esperam para preencher uma resposta de fetch.
  • max.poll.records — quantos registros um consumer processa por poll. Aumente este valor para esvaziar backlogs mais rapidamente, mas mantenha o processamento abaixo de max.poll.interval.ms, caso contrário, o consumer será removido do grupo.
const producer = kafka.producer({
  linger: { ms: 20 },
  compression: CompressionTypes.GZIP,
});

const consumer = kafka.consumer({
  groupId: "billing",
  maxBytesPerPartition: 1_048_576,
  maxWaitTimeInMs: 500,
});

Um producer otimizado para throughput pode usar linger.ms=20, um batch de um megabyte e zstd. Um producer que precisa publicar em milissegundos de um único dígito utiliza linger.ms=0. Não existe uma única resposta correta; existe apenas o trade-off que você escolhe e o p99 com o qual você consegue conviver.

Escolhendo entre Kafka, RabbitMQ e Redis

Esses três são frequentemente comparados como se fossem intercambiáveis. Eles não são.

  • Kafka é um log reproduzível. Escolha-o quando precisar de alto throughput, histórico durável, fan-out para muitos consumidores independentes, event sourcing, processamento de stream ou a capacidade de reprocessar dados. É o mais pesado para operar.
  • RabbitMQ é um message broker. Escolha-o quando o roteamento for importante: exchanges, routing keys, confirmações (acknowledgements) por mensagem, prioridades e dead-lettering. É a escolha natural para distribuição de tarefas e roteamento complexo, sendo mais leve que o Kafka em volumes modestos. Veja o guia do RabbitMQ.
  • Redis é um armazenamento de estruturas de dados em memória que também funciona como uma fila de jobs. Escolha-o quando o trabalho for simples, o volume for moderado e você já utilize Redis. O BullMQ no Redis oferece retentativas, agendamento e uma UI com quase nenhum overhead operacional. Veja Redis Queues.

Uma regra útil: se os consumidores precisarem reproduzir mensagens, fazer fan-out para muitos grupos ou ler o histórico, use Kafka. Se uma mensagem deve ser roteada para um worker específico e depois descartada, use RabbitMQ. Se for um job de background com política de retentativa, use Redis.

Segurança e controle de acesso

O Kafka geralmente contém os dados mais valiosos da infraestrutura, portanto, proteja-o rigorosamente.

  • Criptografia — habilite TLS para o tráfego entre brokers e entre cliente e broker. Kafka em texto simples dentro de uma VPC continua sendo texto simples.
  • Autenticação — utilize SASL/SCRAM ou mTLS para clientes. Evite listeners não autenticados em qualquer cenário que não seja um setup local descartável.
  • Autorização — use ACLs para conceder permissões de leitura, escrita ou criação em tópicos e grupos específicos. Um serviço não deve ser capaz de ler todos os tópicos apenas por conseguir se conectar.
  • Quotas — limite a largura de banda de produtores e consumidores por cliente, para que um serviço descontrolado não esgote os recursos do cluster.
  • Secrets — nunca incorpore credenciais diretamente em um produtor; injete-as via variáveis de ambiente ou por um gerenciador de secrets.

Monitoramento e consumer lag

O Kafka falha silenciosamente se você permitir. Quatro sinais cobrem a maior parte da realidade operacional.

  • Consumer lag — a diferença entre o offset mais recente e o offset commitado de um grupo, por partição. O aumento do lag é o primeiro sinal de que os consumers não estão conseguindo acompanhar o ritmo.
  • Under-replicated partitions — uma contagem diferente de zero significa que um broker está fora do ar ou lento, e a durabilidade dos dados está em risco.
  • Latência de requisição e de rede — os percentis do lado do broker revelam quando os discos ou a rede são o gargalo.
  • Uso de disco e contagem de segmentos — o Kafka é limitado pelo disco, e a replicação multiplica cada byte pelo fator de replicação.

Exponha o lag no mesmo dashboard que seus serviços e configure alertas baseados na tendência, não em um pico isolado. Um lag que cresce constantemente ao longo do dia é um problema de capacidade; um lag que tem um pico e depois se recupera geralmente é resultado de um rebalance ou de um deploy lento.

Melhores práticas

  • Projete as chaves com base na ordenação que você realmente precisa; lembre-se que a ordenação é por partição.
  • Use acks=all, min.insync.replicas=2 e um fator de replicação de três para tópicos importantes.
  • Ative o produtor idempotente e prefira commits de offset manuais.
  • Torne os consumidores idempotentes, pois o contrato padrão é “at-least-once”.
  • Mantenha o processamento de max.poll.records abaixo de max.poll.interval.ms para evitar tempestades de rebalanceamento.
  • Use log compaction para estados com chave e tombstones para exclusões.
  • Monitore o consumer lag por partição, e não apenas o total do cluster.
  • Trate os schemas dos registros como um contrato versionado, utilizando um registry quando mais de um time estiver envolvido.
  • Separe os tópicos por ciclo de vida e retenção, em vez de colocar tudo em um só.
  • Defina a retenção com base em um requisito real; sete dias é um padrão, não uma política.
  • Prefira Kafka gerenciado, a menos que operá-lo seja genuinamente o core business da sua empresa.
  • Use Connect para integrações padrão e Streams para processamento dentro do cluster antes de escrever código de integração customizado.

Erros comuns

  • Assumir ordenação global quando o Kafka ordena apenas dentro de uma partição.
  • Usar uma chave aleatória ou nula para eventos que devem permanecer ordenados.
  • Fazer o auto-commit de offsets antes que o trabalho seja concluído, resultando na perda de registros em caso de crash.
  • Criar handlers que não sejam idempotentes, duplicando efeitos colaterais durante um rebalance.
  • Criar um tópico com apenas uma partição e questionar por que os consumers não conseguem escalar.
  • Aumentar o número de partições posteriormente e alterar silenciosamente o roteamento de chaves.
  • Tratar o Kafka como uma API de request/response com respostas por mensagem.
  • Definir max.poll.records com um valor muito alto e exceder max.poll.interval.ms.
  • Esquecer que a replicação multiplica os custos de armazenamento e rede.
  • Deixar o consumer lag fora do dashboard até que o backlog acumule horas de atraso.
  • Rodar ZooKeeper em um cluster novo em 2026.

Próximos passos

O Kafka é o log no centro de um sistema orientado a eventos. O backend roadmap aborda arquitetura orientada a eventos e microservices, detalhando como eventos, producers e consumers se encaixam e como um log se torna a espinha dorsal entre os serviços. Se você está comparando brokers, os guias de RabbitMQ e Redis Queues mostram as alternativas focadas em roteamento e processamento de jobs, respectivamente. E como os consumers são basicamente workers de background com offsets, o guia de Batch Processing cobre as retentativas, idempotência e observabilidade que também se aplicam aqui.

Na pratica

Produzir, consumir, commitar, administrar

As quatro operações que compõem quase toda aplicação Kafka.

producer.ts
import { Kafka } from "kafkajs";

const kafka = new Kafka({
  clientId: "orders",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer({
  idempotent: true,
  maxInFlightRequests: 1,
});

await producer.connect();

await producer.send({
  topic: "orders.created",
  acks: -1, // wait for all in-sync replicas
  messages: [
    { key: order.customerId, value: JSON.stringify(order) },
  ],
});

await producer.disconnect();

Log do Kafka vs fila clássica

Uma fila deleta a mensagem assim que ela é reconhecida. Um log do Kafka a mantém, permitindo que qualquer grupo reproduza o histórico e vários leitores compartilhem o mesmo stream.

Kafka log
// A second group can read the same history
// from the beginning, months later.
await consumer.subscribe({
  topic: "orders.created",
  fromBeginning: true,
});
Fila clássica
// The message was acknowledged by the first
// worker and removed; it cannot be replayed
// or read by a second independent consumer.
await channel.ack(msg);

Auto commit vs manual commit

O auto-commit roda em um timer e pode avançar por registros que ainda estão sendo processados. Commitar após a conclusão do trabalho torna o modo de falha explícito.

Manual
await consumer.run({
  autoCommit: false,
  eachMessage: async ({ topic, partition, message }) => {
    await handle(message);
    await consumer.commitOffsets([
      {
        topic,
        partition,
        offset: String(Number(message.offset) + 1),
      },
    ]);
  },
});
Auto
await consumer.run({
  // Offsets are committed on a timer, so a
  // crash can skip records that were never
  // processed.
  eachMessage: async ({ message }) => {
    await handle(message);
  },
});

Trade-offs

O Kafka vale o peso operacional?

Kafka resolve problemas que filas não conseguem, mas exige infraestrutura real e um modelo mental diferente. Escolha-o pelas garantias, não pelo hype.

Strengths

  • O replay muda tudo

    Como o histórico é retido, você pode inicializar um novo serviço a partir do log, reconstruir uma projeção após a correção de um bug e executar analytics em eventos brutos muito tempo depois de terem ocorrido.

  • Fan-out sem coordenação

    Qualquer número de consumer groups independentes lê o mesmo tópico em seu próprio ritmo. Adicionar um leitor não custa nada ao producer e nunca perturba os consumers existentes.

  • Throughput com escala horizontal

    Escritas sequenciais, batching e partições permitem que um pequeno cluster absorva milhões de registros por segundo, e você escala adicionando brokers e partições.

Trade-offs

  • É uma plataforma, não uma biblioteca

    Brokers, replicação, rebalances, dimensionamento de disco e consumer lag são todos de sua responsabilidade monitorar. Serviços gerenciados ajudam, mas Kafka nunca é uma dependência única que você simplesmente esquece.

  • Roteamento não é o seu trabalho

    Kafka não possui exchanges ou chaves de roteamento. Consumers leem tópicos inteiros, portanto, o roteamento seletivo por mensagem deve ser construído no design do tópico ou no código da aplicação.

  • Ordenação apenas dentro de uma partição

    A ordem global não está disponível. Garantir a ordenação correta significa escolher chaves e contagens de partições cuidadosamente, e você não pode adicionar partições posteriormente sem quebrar o roteamento de chaves.

Perguntas frequentes

Perguntas frequentes

Keep learning

Related topics from the roadmap.

$ comecar a aprender

Pronto para aprender Apache Kafka?

Nosso tutorial interativo te guia por Apache Kafka passo a passo — com quizzes e codigo real que voce pode executar no navegador.