-
Notifications
You must be signed in to change notification settings - Fork 0
Componentes
Catálogo completo de cada script, asset e artefato do pipeline, com exemplos de uso e explicação de funcionamento.
Função: Interface única para SQLite (dev/testes) e PostgreSQL (produção).
from src.database import get_connection, insert_sql, DB_PATH
# Conecta ao backend correto automaticamente
with get_connection() as conn:
cursor = conn.cursor()
cursor.executemany(
insert_sql("posicoes", columns),
[(ts, id_bus, linha, lat, lon, ts_gps)]
)O que abstrai:
- Formato de placeholders:
?(SQLite) vs%s(PostgreSQL) - INSERT idempotente:
INSERT OR IGNOREvsON CONFLICT DO NOTHING - Tipos de primary key:
INTEGER AUTOINCREMENTvsSERIAL - Gerenciamento de transações:
commit()erollback()automáticos
Quando usar: Sempre que precisar escrever no banco. Nunca use sqlite3 ou psycopg2 diretamente.
Função: Coleta coordenadas GPS dos ônibus em tempo real.
Frequência: A cada 5 minutos (via Dagster) ou 30 minutos (modo standalone).
Fluxo:
- Lê
config.ini→ token + linhas alvo - Carrega catálogo de linhas → obtém letreiros
- Autentica na API Olho Vivo
- Faz GET em
/Posicao→ recebe JSON com todos os ônibus - Filtra apenas os letreiros de interesse
- Insere no banco com
INSERT OR IGNORE
Destaque técnico — Filtro Inteligente:
# A API retorna todos os ônibus de SP.
# Filtramos apenas nossas linhas alvo para economizar espaço.
for linha in dados["l"]:
if linha.get("c") in letreiros_alvo: ← FILTRO
registros_para_salvar.append(...)Execução:
python src/coleta_sptrans.pyFunção: Coleta previsões de chegada dos ônibus nos pontos.
Frequência: A cada 15 minutos (via Dagster).
Diferença do coletor de posições:
- Requer
id_linha(não letreiro) para consultar previsões - Chama
/Previsao/Linha?id={id_linha}para cada linha alvo - Salva
id_paradaehorario_previsao
Por que menos frequente? Previsões de chegada mudam menos que posições GPS — não faz sentido coletar a cada 5 minutos.
Execução:
python src/coleta_previsoes.pyFunção: Cria as tabelas e índices necessários no banco.
python src/inicializar_banco.py # SQLite
DATABASE_URL=postgresql://... python src/inicializar_banco.py # PostgreSQLO que cria:
-
posicoes— posições GPS dos ônibus -
previsoes— previsões de chegada -
resultados_analise— resultados de análises executadas - Índices UNIQUE para dedup automático
Idempotente: Pode executar várias vezes — CREATE TABLE IF NOT EXISTS e CREATE INDEX IF NOT EXISTS.
Função: Transfere dados do SQLite para Parquet particionado via DuckDB.
# Exportar tudo
python src/compactar_parquet.py
# Exportar dia específico
python src/compactar_parquet.py --date 2025-08-18Mágica do DuckDB:
# DuckDB lê SQLite diretamente! Não precisa de dump intermediário.
con.execute("""
SELECT *, CAST(timestamp_coleta AS DATE) AS dt
FROM sqlite_scan('data/sptrans_data.db', 'posicoes')
""")
# Escreve Parquet particionado — cada data vira uma pasta
con.execute("""
COPY __temp TO 'data/parquet/posicoes'
(FORMAT PARQUET, PARTITION_BY (dt), OVERWRITE_OR_IGNORE)
""")OVERWRITE_OR_IGNORE: Se o Parquet do dia já existe, sobrescreve sem duplicar. Idempotente.
Função: Gera métricas operacionais: ônibus agrupados, stuck buses, etc.
python src/analise_onibus.py --mode parquet # via DuckDB (rápido)
python src/analise_onibus.py --mode sqlite # via SQLite (fallback)Métricas geradas:
- Bunched buses: dois ônibus da mesma linha a menos de 500m → indica falha de headway
- Stuck buses: ônibus que não se moveram por mais de 15 minutos
- Enriquecimento: cruza posições com previsões para montar tabela consolidada
Resultados salvos em resultados_analise no banco.
Função: Visualização interativa dos dados da frota.
streamlit run src/dashboard_sptrans.pyCaracterísticas técnicas:
- Fallback automático: tenta Parquet (rápido) → se falhar, usa SQLite
-
Cache:
@st.cache_dataevita reprocessar dados iguais - Sideband status: mostra qual modo está ativo (Parquet ou SQLite)
- Métricas visuais: agrupamento, stuck buses, mapa de calor
Função: Remove registros antigos para manter o banco leve.
python src/expurgar_sqlite.py # padrão 7 dias
python src/expurgar_sqlite.py --dias 30 # custom
python src/expurgar_sqlite.py --dry-run # simulaçãoPor que expurgar? O SQLite cresce ~4MB por dia. Em 1 ano seriam ~1.4GB. O Parquet mantém o histórico completo de forma mais eficiente.
Segurança: --dry-run mostra quantos registros seriam removidos sem deletar.
Função: Remove duplicatas existentes e cria índices UNIQUE em bancos legados.
python src/migrar_dedup.pyQuando usar: Se você já estava coletando dados ANTES da implementação da idempotência. Este script:
- Remove duplicatas mantendo apenas o registro com menor
id - Cria os índices UNIQUE
Função: Transfere dados do SQLite para PostgreSQL em lotes.
DATABASE_URL="postgresql://sptrans:sptrans_local@localhost:5432/sptrans" \
python src/migrar_postgres.pyCaracterísticas:
- Lê em lotes de 1000 registros (não carrega tudo na RAM)
- Usa
ON CONFLICT DO NOTHINGpara evitar duplicatas - Verifica contagens após migração
Assets definidos:
| Asset | O que faz | Grupo |
|---|---|---|
posicoes_sptrans |
Chama coleta_sptrans.job(letreiros_alvo)
|
coleta |
previsoes_sptrans |
Chama coleta_previsoes.job(session, linhas)
|
coleta |
Arquitetura: os assets envolvem as funções de src/ sem modificá-las. Isso permite que os scripts rodem tanto standalone quanto orquestrados.
@asset(group_name="coleta")
def posicoes_sptrans():
config = coleta_sptrans.get_config()
letreiros = coleta_sptrans.get_letreiros_alvo(
coleta_sptrans.get_linhas_alvo_ids(config)
)
coleta_sptrans.job(letreiros) # ← chama a função originalAssets definidos:
| Asset | Depende de | O que faz | Grupo |
|---|---|---|---|
compactar_posicoes |
posicoes_sptrans |
Exporta posições → Parquet | processamento |
compactar_previsoes |
previsoes_sptrans |
Exporta previsões → Parquet | processamento |
expurgar_posicoes |
compactar_posicoes |
Limpa posições antigas | processamento |
expurgar_previsoes |
compactar_previsoes |
Limpa previsões antigas | processamento |
Cadeia de dependências:
coleta → compactação → expurgo
(dados crus) → (Parquet) → (janela quente)
Ponto de entrada do Dagster. Exporta:
- 6 assets → coletores + processamento
- 4 jobs → coleta_posicoes, coleta_previsoes, compactacao, expurgo
- 4 schedules → executam os jobs nos intervalos definidos
defs = Definitions(
assets=[posicoes_sptrans, previsoes_sptrans, ...],
schedules=[posicoes_schedule, previsoes_schedule, ...],
)Schedule breakdown:
| Schedule | Cron | Executa |
|---|---|---|
| posicoes | */5 * * * * |
A cada 5 minutos |
| previsoes | */15 * * * * |
A cada 15 minutos |
| compactacao | 0 2 * * * |
02:00 diário |
| expurgo | 0 3 * * * |
03:00 diário |
| Arquivo | Testes | O que cobre |
|---|---|---|
test_inicializar_banco.py |
3 | Schema, índices, INSERT OR IGNORE |
test_coleta_sptrans.py |
9 | Autenticação, coleta, job com mocks |
test_coleta_previsoes.py |
7 | Autenticação, coleta por linha |
test_migrar_dedup.py |
1 | Remoção de duplicatas |
test_expurgar_sqlite.py |
4 | Expurgo, dry-run, tabela vazia |
test_compactar_parquet.py |
3 | Exportação Parquet com MonkeyPatch |
test_analise_onibus.py |
3 | Merge, rename, dados vazios |
test_dashboard.py |
5 | Stuck buses, bunched buses, enrich |
test_assets.py |
8 | Dagster: load, keys, deps, schedules, jobs |
test_postgres.py |
5 | Conexão, schema, migração (skipped sem DATABASE_URL) |
test_contracts.py |
9 | Pydantic models: tipos, ranges, validação em lote |
test_linhagem.py |
4 |
registrar_linhagem() + metadata nos assets |
test_checks.py |
4 | AssetCheck reconciliação Bronze↔Silver |
Total: 67 testes (62 ativos + 5 condicionais para PostgreSQL)
data/
├── sptrans_data.db ← SQLite (janela quente)
├── parquet/
│ ├── posicoes/
│ │ ├── dt=2025-08-13/
│ │ ├── dt=2025-08-14/ ← Parquet particionado
│ │ └── ...
│ └── previsoes/
│ └── ...
├── resultados_analise/ ← Parquet de análises (se exportado)
├── todas_as_linhas.csv ← Catálogo de linhas SPTrans
Função: Pydantic models que definem o schema esperado de cada camada (Bronze e Silver), com restrições de tipo e valor.
from src.contracts import PosicaoBronze, PrevisaoBronze, validar_lote
# Validação unitária
rec = PosicaoBronze(
timestamp_coleta=datetime(2025, 8, 13),
id_onibus=12345,
latitude=-23.55,
longitude=-46.63,
)
# Validação em lote (retorna lista de erros)
erros = validar_lote(PosicaoBronze, lista_de_dicts)
if erros:
logger.warning("%d registros rejeitados pelo contrato", len(erros))Restrições automáticas:
-
id_onibus > 0,id_linha > 0 -
latitude ∈ [-90, 90],longitude ∈ [-180, 180] -
letreiro_linha ≤ 20 chars,horario_previsao ≤ 30 chars -
dt(silver) no formatoYYYY-MM-DD(regex) - Campos extras são ignorados (
extra="ignore")
Função: Dagster AssetCheck que reconcilia contagem Bronze ↔ Silver após cada compactação. Tolerância de 5%.
@asset_check(asset=compactar_posicoes)
def check_posicoes_bronze_silver():
# Verifica se a quantidade exportada para Parquet
# corresponde ao que está no SQLite (Bronze)
...Saída na UI Dagster: badge ✅/❌ ao lado do asset, com descrição do diff.
Função: Histórico consultável de cada execução de asset (Bronze e Silver).
CREATE TABLE lineage_audit (
id INTEGER PRIMARY KEY AUTOINCREMENT,
asset_name TEXT NOT NULL,
table_name TEXT NOT NULL,
layer TEXT NOT NULL, -- 'bronze' ou 'silver'
run_timestamp TEXT NOT NULL,
row_count INTEGER NOT NULL,
status TEXT NOT NULL -- 'ok', 'warning', 'error'
);Como consultar:
-- Últimas 20 execuções
SELECT * FROM lineage_audit ORDER BY run_timestamp DESC LIMIT 20;
-- Contagem total por asset/camada
SELECT asset_name, layer, SUM(row_count) AS total
FROM lineage_audit
WHERE status = 'ok'
GROUP BY asset_name, layer;