basedosdados / basedosdados/pipelines

[chore] Dividir flows monolíticos em estágios (check_update / flow_download / mat_test)

Open
#1,867 0 comments 0 reactions 1 assignee Claimed by @Winzen View on GitHub
Dominant language
Python
Stars
49
Forks
22
Avg merge
19h 31m
Merged PRs (30d)
165

Description

> ⚠️ **Rascunho original (2026-08-28)** — proposta inicial pra divulgar e discutir a ideia. O restante deste documento é o texto original, mantido como registro — ver "Status atual" abaixo pra saber o que efetivamente mudou/foi decidido desde então.

---

## Status atual (2026-09-03)

O mecanismo central já está **implementado e testado de ponta a ponta em dev e em produção reais** (dados reais, BigQuery real, backend real — não simulado). PR: #1932.

**Mudança em relação à proposta original**: as etapas não são mais encadeadas por **Automação do Prefect 3** — são disparadas por uma chamada direta de `run_deployment()` no código do flow upstream (`timeout=0`, não bloqueia; `as_subflow=True`, aparece com lineage real no Prefect UI). Mesmo isolamento de recurso entre pods que a automação dava, mas sem precisar manter um mecanismo separado de sincronização de automações conforme datasets forem migrados, e com validação automática de parâmetro (sem o workaround de serializar tudo em JSON pra contornar o Jinja da automação). Comparação completa e motivo da troca documentados no PR.

`mat_test` é genérico, como a proposta original já sugeria — um deployment só, reaproveitado por todos os datasets. O boilerplate de `check_update`/`flow_download` (rename do flow run, poll/commit de coverage, dispatch pro próximo estágio) também foi encapsulado numa classe genérica reutilizável (`CheckThenDownloadPipeline`, em `pipelines/utils/stage_dispatch.py`) — cada dataset novo só fornece a lógica específica (como checar a fonte, como baixar).

**Pendências desta issue, atualizadas**:
- [x] Como o disparo passa parâmetros ao flow downstream — resolvido: `run_deployment(parameters={...})` aceita dict nativo, validado automaticamente pelo tipo do parâmetro do flow (não precisa mais de solução de payload/Jinja, já que a automação saiu da equação).
- [x] Estado de saída do `check_update` quando não há dado novo — resolvido: simplesmente não dispara nada.
- [x] Granularidade das tags — implementada a nível de dataset; datasets multi-tabela podem exigir revisitar.
- [x] Decisão: `mat_test` genérico ou por dataset — **genérico**, um deployment só.
- [ ] Impacto no CI de deploy com 3 flows por dataset — investigado e medido (deploy de prod hoje é `--all`, ~21min, dobraria com a migração completa); correção proposta (deploy seletivo, como o staging já faz) em #1943.

