Pular para o conteúdo

apps/messaging-streaming/aggregator

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

Aggregator

Combina uma série de mensagens relacionadas (mesmo correlationId) em uma única mensagem consolidada, usando uma condição de conclusão — contagem, timeout, ou ambos. É o padrão inverso do Splitter: aqui muitas mensagens viram uma; no Splitter, uma vira muitas.

passos
6
arquivos
8
testes
2
tecnologias
5
Infraestrutura realTypeScriptBunElysiaRabbitMQDocker
Baixar cartão

Cenário

Ao final do dia, cada transação liquidada (PIX, TED, boleto) gera uma confirmação publicada na exchange settlement, identificada por businessDate — o correlationId do lote. O worker agregador junta as confirmações de cada data até chegar o total esperado ou até a janela vencer, e publica um relatório de fechamento (EOD) em settlement-report.

Planta

Sequência
6/6
confirmação 1..5 (businessDate = correlationId)1entrega; ack ao entrar na janela2janela por businessDate, dedup por transactionId3relatório EOD ao completar a contagem (completedBy: count)4outro businessDate com 3/5: fecha no timeout com 2 pendentesmensagem que não é confirmação5nack sem requeue → DLX → DLQ6Liquidaçãosettlement-confirmationworker_aggregatorsettlement-reportsettlement-confirmation-dlq
6

6 passos — reproduza para seguir o fluxo

  1. 01

    Cada confirmação traz businessDate, transactionId, total e amountCents, validada por schema — a inválida vai para a DLQ

  2. 02

    A primeira confirmação de uma data abre a janela e o timer; a confirmação é confirmada (ack) assim que entra na janela

  3. 03

    Reentrega com o mesmo transactionId substitui a si mesma em vez de inflar a contagem

  4. 04

    Ao chegar total, a janela fecha por count e cancela o timer; se o timer vence antes, fecha por timeout com transactionsPending

  5. 05

    O relatório sai uma única vez, por confirm channel com mandatory, e o estado da data é descartado

  6. 06

    O worker roda em processo próprio, com recovery do amqplib e probes /health/live e /health/ready

apps/messaging-streaming/aggregator

8 arquivos

src/

  • aggregation_settlement.tsMáquina de estado da janela: contagem, timeout e deduplicação por transactionId
  • worker_aggregator.tsWorker: consome confirmações, publica relatórios, probes e desligamento limpo
  • topology_settlement.tsExchange, filas, DLX, schemas e publish confirmado com mandatory
  • config_settlement.tsLê e valida a configuração no boot
  • demo.tsdemoSobe o worker e fecha uma data por contagem, outra por timeout e manda uma mensagem inválida à DLQ
  • aggregation_settlement.test.tstesteUnitário: fronteira da contagem e da janela, relatório único, reentrega
  • worker_aggregator.test.tstesteIntegração com RabbitMQ real: relatório, timeout, DLQ e mandatory

./

  • docker.shinfraSobe o RabbitMQ na porta 5672

Executar · com Docker

  1. ./docker.sh up# sobe RabbitMQ
  2. cp .env.example .env# variáveis de ambiente
  3. bun install# dependências
  4. bun run demo# roda o cenário

Testes: bun run test.

Requisitos
BunDocker
Sobe junto
RabbitMQ

Por que se relacionam

Teste rápido

Qual é o próximo passo depois de Aggregator?

Próximo projeto · Comunicação & MensageriaContent Enricher
Esc

↑ ↓ navegarEnter abrir191 resultados