As Transações Kafka no Node.js permitem publicar várias mensagens e confirmar offsets de consumo como uma única operação atômica dentro do Kafka. Se algo falhar antes do commit, as mensagens produzidas podem ser abortadas e consumidores configurados corretamente não as enxergam.
Esse recurso é útil no padrão consume-transform-produce: um serviço lê um evento, gera novos eventos e confirma o offset somente se a produção for concluída. Porém, transações Kafka não incluem bancos, APIs externas ou outros brokers. Elas oferecem atomicidade dentro do ecossistema Kafka, não uma transação distribuída universal.
Neste guia, você aprenderá a configurar um produtor transacional com KafkaJS, escolher transactionalId, enviar mensagens, abortar, confirmar offsets, configurar consumidores read_committed, evitar produtores zumbis e entender limites de exactly-once semantics.
O que são transações Kafka?
A documentação de transações do KafkaJS descreve uma interface com producer.transaction(), send, sendOffsets, commit e abort.
A documentação do Apache Kafka sobre semântica de entrega explica como idempotent producer, transações e consumers com isolation level apropriado participam de exactly-once.
O problema do consume-transform-produce
Sem transação:
- consumidor lê mensagem A;
- produz mensagem B;
- processo falha antes de confirmar offset;
- A é lida novamente;
- B é produzida novamente.
Ou, se confirmar primeiro:
- consumidor confirma A;
- processo falha antes de produzir B;
- A não volta;
- B é perdida.
Uma transação une a produção de B e o commit do offset de A.
Configurando o produtor
import { Kafka } from 'kafkajs';
const kafka = new Kafka({
clientId: 'orders-enricher',
brokers: process.env.KAFKA_BROKERS.split(',')
});
const producer = kafka.producer({
transactionalId: 'orders-enricher-0',
idempotent: true,
maxInFlightRequests: 1,
retry: {
retries: Number.MAX_SAFE_INTEGER
}
});
await producer.connect();KafkaJS ajusta opções necessárias quando o produtor é idempotente e transacional. Confirme a versão do cliente e broker usada.
transactionalId
O transactionalId identifica uma linha lógica de produção. Kafka usa o ID para fencing: se uma instância antiga reaparece, o broker rejeita suas escritas porque uma geração mais nova assumiu o mesmo ID.
Não gere um UUID aleatório a cada reinício. Isso impede fencing consistente e acumula produtores lógicos.
ID por partição
No consume-transform-produce, uma estratégia é incluir tópico e partição:
orders-enricher-orders.events.v1-3Cada partição mantém um produtor transacional estável. A implementação precisa coordenar atribuição e lifecycle durante rebalances.
Primeira transação
const transaction = await producer.transaction();
try {
await transaction.send({
topic: 'orders.enriched.v1',
messages: [{
key: order.orderId,
value: JSON.stringify(order)
}]
});
await transaction.commit();
} catch (error) {
await transaction.abort();
throw error;
}Depois do abort, as mensagens ficam marcadas como abortadas e consumidores read_committed não as recebem.
Várias mensagens
await transaction.sendBatch({
topicMessages: [
{
topic: 'orders.enriched.v1',
messages: [enrichedMessage]
},
{
topic: 'audit.events.v1',
messages: [auditMessage]
}
]
});Ou todas entram no commit ou todas são abortadas, dentro do Kafka.
Enviando offsets
await transaction.sendOffsets({
consumerGroupId: 'orders-enricher',
topics: [{
topic,
partitions: [{
partition,
offset: String(BigInt(message.offset) + 1n)
}]
}]
});O offset deve representar a próxima mensagem. O commit ocorre junto com as saídas.
Consumidor sem auto commit
await consumer.run({
autoCommit: false,
eachBatchAutoResolve: false,
eachBatch: async ({ batch, heartbeat, isRunning, isStale }) => {
for (const message of batch.messages) {
if (!isRunning() || isStale()) break;
await processTransactionally(batch, message);
await heartbeat();
}
}
});Não deixe o consumer confirmar offsets fora da transação, pois isso quebra a atomicidade desejada.
Fluxo completo
async function processTransactionally(batch, message) {
const transaction = await producer.transaction();
try {
const input = InputSchema.parse(JSON.parse(message.value));
const output = transform(input);
await transaction.send({
topic: 'orders.enriched.v1',
messages: [{
key: message.key,
value: JSON.stringify(output)
}]
});
await transaction.sendOffsets({
consumerGroupId: 'orders-enricher',
topics: [{
topic: batch.topic,
partitions: [{
partition: batch.partition,
offset: String(BigInt(message.offset) + 1n)
}]
}]
});
await transaction.commit();
} catch (error) {
await transaction.abort();
throw error;
}
}Consumidor read_committed
Consumidores que devem ignorar mensagens abortadas precisam usar isolamento committed:
const consumer = kafka.consumer({
groupId: 'billing-enriched-orders',
readUncommitted: false
});O valor padrão do KafkaJS é adequado para não retornar mensagens transacionais não confirmadas, mas valide a versão.
Last Stable Offset
Um consumer read_committed não avança além do Last Stable Offset enquanto existe transação aberta anterior. Transações muito longas podem aumentar latência para outros consumidores.
Timeout da transação
Kafka limita a duração. Não execute chamadas externas lentas dentro da transação. Mantenha o trabalho curto e previsível.
Se precisa buscar dados externos, considere pré-carregar, usar cache ou aceitar at-least-once com Inbox Pattern.
Exactly-once tem escopo
Kafka consegue oferecer exactly-once no ciclo:
Kafka → processamento determinístico → KafkaEle não torna atômico:
Kafka → PostgreSQL
Kafka → gateway de pagamento
Kafka → e-mail
Kafka → SQSNesses casos, use idempotência, Inbox, Outbox e reconciliação.
Banco de dados
Se uma transação Kafka confirma e o banco falha, ou vice-versa, não existe atomicidade conjunta. Opções:
- Inbox Pattern no consumidor;
- Outbox Pattern no serviço que grava;
- CDC;
- idempotency key em APIs externas;
- Saga com compensações.
Veja Inbox Pattern no Node.js e Outbox Pattern no Node.js.
Idempotent producer
O produtor idempotente evita duplicatas causadas por retries do próprio producer na mesma sessão. Ele usa producer ID e sequência por partição.
Isso não deduplica eventos semanticamente iguais enviados pela aplicação com IDs diferentes.
acks
Transações exigem confirmação de todas as réplicas em sincronia. O producer precisa operar com acks=-1. KafkaJS gerencia isso no modo transacional.
maxInFlightRequests
Limitar a um request em voo simplifica ordenação e exatamente uma vez. Isso pode reduzir throughput. Meça antes de adotar transações para todos os fluxos.
Retries
Retries do cliente precisam ser amplos para preservar a sessão transacional durante falhas recuperáveis. Alguns erros são abortáveis; outros exigem encerrar e recriar o producer.
Não repita cegamente uma transação com efeitos externos já realizados.
Fencing
Quando duas instâncias usam o mesmo transactionalId, a mais nova invalida a antiga. A antiga recebe erro de producer fenced e deve encerrar.
Isso é uma proteção, não um erro para ignorar e tentar infinitamente.
Rebalances
Se uma partição muda de consumidor, a nova instância deve usar o transactional ID associado à partição. A antiga será fenced.
Integre esse lifecycle com os eventos de group join e shutdown. Consulte Consumer Groups no Kafka.
Uma transação por mensagem
É simples, porém cria overhead de begin e commit. Em alto volume, processe um batch por transação:
const transaction = await producer.transaction();
try {
const outputs = batch.messages.map(transformMessage);
await transaction.send({ topic: outputTopic, messages: outputs });
await transaction.sendOffsets({ ...offsets });
await transaction.commit();
} catch (error) {
await transaction.abort();
throw error;
}Um erro reprocessa o batch inteiro. Mantenha idempotência e limite tamanho.
Mensagens inválidas
Uma poison message não deve causar abort eterno. Defina política:
- valide;
- publique a mensagem e erro em tópico de quarentena dentro da transação;
- confirme o offset;
- alerte a equipe.
Quarentena transacional
await transaction.send({
topic: 'orders.invalid.v1',
messages: [{
key: message.key,
value: JSON.stringify({
originalTopic: topic,
partition,
offset: message.offset,
reason: validationError.code,
payload: safePayload
})
}]
});
await transaction.sendOffsets({ ... });Observabilidade
Meça:
- transações iniciadas;
- commits;
- aborts;
- duração;
- producer fenced;
- erros abortáveis;
- lag do consumidor;
- mensagens por transação;
- tempo bloqueado no Last Stable Offset.
Logs
logger.info({
transactionalId,
topic,
partition,
firstOffset,
lastOffset,
messageCount
}, 'Transação Kafka confirmada');Não registre payloads sensíveis.
Tracing
Uma transação pode produzir vários eventos. Crie spans para transformação, envio e commit. O trace não substitui o transaction ID do Kafka.
Shutdown
Pare de iniciar novas transações, finalize ou aborte a atual, desconecte consumer e producer:
await consumer.stop();
await producer.disconnect();
await consumer.disconnect();Defina timeout externo. Veja Graceful Shutdown no Node.js.
Testes
Use Kafka real em container e cubra:
- commit com múltiplos tópicos;
- abort antes do commit;
- consumidor read_committed;
- offset confirmado na transação;
- falha após produção e antes do commit;
- producer fencing;
- rebalance;
- poison message;
- batch parcial.
Testcontainers
Consulte Testcontainers no Node.js para criar ambiente isolado no CI.
Quando usar
Transações são adequadas quando:
- entrada e saídas estão no Kafka;
- offset precisa ser confirmado junto às saídas;
- consumidores leem committed;
- o custo adicional é aceitável;
- transactional IDs podem ser estáveis.
Quando não usar
- apenas produção simples com idempotent producer suficiente;
- efeito principal está em banco externo;
- processamento demora muito;
- throughput máximo é prioridade;
- equipe não consegue operar fencing e rebalances.
Erros comuns
- ID aleatório: fencing perde valor.
- Consumer read_uncommitted: vê mensagens abortadas.
- Auto commit ativo: offset sai da transação.
- Chamada externa dentro da transação: duração cresce.
- Confundir com 2PC: banco continua fora.
- Ignorar producer fenced: instância zumbi continua tentando.
- Transação gigante: latência e abort aumentam.
- Sem política para poison message: partição trava.
Conclusão
As Transações Kafka no Node.js tornam atômico o ciclo de consumir, produzir e confirmar offsets dentro do Kafka. KafkaJS fornece uma API direta para iniciar, enviar, confirmar offsets, cometer ou abortar.
Use transactional IDs estáveis, consumers read_committed e handlers curtos. Entenda que exactly-once vale dentro do Kafka; bancos e APIs externas continuam exigindo Inbox, Outbox e idempotência. Com esse escopo claro, transações evitam duplicatas sem prometer uma atomicidade que o sistema não possui.