**Novidades não previstas na proposta original**:
- Dois bugs reais achados e corrigidos no caminho — um específico do piloto (coverage não coagido corretamente), outro **repo-wide**, afetando ~63 flows reais já em produção (`rename_flow_run_dataset_table` nunca executava de verdade) — issue #1940.
- Resolvido também o problema (não antecipado na proposta original) do staging dev/prod ficar partido entre pods diferentes quando `flow_download` e `mat_test` rodam em pools distintos — reaproveitando `transfer_files_to_prod_flow`, que já existia no repo mas nunca era chamada.
- Testada pela primeira vez a variante com **dados particionados** (`ano=/mes=`), exercitando `DownloadResult.partition_folders`/`transfer_files_to_prod_flow(folders=...)` — caminho que o piloto original (arquivo único) nunca cobriu. 4 problemas reais encontrados e corrigidos no caminho, 2 deles genéricos e não específicos deste piloto: (1) `poll_source_for_update` quebra pra **qualquer** tabela sem `RawDataSource` vinculado — o filtro GraphQL `rawDataSource_Id: null` é ignorado em vez de filtrar por nulo, retornando todos os `Poll`s do backend (ainda sem issue própria); (2) `bd.Table.create()` só registra o ponteiro externo do BigQuery em `basedosdados-staging`, nunca em `basedosdados-dev` — o `target=dev` do dbt espera achar a tabela lá, exigindo criação manual não documentada até então (issue #1967).
- Os dois pilotos de teste (`test_event_pipeline`, `test_event_pipeline_partitioned`) foram consolidados dentro de `pipelines/datasets/test_dataset/` (mesmo `dataset_id` de sandbox no backend — não fazia sentido cada um numa pasta de topo própria). Isso expôs e corrigiu uma limitação real de `deployment_name()`/`CheckThenDownloadPipeline`, que assumia um dataset por `flows.py`; agora suporta múltiplos pilotos/tabelas no mesmo arquivo via um parâmetro `flow_download_deployment` derivado do nome real da função (`.fn.__name__`), não uma string repetida solta.
- A variante `check_and_download` (49% dos ~82 datasets, conforme o levantamento abaixo) segue **não implementada** — só o caminho de 3 flows tem código real até agora.

Documentação completa (log cronológico da implementação, comparações arquiteturais, fluxogramas): PR #1932 e branch `feat/event-pipeline-automations-poc`.

---

## Contexto

Os flows atuais são monolíticos: um único flow executa check de atualização, download, upload para GCS, materialização em dev, testes, materialização em prod e atualização de metadados — tudo em sequência. Isso gera quatro problemas:

1. **Desperdício de recurso** — pods de 4 GB sobem só para verificar se o dado precisa ser atualizado. Na maioria das execuções a resposta é "não", sem justificar o custo.
2. **Recursos idênticos para cargas diferentes** — check (leve), download (médio) e materialização dbt (pesado) rodam no mesmo pod com a mesma alocação.
3. **Cota BQ compartilhada** — o gate de dev consome a cota de `basedosdados-dev` junto com o desenvolvimento manual (issue #1767).
4. **Reruns custosos** — falhar em prod obriga a rerrodar tudo desde o check, mesmo que o dado já esteja no storage.

Além disso, todo flow repete o mesmo boilerplate (check, gate dev, materialização prod, metadata) — código que não tem nada de específico do dataset em questão.

## Proposta

Quebrar o flow monolítico em **três flows independentes**, conectados por automação de disparo (originalmente proposto como Automação do Prefect 3 disparada por evento — na implementação final, trocado por `run_deployment()` direto no código, ver "Status atual" acima):

```
[check_update]
│ sucesso → dispara o próximo

[flow_download]
│ sucesso → dispara o próximo

[mat_test] ← dbt run+test em dev → dbt run+test em prod → atualiza metadados
```

## Identificação por tags

Cada flow recebe duas tags no deploy:
- **Etapa:** `check_update`, `flow_download`, `mat_test`
- **Dataset:** ex: `dataset:br_tse_filiacao_partidaria`

Na proposta original, as automações filtravam pelo par `(etapa_downstream, dataset)` pra disparar o flow certo. Na implementação final, essas tags não são mais necessárias pro disparo em si (`run_deployment()` resolve o deployment por nome, não por tag) — continuam existindo só como convenção de organização/descoberta no Prefect UI.

## Responsabilidade de cada flow

| Flow | O que faz | Recursos |
|---|---|---|
| `check_update` | Consulta o backend/fonte para ver se há dado novo | Mínimo (0.5 CPU / 512 MB) |
| `flow_download` | Baixa os dados e faz upload para GCS staging | Médio (1 CPU / 2 GB) |
| `mat_test` | `dbt run+test` em dev → `dbt run+test` em prod → atualiza metadados | Alto (2 CPU / 4 GB) |

**Decisão de design:** dev e prod rodam no mesmo pod dentro do `mat_test`. Manter dois flows separados (`dev_mat_test` e `prod_mat_test`) permitiria rerrodar só prod se prod falhar, mas esse caso é raro na prática. O overhead de um startup de pod extra em todo pipeline com dado novo não se justifica. A isolação de cota BQ entre dev e prod é controlada pelo `execution_project` do dbt target — não depende de pods separados.

## Flows genéricos vs. específicos

O único flow que varia por dataset é o `flow_download` — onde fica a lógica de como baixar aquele dado específico. Os demais (`check_update`, `mat_test`) podem ser **flows genéricos reutilizáveis**, parametrizados por `dataset_id` e `table_id`.

Isso reduz drasticamente o número de deployments: em vez de N flows completos, passamos a ter N `flow_download` + 2 flows genéricos compartilhados.

(Na implementação final, `check_update` acabou ficando um flow por dataset, não genérico — só `mat_test` é de fato compartilhado. Ver "Status atual".)

## Variantes do pipeline

A arquitetura suporta duas variantes, pois alguns datasets precisam baixar o arquivo para checar se há atualização:

### Variante padrão
Para datasets onde o check é uma chamada leve (HEAD request, API de metadata, listagem FTP):
```
[check_update] → [flow_download] → [mat_test]
```

### Variante check_and_download
Para datasets que precisam baixar o arquivo da fonte para descobrir se há dado novo (não existe endpoint de versão/hash):
```
[check_and_download] → [mat_test]
```
O `check_and_download` funde check + download em um único flow: baixa, verifica se é novo, e faz upload para GCS. Se não for novo, conclui sem disparar o downstream.

**Tradeoff:** separar em dois flows forçaria o `flow_download` a re-baixar do GCS (não da fonte), o que é mais rápido e confiável, mas adiciona um step desnecessário. Para esses casos, fundir é mais limpo.

**Status**: não implementada ainda — só a variante padrão (3 flows) tem código real, ver "Status atual" acima.

## Levantamento dos crawlers existentes

> Levantamento original (2026-08-28), mantido como registro histórico: analisados ~82 datasets reais do diretório `datasets/`, contagem por variante — Padrão (check via API/metadata leve) 28 (51%), check_and_download (baixa pra checar) 27 (49%), Sem check (sempre executa) 13, Inativos (sem flows.py) 12. **As duas variantes principais são igualmente comuns — ambas precisam ser tratadas como cidadãs de primeira classe.**

### Recontagem nominal (2026-09-03)

Refeito o levantamento pra ter a lista nominal completa (não só contagens), necessária pra escolher candidatos reais de migração. Números não batem exatamente com o original (33/26/10/12 nesta recontagem vs. 28/27/13/12) — provavelmente por datasets `au_*` adicionados ao repo depois de 28/08, e por algumas funções com nome de check leve que na verdade baixam o arquivo inteiro antes de extrair a data. Tratando esta recontagem como a mais confiável.

#### Padrão — check leve, sem baixar o arquivo completo (33)

Candidatos primários de migração — o check já é barato, só falta encapsular em `check_update`/`flow_download` via `CheckThenDownloadPipeline`.

| Dataset | Técnica de check | Onde |
|---|---|---|
| `au_ato_abr` | HTTP HEAD (`source_last_modified`) | `tasks.py:15` |
| `au_geoscape_gnaf` | API CKAN (`resolve_source`) | `tasks.py:15` |
| `br_ans_beneficiario` | scrape leve de listagem HTML/FTP (`extract_links_and_dates`) | `crawler/ans_beneficiario/tasks.py:26` |
| `br_bcb_agencia` | API de metadados BCB (`get_documents_metadata`) | `crawler/bcb_agencia/tasks.py:38` |
| `br_bcb_estban` | API de metadados BCB (`get_documents_metadata`) | `crawler/bcb_estban/tasks.py:35` |
| `br_bcb_ifdata` | API de índice de competências (`source_max_period`/`fetch_index`) | `datasets/br_bcb_ifdata/utils.py:79` |
| `br_bcb_sicor` | listagem de links da fonte antes do download (`search_sicor_links`) | `crawler/bcb/flows.py::_run_bcb_sicor` |
| `br_bndes_operacoes_contratadas` | API/metadado (`get_source_max_date`) | `crawler/bndes/flows.py:68` |
| `br_camara_dados_abertos` | checagem de URL (`check_if_url_is_valid`) | `crawler/camara_dados_abertos/flows.py:35` |
| `br_cnj_improbidade_administrativa` | contagem BQ + scrape leve de página (`is_up_to_date`) | `crawler/cnj_improbidade_administrativa/tasks.py:238` |
| `br_cvm_fi` | scrape de listagem (`extract_links_and_dates`) | `crawler/cvm/flows.py:43` |
| `br_denatran_frota` | API do próprio backend BD (`get_api_most_recent_date`) | `tasks.py:325` |
| `br_ibge_inpc` | API IBGE (`get_date_api`) | `crawler/ibge_inflacao/tasks.py:21` |
| `br_ibge_ipca` | API IBGE (`get_date_api`) | `crawler/ibge_inflacao/tasks.py:21` |
| `br_ibge_ipca15` | API IBGE (`get_date_api`) | `crawler/ibge_inflacao/tasks.py:21` |
| `br_inmet_bdmep` | listagem/API leve (`extract_last_date_from_source`) | `flows.py` |
| `br_me_caged` | API/metadado leve (`get_source_last_date`) | `flows.py` |
| `br_me_cnpj` | leitura de índice da API, sem baixar arquivos (`data_url`) | `crawler/me_cnpj/tasks.py:29` |
| `br_me_comex_stat` | scrape leve de página (`parse_last_date`) | `crawler/me_comex_stat/tasks.py:27` |
| `br_mf_divida_ativa` | probe leve da fonte (`latest_available_quarter`) | `tasks.py:21` |
| `br_mp_pep` | checagem de página via Selenium, sem baixar dado (`is_up_to_date`) | `crawler/mp_pep/tasks.py:431` |
| `br_ms_cnes` | listagem FTP (nomes de arquivo, sem baixar conteúdo) | `crawler/datasus/tasks.py:114` |
| `br_ms_sia` | listagem FTP (nomes de arquivo, sem baixar conteúdo) | `crawler/datasus/tasks.py:114` |
| `br_ms_sih` | listagem FTP (nomes de arquivo, sem baixar conteúdo) | `crawler/datasus/tasks.py:114` |
| `br_ms_sinan` | listagem FTP (nomes de arquivo, sem baixar conteúdo) | `crawler/datasus/tasks.py:114` |
| `br_rf_cafir` | API de metadados (`task_parse_api_metadata`/`task_get_last_update_date`) | `flows.py` |
| `br_rf_cno` | check leve (`check_need_for_update`) | `crawler/rf/flows.py:43` |
| `br_rf_cnpj` | leitura de índice da API (`data_url`) | `crawler/rf_cnpj/tasks.py:32` |
| `br_sfb_sicar` | 1 page fetch — comentário explícito no código ("Cheap... before downloading gigabytes") | `flows.py` |
| `us_bls_oes` | leitura de página HTML (`resolve_latest_year`) | `tasks.py:17` |
| `us_cfpb_hmda` | leve (`resolve_years`/`latest_source_year`) | `tasks.py:12` |
| `us_sec_edgar` | API/listagem leve (`resolve_latest_quarter`) | `flows.py` |

#### check_and_download — precisa baixar o arquivo pra checar (26)

**Critério de reclassificação, limiar de 5 GB**: dentre os `check_and_download`, qualquer um cujo download de check fica **abaixo de 5 GB** é barato o suficiente pra ser tratado como candidato de migração tão leve quanto a categoria "Padrão" — baixar o arquivo e conferir o que aconteceu nele não justifica tratamento especial. Só os que de fato passam de 5 GB (ou onde não dá pra confirmar que ficam abaixo disso) continuam como "pesados de verdade", exigindo desenho mais cuidadoso (`check_and_download` fundido, streaming parcial, ou HEAD/ETag na origem antes de baixar). Onde a evidência é textual ("multi-GB", "several hundred MB") sem número exato, marcado como "verificar" em vez de assumir.

| Dataset | Técnica de check | Tamanho (evidência) | < 5 GB? |
|---|---|---|---|
| `br_anatel_banda_larga_fixa` | unzip pra checar data | ~1 GB (README) | ✅ sim |
| `br_anatel_telefonia_movel` | unzip pra checar data | `memory_limit: 8Gi` no job (mesma família, ~1 GB provável) | ⚠️ verificar |
| `br_cgu_beneficios_cidadao` | baixa antes do poll (ZIP assíncrono, Portal da Transparência) | não documentado | ⚠️ verificar |
| `br_cgu_cartao_pagamento` | baixa antes do poll (mesma família de portal) | não documentado | ⚠️ verificar |
| `br_cgu_licitacao_contrato` | baixa antes do poll (mesma família de portal) | não documentado | ⚠️ verificar |
| `br_cgu_sancoes` | baixa antes do poll, snapshot cumulativo | ZIP inteiro em memória, sem tamanho documentado | ⚠️ verificar |
| `br_tse_eleicoes` | baixa ZIPs pra extrair data máxima | ZIPs nacionais por tipo de dado, escopo Brasil inteiro | ⚠️ verificar |
| `br_stf_corte_aberta` | Selenium dispara download automático | não documentado | ⚠️ verificar |
| `us_bls_qcew` | baixa ZIP trimestral | "multi-GB CSVs" no total; ZIP trimestral recente > 200 MB | ⚠️ verificar |
| `us_fec_campaign_finance` | baixa arquivo de contribuições | atual ~2 GB comprimido; ciclo completo (`indiv20`) 5.6 GB | ❌ não |
| `us_bea` | API JSON por série | não é bundle, é por série individual | ✅ sim |
| `us_fed_fred` | API JSON por série (`file_type=json`) | mesmo padrão do `us_bea` | ✅ sim |
| `world_cricsheet` | baixa bundle compactado | download ~114 MB (extração em disco chega a "several GB", pós-download) | ✅ sim |
| `world_wb_wdi` | baixa arquivo | ~270 MB | ✅ sim |
| `au_abs_cpi` | baixa `.xlsx` por tabela | não documentado (provável pequeno) | ⚠️ verificar |
| `au_abs_labour_force` | baixa SDMX + Excel | ~38 MB | ✅ sim |
| `au_ato_taxation_statistics` | baixa dado real antes de decidir | ~92 MB | ✅ sim |
| `au_rba_statistical_tables` | baixa múltiplos CSVs individuais | não documentado, mas CSVs individuais | ✅ sim (provável) |
| `br_anp_precos_combustiveis` | baixa recorte de 4 semanas (`ultimas-4-semanas-*.csv`) | recorte pequeno, não histórico | ✅ sim |
| `br_bcb_taxa_cambio` | API JSON por intervalo de datas | pequeno | ✅ sim |
| `br_bcb_taxa_selic` | mesmo padrão de API JSON/CSV do `taxa_cambio` | pequeno | ✅ sim |
| `br_cgu_emendas_parlamentares` | baixa `emendas_parlamentares.zip` | não documentado | ⚠️ verificar |
| `br_cgu_servidores_executivo_federal` | ZIP assíncrono por subsistema/mês (README confirma) | não documentado | ⚠️ verificar |
| `br_sedec_desastres` | concatenação de 27 downloads (1 export CSV por UF) | individualmente pequeno | ✅ sim |
| `mx_sesnsp_incidencia_delictiva` | baixa dado real antes de decidir | "several hundred MB" | ✅ sim |
| `us_bls_cpi` | baixa dado real antes de decidir | "several hundred MB" | ✅ sim |

**Resumo pós-reclassificação**: **18 dos 26** ficam com evidência de estarem abaixo de 5 GB (13 confirmados + 5 prováveis) e passam a ser candidatos de migração tão prioritários quanto a categoria "Padrão". Restam **8 incertos ou genuinamente pesados**, que exigem confirmação (teste real ou inspeção manual do portal) antes de decidir a estratégia: `br_anatel_telefonia_movel`, `br_cgu_beneficios_cidadao`, `br_cgu_cartao_pagamento`, `br_cgu_licitacao_contrato`, `br_cgu_sancoes`, `br_tse_eleicoes`, `br_stf_corte_aberta`, `us_fec_campaign_finance` (único com evidência concreta de passar dos 5 GB).

#### Sem check — sempre roda, sem gate de novidade (10)

`br_bd_indicadores`, `br_bd_siga_o_dinheiro`, `br_cgu_pessoal_executivo_federal`, `br_cvm_administradores_carteira`, `br_cvm_oferta_publica_distribuicao`, `br_me_rais`, `br_me_siconfi` (poll não-bloqueante, sempre reconstrói), `br_poder360_pesquisas`, `br_senado_dados_abertos`, `fundacao_lemann`.

#### Indeterminado (1)

`br_rj_isp_estatisticas_seguranca` — usa `get_count_lines`; não ficou claro se conta linhas de um arquivo já local ou baixa pra contar. Precisa inspeção manual.

#### Inativos — sem `flows.py` (13)

`br_b3_cotacoes`, `br_mercadolivre_ofertas`, `br_mg_belohorizonte_smfa_iptu`, `br_mp_pep_cargos_funcoes`, `br_ons_avaliacao_operacao`, `br_ons_estimativa_custos`, `br_senado_dados_abertos_administrativos`, `br_sp_saopaulo_dieese_icv`, `br_tse_filiacao_partidaria`, `mundo_transfermarkt_competicoes`, `mundo_transfermarkt_competicoes_internacionais`, `world_sofascore_competicoes_futebol`, `world_wil_wid`.

#### Candidatos sugeridos pra primeira migração

Baixo risco pra começar: dataset único (sem multi-tabela), check simples, sem particionamento — `br_ibge_ipca`/`br_ibge_ipca15`/`br_ibge_inpc` (mesma API IBGE, `check_fn` quase copiar-colar), `br_denatran_frota` (API do próprio backend BD), `us_sec_edgar` (API/listagem leve).

Depois do limiar de 5 GB, o universo de candidatos leves cresce bastante: os 18 `check_and_download` reclassificados também viram elegíveis, não só os 33 "Padrão".

## Impacto na organização das pastas

Esta refatoração deve ser feita **antes** da reorganização de `crawler/ → datasets/` (issue #1705), pois define o formato final de como cada flow vai ser escrito. Mover as pastas antes resultaria em retrabalho.

## Relacionado

- #1767 — isolamento de cota BQ entre desenvolvimento e gate de prod
- #1705 — reorganização da pasta `crawler/`
- #1768 — logs de `_upload_to_gcs` não indicam o ambiente
- #1940 — bug repo-wide: `rename_flow_run_dataset_table` nunca renomeia o flow run
- #1943 — deploy de prod deveria ser seletivo, não `--all` em todo push

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.