Skip to content

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

Description

@laura-l-amaral

⚠️ 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-17)

Substitui o "Status atual (2026-09-03)" abaixo, mantido como histórico. PR #1932 atualizado com o estado atual completo — esta seção é um resumo, detalhes técnicos e tabela de testes estão lá.

  • 7 datasets migrados de verdade (54 flows): br_ans_beneficiario, br_ibge_ipca, br_inmet_bdmep, br_me_caged, br_me_comex_stat, br_ms_cnes, us_cfpb_hmda. Nome de deployment sem sufixo _flow (<dataset>_<table>_check_update, não ..._check_update_flow).
  • 3 datasets revertidos, fora do escopo deste PR: br_me_cnpj (decisão explícita, sem pendência aberta), br_sfb_sicar (desvio deliberado do padrão — 9 tabelas com teste dbt cruzado — ainda sem a revisão humana necessária antes de qualquer deploy real), br_denatran_frota (renomeado pra br_senatran_estatisticas em produção por outro PR, sem relação com este trabalho).
  • Branch sincronizada com main (estava 127 commits atrás) — trouxe junto a correção de uma colisão real de nomes de deployment entre os pools basedosdados/basedosdados-dev que chegou a esvaziar produção (bug: deploy de PR (label deploy-flow) rouba deployments de produção pro pool dev #2079, já corrigido em fix: corrige colisão de nomes dev/prod e paraleliza deploy_flows.py #2081) e a paralelização do deploy_flows.py (~3.4x mais rápido).
  • Pendências que continuam de pé: impacto no CI de deploy com o crescimento de flows por dataset (ver checklist original abaixo — parcialmente endereçado por Deploy de prod (cd-prefect3.yaml) deveria ser seletivo, não --all em todo push #1943/fix: corrige colisão de nomes dev/prod e paraleliza deploy_flows.py #2081, mas ainda não medido pra escala completa dos ~82 datasets); br_bcb_sicor continua bloqueado (comparação por tamanho de arquivo, não por data — stage_dispatch.py não suporta isso ainda); variante check_and_download continua não implementada.

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:

  • 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).
  • Estado de saída do check_update quando não há dado novo — resolvido: simplesmente não dispara nada.
  • Granularidade das tags — implementada a nível de dataset; datasets multi-tabela podem exigir revisitar.
  • 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 Deploy de prod (cd-prefect3.yaml) deveria ser seletivo, não --all em todo push #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 rename_flow_run_dataset_table nunca renomeia o flow run (async task chamada sem await) #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 Polls 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 bd.Table.create() só registra o ponteiro externo do BigQuery em basedosdados-staging, nunca em basedosdados-dev #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 [chore] Isolamento de cota BQ entre desenvolvimento e gate de prod #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

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

No labels
No labels

Type

No type

Projects

Milestone

No milestone

Relationships

None yet

Development

No branches or pull requests

Issue actions