Streaming para alerta operacional
Alerta de fila e atraso em streaming: Kafka, Spark e o mesmo grão da auditoria.
- Fonte
- Kafka
- Spark
- Bronze
- Silver
- Gold
- Alertas
Competências demonstradas
- Python
- Spark
- Kafka
- Data Engineering
- Kafka
- Spark
- Python
Recorte medido
| 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.
- Batch horário 60
- Streaming (alvo) 2
| Batch horário | 60 min |
| Streaming (alvo) | 2 min |
Alvo do desenho, não benchmark.
-
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.
-
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.
-
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.
-
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.
-
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:
- um tópico e um consumidor que persiste o último estado por pedido;
- regra: mais de N pedidos em preparo há M minutos → alerta;
- 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.