Pular para o conteúdo

apps/messaging-streaming/scatter-gather

068 · Comunicação & Mensageria · ≈ 6 min de estudo

Scatter-Gather

Envia a mesma requisição em paralelo para múltiplos destinatários (scatter) e consolida as respostas em uma única decisão, aguardando todas chegarem ou um timeout — o que ocorrer primeiro (gather). Usado quando uma decisão depende de várias fontes independentes e nenhuma delas pode travar o processo inteiro se ficar fora do ar.

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

Cenário

Uma fintech consulta o score de um cliente em três bureaus (Serasa, Boa Vista, SPC) antes de conceder crédito. Cada bureau responde num tempo diferente e às vezes está fora do ar. A análise não pode esperar para sempre: decide com quem respondeu dentro do prazo e diz de quem ficou sem resposta.

Planta

Sequência
6/6
1 publish, correlationId c1, replyTo R, expiration 15001cópia2cópia3cópia (fica na fila e expira)4score, c15score, c16timeout vence: decide com 2 de 3 e devolve bureausMissing = [spc]API de crédito (gatherer)credit.score.request (fanout)SerasaBoa VistaSPC (fora do ar)fila de resposta do gatherer
6

6 passos — reproduza para seguir o fluxo

Como funciona

6 passos

Dispara para muitos, consolida respostas.

Esconde cada passo: lembre antes de tocar para revelar.

  1. 01

    A API publica a consulta uma vez numa exchange fanout; cada bureau tem sua fila ligada a ela e recebe uma cópia

  2. 02

    A consulta leva replyTo (uma fila exclusiva do gatherer, única para todos os lotes), um correlationId novo por lote e expiration igual ao timeout

  3. 03

    Cada bureau, um processo próprio, responde em replyTo ecoando o correlationId

  4. 04

    O gatherer junta as respostas do lote até chegarem todas ou vencer o timeout — o que vier primeiro — e limpa o estado do lote

  5. 05

    A consolidação é outra funçãoMédia de quem respondeu, veto se algum bureau negou, sem resposta nenhuma é negado

  6. 06

    A resposta traz bureausMissing; consulta presa na fila de um bureau parado expira e não é respondida quando ele volta

apps/messaging-streaming/scatter-gather

10 arquivos

src/

  • topology_credit.tsExchange fanout, filas dos bureaus e formatos de pedido e resposta
  • gather_credit.tsScatter e gather: fila de resposta única, lotes por correlationId, timeout e bureausMissing
  • consolidation_credit.tsRegra de decisão sobre as respostas recebidas
  • bureau_credit.tsBureau como processo próprio, com latência e score determinísticos
  • api_credit.tsAPI interna Elysia que roda a consulta e consolida
  • config_credit.tsVariáveis de ambiente validadas no boot
  • consolidation_credit.test.tstesteRegra: limite do score, veto, parcial, nenhuma resposta e ordem de chegada
  • gather_credit.test.tstesteIntegração com RabbitMQ real: todos respondem, bureau fora do ar, lotes concorrentes e expiração
  • demo.tsdemoTrês bureaus e a API como processos; consulta completa e consulta com o SPC parado

./

  • docker.shinfraSobe o RabbitMQ 3 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
  5. bun run test# integração contra o serviço real
Requisitos
BunDocker
Sobe junto
RabbitMQ

Por que se relacionam

Teste rápido

Qual destes combina com Scatter-Gather?

Próximo projeto · Comunicação & MensageriaEvent Streaming
Esc

↑ ↓ navegarEnter abrir191 resultados