Pular para o conteúdo

apps/messaging-streaming/event-streaming

069 · Comunicação & Mensageria · ≈ 5 min de estudo

Event Streaming

Log imutável e ordenado de eventos que vários consumer groups leem de forma independente. Diferente de uma fila, o evento não some ao ser consumido: cada grupo guarda o próprio offset, anda no próprio ritmo e pode voltar atrás para reprocessar.

passos
6
arquivos
7
teste
1
tecnologias
4
Infraestrutura realTypeScriptBunKafkaDocker
Baixar cartão

Cenário

Um banco digital grava cada lançamento de conta (PIX, TED, cartão) num tópico Kafka. A projeção de saldo e a trilha de auditoria leem o mesmo log, cada uma no seu consumer group. Lançamentos da mesma conta precisam chegar em ordem; um lançamento corrompido não pode travar a partição nem sumir.

Planta

Fluxo
5/5
Produtorkey = conta, acks = allbank-transaction6 partiçõesGrupo balance-projectionGrupo transaction-auditbank-transaction-dlqoffset própriooffset própriopermanente, ilegívelou retry esgotadoseek para reprocessaruma partição
5

5 passos — reproduza para seguir o fluxo

Como funciona

6 passos

Log particionado, retido e reprocessável.

Esconde cada passo: lembre antes de tocar para revelar.

Vocabulário compartilhado

  1. 01

    Tópico e DLQ criados com partições explícitas; auto-criação desligada no broker (daria uma partição só)

  2. 02

    O produtor publica com key = accountId e acks: -1: mesma conta, mesma partição, ordem garantida

  3. 03

    Cada grupo roda com autoCommit: false e partitionsConsumedConcurrently: partições em paralelo, sequencial dentro de cada uma

  4. 04

    Erro transitório → retry curto no lugar com heartbeat(); permanente, JSON ilegível ou retry esgotado → DLQ com erro, origem e grupo

  5. 05

    Depois do efeito ou da DLQ, commit de offset + 1 em BigInt; lag = fim do log − offset commitado

  6. 06

    seek reposiciona uma partição para reprocessar; desligamento é pause → drenar → disconnect

apps/messaging-streaming/event-streaming

7 arquivos

src/

  • topic_transaction.tsCliente, tópicos com partições explícitas, publicação com chave e cálculo de lag
  • consumer_group.tsGrupo com commit manual, retry curto, DLQ e desligamento limpo
  • projection_transaction.tsProjeção de saldo idempotente por sequência e trilha de auditoria
  • config_transaction.tsVariáveis de ambiente validadas no boot
  • consumer_group.test.tstesteIntegração com Kafka real: grupos independentes, ordem por conta, retry, DLQ, restart sem reprocesso, replay e drenagem
  • demo.tsdemoDois grupos no mesmo log, lançamento corrompido na DLQ e replay de uma partição

./

  • docker.shinfraSobe o Kafka 3.7 em KRaft na porta 9092, sem auto-criação de tópico

Executar · com Docker

  1. ./docker.sh up# sobe Kafka
  2. cp .env.example .env# variáveis de ambiente
  3. bun install# dependências
  4. bun run demo# roda o cenário
  5. bun run test# integração contra o serviço real

Sob Bun o KafkaJS imprime TimeoutNegativeWarning, e com o broker recém-criado registra erros de coordinator ainda carregando que ele mesmo repete — os dois são inofensivos.

Requisitos
BunDocker
Sobe junto
Kafka
Esc

↑ ↓ navegarEnter abrir191 resultados