← Projetos
Em progresso

Streaming para alerta operacional

Alerta de fila e atraso em streaming: Kafka, Spark e o mesmo grão da auditoria.

  1. Fonte
  2. Kafka
  3. Spark
  4. Bronze
  5. Silver
  6. Gold
  7. Alertas

Competências demonstradas

  • Python
  • Spark
  • Kafka
  • Data Engineering
  • Kafka
  • Spark
  • Python
7 min de leitura

Recorte medido

Pedidos a cada 5 minutos
pedidos
Pedidos a cada 5 minutos 12:00 12:10 12:20 12:30 12:40 12:50
Pedidos a cada 5 minutos
12:00 18 pedidos
12:10 20 pedidos
12:20 22 pedidos
12:30 49 pedidos
12:40 28 pedidos
12:50 17 pedidos

Desenho: o job horário achata o pico. Não é medição de produção.

Quando o alerta chega
min
  • Batch horário 60
  • Streaming (alvo) 2
Quando o alerta chega
Batch horário 60 min
Streaming (alvo) 2 min

Alvo do desenho, não benchmark.

Arquitetura
  1. Origem

    API / CDC

    Por quê / trade-offs
    Por que foi escolhido
    O mock de delivery publica eventos em paralelo ao batch Parquet.
    Problema que resolve
    Batch horário não observa pico de 20 minutos.
    Trade-offs
    CDC real exige acesso à origem; no MVP, produtor ao lado do FastAPI.
  2. Event bus

    Kafka

    Por quê / trade-offs
    Por que foi escolhido
    Tópicos orders/status desacoplam produtor e consumidores.
    Problema que resolve
    Replay e múltiplos leitores (alerta vs. persistência).
    Trade-offs
    Operar Kafka tem custo. Redpanda local ou cluster gerenciado no início.
  3. Processamento

    Spark Structured Streaming

    Por quê / trade-offs
    Por que foi escolhido
    Agregações em janela (ex.: 5 min) e lag origem → processamento.
    Problema que resolve
    Contagem por minuto por loja e SLA projetado não cabem em job diário.
    Trade-offs
    Spark cobre janela e escala; o primeiro incremento pode ser um consumidor Python.
  4. Estado / sink

    Redis ou Postgres

    Por quê / trade-offs
    Por que foi escolhido
    Último estado por pedido e agregados para alerta.
    Problema que resolve
    Sem sink, o stream só emite log.
    Trade-offs
    Redis para hot state; Postgres se o consumo analítico precisar de SQL.
  5. Alertas

    Slack / e-mail

    Por quê / trade-offs
    Por que foi escolhido
    Sinal operacional com latência baixa.
    Problema que resolve
    UI complexa atrasa o aprendizado de idempotência e lag.
    Trade-offs
    Grafana ou WebSocket depois do alerta mínimo.

01 — Contexto

Segunda fase da plataforma de delivery: ingestão de eventos em fluxo para monitorar pedidos, fila e violações de SLA enquanto acontecem. A base batch está no Delivery Audit.

02 — Problema

  • Batch de hora em hora não observa pico de 20 minutos.
  • A operação precisa de alerta imediato (fila, loja saturada, atraso em cascata).
  • As mesmas regras de auditoria, em janela, mudam o tempo da decisão — não a semântica.

03 — Arquitetura

API / CDC → Kafka (orders, status) → Spark Streaming
    → agregações em janela (5 min) → sink (Redis/Postgres) → alertas

Idempotência e exactly-once são decisões explícitas. No MVP: at-least-once com deduplicação por order_id.

04 — Stack

  • Kafka (ou Redpanda no protótipo local) — log de eventos.
  • Spark Structured Streaming — janelas e lag; o primeiro incremento pode ser um consumidor Python.
  • Python — produtor ao lado do mock FastAPI.

05 — Pipeline

MVP:

  1. um tópico e um consumidor que persiste o último estado por pedido;
  2. regra: mais de N pedidos em preparo há M minutos → alerta;
  3. Slack ou e-mail.

O mock API publica em Kafka em paralelo ao batch Parquet. Regras de auditoria podem rodar em modo janela, versão simplificada.

06 — Modelagem

Estado por order_id (último status, timestamps da jornada) e agregados por loja em janela. O grão analítico do Gold batch permanece: streaming alimenta alerta; o fato diário continua sendo a fonte para tendência.

07 — Data Quality

No stream, qualidade começa por deduplicação e ordem aproximada. Sequência inválida continua relevante; watermark e atraso de evento passam a ser modos de falha.

08 — Observabilidade

Métricas de desenho: lag origem → processamento, tamanho da fila por loja, taxa de alerta.

09 — Trade-offs

  • Kafka fora do batch do Delivery Audit — o problema daquele recorte era auditoria em camadas, não latência de segundos.
  • Spark vs. consumidor Python — Spark justifica janela e escala; o primeiro passo é o consumidor mínimo, para não operar cluster sem regra.
  • Exactly-once vs. at-least-once + dedup — o segundo é operacionalmente mais barato no MVP.
  • Custo de infra — Redpanda local no protótipo.

10 — Resultado

Arquitetura, recorte de MVP e riscos documentados.

11 — O que deu errado

Tratar streaming como troca de ferramenta sobre o batch. O recorte correto é lag, estado por pedido e alerta — não replicar Gold em tempo real no dia um.

12 — Código

A base batch e o mock que vira produtor estão em Delivery-Projeto.