-
Notifications
You must be signed in to change notification settings - Fork 0
Arquitetura
O pipeline foi desenhado seguindo o padrão Medallion (também conhecido como Multi-Hop Architecture), um padrão consagrado em engenharia de dados que organiza os dados em camadas de maturidade crescente.
┌──────────────────────────────────────────────────┐
│ ORQUESTRAÇÃO │
│ Dagster + Schedules │
│ ┌──────┐ ┌──────────┐ ┌──────────┐ ┌──────┐ │
│ │ */5 │ │ */15 │ │ 0 2 │ │ 0 3 │ │
│ │ pos. │ │ prev. │ │ compact. │ │ exp. │ │
│ └──┬───┘ └────┬─────┘ └────┬─────┘ └──┬───┘ │
└─────┼───────────┼──────────────┼───────────┼──────┘
│ │ │ │
▼ ▼ ▼ ▼
┌─────────────────────────────────────────────────────┐
│ BRONZE (Raw Storage) │
│ │
│ ┌───────────────────────────────────────────────┐ │
│ │ SQLite (padrão) / PostgreSQL (produção) │ │
│ │ • dados crus como chegam da API │ │
│ │ • INSERT OR IGNORE → idempotente │ │
│ │ • UNIQUE INDEX → sem duplicatas │ │
│ │ • janela quente: 7 dias (expurgo automático) │ │
│ └───────────────────────────────────────────────┘ │
└─────────────────────┬───────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────┐
│ SILVER (Camada Analítica) │
│ │
│ ┌───────────────────────────────────────────────┐ │
│ │ Parquet particionado por data (dt=YYYY-MM-DD)│ │
│ │ • DuckDB para transformação │ │
│ │ • OVERWRITE_OR_IGNORE → idempotente │ │
│ │ • schema consistente, colunas padronizadas │ │
│ │ • otimizado para consultas analíticas │ │
│ └───────────────────────────────────────────────┘ │
└─────────────────────┬───────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────┐
│ GOLD (Consumo) │
│ │
│ ┌──────────────┐ ┌──────────────────────────┐ │
│ │ analise_ │ │ dashboard_sptrans.py │ │
│ │ onibus.py │ │ (Streamlit) │ │
│ │ (CLI) │ │ • modo Parquet (padrão) │ │
│ │ • métricas │ │ • fallback SQLite │ │
│ │ • relatórios │ │ • visualização interativa│ │
│ └──────────────┘ └──────────────────────────┘ │
└─────────────────────────────────────────────────────┘
Cada camada resolve um problema diferente:
| Camada | Problema | Solução |
|---|---|---|
| Bronze | Dados crus precisam ser preservados para auditoria e reprocessamento | SQLite/PostgreSQL mantém o registro original |
| Silver | Consultar 555k linhas por dia em SQLite é lento para dashboards | Parquet + DuckDB é ~10× mais rápido em consultas analíticas |
| Gold | Consumidores precisam de dados prontos, sem transformação | Queries prontas e dashboard com fallback |
Idempotência significa que executar a mesma operação duas vezes produz o mesmo resultado. É um dos princípios mais importantes em pipelines de dados — sem ele, um restart de job ou uma falha de rede pode duplicar registros silenciosamente.
┌─────────────────────────────┐
│ IDEMPOTÊNCIA │
│ (executar N× = mesmo resultado) │
└─────────────────────────────┘
│
┌─────────────────┼─────────────────┐
▼ ▼ ▼
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ INSERÇÃO │ │ COMPACTAÇÃO │ │ EXPURGO │
│ INSERT OR │ │ COPY TO │ │ DELETE por │
│ IGNORE + │ │ OVERWRITE_OR │ │ janela + │
│ UNIQUE INDEX │ │ _IGNORE │ │ backup │
│ → roda N× │ │ → roda N× │ │ → roda N× │
│ sem duplic.│ │ mesmo parq.│ │ mesmo res. │
└──────────────┘ └──────────────┘ └──────────────┘
-- Se o coletor executar duas vezes no mesmo minuto,
-- a segunda execução simplesmente ignora o INSERT duplicado
INSERT OR IGNORE INTO posicoes (timestamp_coleta, id_onibus, ...)
VALUES ('2025-08-18 10:30:00', 12345, ...);Um desafio comum em pipelines de dados é usar SQLite em desenvolvimento (leve, zero config) e PostgreSQL em produção (concorrente, robusto).
O módulo src/database.py resolve isso com uma abstração de backend:
from src.database import get_connection, insert_sql
# Ambiente: DATABASE_URL não definida → SQLite
# Ambiente: DATABASE_URL=postgresql://... → PostgreSQL
with get_connection() as conn:
cursor = conn.cursor()
cursor.executemany(insert_sql("posicoes", columns), dados)Diferenças abstraídas:
| Aspecto | SQLite | PostgreSQL |
|---|---|---|
| Placeholder | ? |
%s |
| INSERT idempotente | INSERT OR IGNORE |
INSERT ... ON CONFLICT DO NOTHING |
| Primary Key | INTEGER PRIMARY KEY AUTOINCREMENT |
SERIAL PRIMARY KEY |
| Conexão | Arquivo local | TCP |
O Dagster coordena a execução dos assets e schedules:
┌─────────────────┐
│ Definitions │
│ (ponto de │
│ entrada) │
└────────┬────────┘
│
┌──────────────┼──────────────┐
▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌──────────┐
│ Assets │ │ Schedules│ │ Jobs │
│ (6) │ │ (4) │ │ (4) │
└──────────┘ └──────────┘ └──────────┘
Dependências entre assets (quem depende de quem):
posicoes_sptrans (coleta)
└── compactar_posicoes (transforma em Parquet)
└── expurgar_posicoes (limpa janela quente)
previsoes_sptrans (coleta)
└── compactar_previsoes (transforma em Parquet)
└── expurgar_previsoes (limpa janela quente)
Cada schedule executa um job no cron especificado:
| Schedule | Cron | Job |
|---|---|---|
| Coleta posições |
*/5 * * * * (a cada 5 min) |
coleta_posicoes_job |
| Coleta previsões |
*/15 * * * * (a cada 15 min) |
coleta_previsoes_job |
| Compactação |
0 2 * * * (02:00 diário) |
compactacao_job |
| Expurgo |
0 3 * * * (03:00 diário) |
expurgo_job |
| Decisão | Alternativa Rejeitada | Motivo |
|---|---|---|
| SQLite como bronze padrão | PostgreSQL desde o início | Simplicidade local; PostgreSQL fica opcional via DATABASE_URL
|
| DuckDB + Parquet como silver | Arrow IPC / Feather | DuckDB lê Parquet nativamente, suporta sqlite_scan() e particionamento |
| Dagster como orquestrador | Apache Airflow | Mais leve, asset-centric, schedules nativas, menos sobrecarga operacional |
| INSERT OR IGNORE para dedup | Remover duplicatas pós-facto | Prevenir é mais barato que remediar — índice UNIQUE impede a entrada |
database.py como abstração |
ORM (SQLAlchemy) | Menos dependências, controle explícito da sintaxe, 190 linhas vs framework inteiro |
A partir do Plano 004, o pipeline passou a registrar linhagem leve e aplicar verificações de integridade entre camadas.
Cada execução de asset registra uma linha na tabela lineage_audit (SQLite/PostgreSQL) com:
| Coluna | Descrição |
|---|---|
id |
Chave primária auto-incremento |
asset_name |
Nome do asset Dagster (ex: compactar_posicoes) |
table_name |
Tabela física envolvida (ex: posicoes) |
layer |
Camada: bronze (SQLite) ou silver (Parquet) |
run_timestamp |
Data/hora de execução |
row_count |
Quantidade de registros processados |
status |
ok, warning, error
|
-- Como consultar (qualquer momento)
SELECT asset_name, table_name, layer, row_count, status, run_timestamp
FROM lineage_audit
ORDER BY run_timestamp DESC
LIMIT 20;Dois AssetCheck reconciliam Bronze ↔ Silver após cada compactação, com tolerância de 5% (alarme, não barreira):
| Check | Verifica | Tolerância |
|---|---|---|
check_posicoes_bronze_silver |
count(posicoes) SQLite ≈ count(*) Parquet |
±5% |
check_previsoes_bronze_silver |
count(previsoes) SQLite ≈ count(*) Parquet |
±5% |
Por que não-bloqueante: dado ruim é logado e registrado em lineage_audit com status='warning', mas o pipeline não aborta — o alerta chega ao dashboard Dagster e aos logs. Se a divergência for sistemática (>5% consistente), vira issue a investigar.
Quatro Pydantic models formalizam o schema esperado por camada:
| Model | Camada | Restrições |
|---|---|---|
PosicaoBronze |
Bronze |
id_onibus > 0, latitude ∈ [-90, 90], longitude ∈ [-180, 180], letreiro_linha ≤ 20 chars
|
PrevisaoBronze |
Bronze |
id_linha > 0, id_onibus > 0, horario_previsao ≤ 30 chars
|
PosicaoSilver |
Silver | Tudo de Bronze + dt no formato ^\d{4}-\d{2}-\d{2}$ (partição) |
PrevisaoSilver |
Silver | Tudo de Bronze + dt no formato partição |
Validação em lote disponível via validar_lote(PosicaoBronze, registros).
| Decisão | Alternativa Rejeitada | Motivo |
|---|---|---|
MetadataValue.int nativo do Dagster |
OpenLineage + Marquez | Stack pesada (Java/Kafka/DB externo) para 2 tabelas |
lineage_audit no mesmo banco (SQLite) |
Catálogo externo (PostgreSQL/DataHub) | Mantém auditoria na mesma transação dos dados; sem dependência extra |
| Pydantic para schema contract | Great Expectations / Soda Core | Já temos Pydantic no projeto; cobre o caso com zero dependência nova |
| AssetCheck não-bloqueante | Bloquear pipeline em falha | Falha silenciosa em log é melhor que pipeline travado — alarme, não barreira |