basedosdados / basedosdados/pipelines
rename_flow_run_dataset_table nunca renomeia o flow run (async task chamada sem await)
- Dominant language
- Python
- Stars
- 49
- Forks
- 22
- Avg merge
- 19h 31m
- Merged PRs (30d)
- 165
Description
## Resumo
`rename_flow_run_dataset_table` (`pipelines/utils/tasks.py`) é uma `@task`
**async**. Em todo lugar do repositório onde é chamada, é chamada de dentro
de um `@flow` **síncrono**, sem `await`, e marcada com um
`# pyrefly: ignore [unused-coroutine]` — esse comentário não é um falso
positivo a suprimir: o pyrefly está certo, a coroutine é criada e
descartada sem nunca rodar.
Resultado: **o rename nunca acontece**. O flow run mantém o nome
auto-gerado do Prefect (ex. `tunneling-aardvark`) em vez do nome esperado
(ex. `"."`). Não quebra a execução (não dá
erro, não aparece nada no log — a task simplesmente não roda), então
passou despercebido: é cosmético, mas está quebrado desde que esse padrão
foi introduzido.
## Causa raiz
Confirmado lendo o código-fonte do Prefect 3.5.0
(`prefect.tasks.Task.__call__` → `prefect.task_engine.run_task`):
```python
if task.isasync and task.isgenerator:
return run_generator_task_async(**kwargs)
elif task.isgenerator:
return run_generator_task_sync(**kwargs)
elif task.isasync:
return run_task_async(**kwargs) # <- retorna uma coroutine, não o resultado
else:
return run_task_sync(**kwargs)
```
Para uma task async (`task.isasync == True`), `Task.__call__` sempre
retorna o que `run_task_async(...)` retorna — que é, ele mesmo, uma
coroutine. Chamar essa task solta, sem `await`, de um flow síncrono
(`@flow def algum_flow(...):`, não `async def`) cria esse objeto
`Coroutine` e descarta — a chamada RPC de verdade
(`client.update_flow_run(...)`, dentro do corpo da task) nunca chega a
executar.
## Evidência empírica
Confirmado contra um flow run real e `Completed` de produção:
`br_bcb_estban__municipio`, run `tunneling-aardvark`
(`01a059bb-dff4-7f63-b36d-96c059dc7c80`, completado
`2026-09-01T01:30:20Z`) — o nome nunca mudou, apesar do flow passar por
`rename_flow_run_dataset_table` no início da execução (mesmo padrão
quebrado que este flow usa).
Também reproduzido e corrigido durante o trabalho da #1867 (piloto de
pipeline orientado a eventos): o `mat_test_flow` genérico
(`pipelines/utils/metadata/flows.py`) tinha o mesmo padrão; sem log nenhum
de conclusão da task de rename, diferente de todas as outras tasks do
mesmo flow run. Depois da correção abaixo, a task passou a aparecer com
`Finished in state Completed()` no log e o flow run foi renomeado de
verdade (`"Mat Test: test_dataset.test_event_pipeline"`).
## Correção
Prefect já fornece o utilitário sancionado pra rodar uma coroutine de
dentro de código síncrono e esperar o resultado:
`prefect.utilities.asyncutils.run_coro_as_sync`.
Trocar o padrão atual:
```python
# pyrefly: ignore [unused-coroutine]
rename_flow_run_dataset_table(
prefix="...", dataset_id=dataset_id, table_id=table_id
)
```
por:
```python
from prefect.utilities.asyncutils import run_coro_as_sync
run_coro_as_sync(
rename_flow_run_dataset_table(
prefix="...", dataset_id=dataset_id, table_id=table_id
)
)
```
Já aplicado e validado em produção real em dois flows, como parte do
trabalho da #1867:
- `pipelines/utils/metadata/flows.py::mat_test_flow`
- `pipelines/utils/materialize_prod/flows.py::transfer_files_to_prod_flow`
## Escopo do que falta
O mesmo padrão quebrado (chamada solta + `# pyrefly: ignore
[unused-coroutine]`) ainda está presente em **63 arquivos** (76 call
sites no total, 73 ainda não corrigidos):
```
pipelines/crawler/anatel/banda_larga_fixa/flows.py
pipelines/crawler/anatel/telefonia_movel/flows.py
pipelines/crawler/bcb/flows.py
pipelines/crawler/bndes/flows.py
pipelines/crawler/camara_dados_abertos/flows.py
pipelines/crawler/cgu/flows.py
pipelines/crawler/cvm/flows.py
pipelines/crawler/cvm_administradores_carteira/flows.py
pipelines/crawler/datasus/flows.py
pipelines/crawler/fgv_igp/flows.py
pipelines/crawler/ibge_inflacao/flows.py
pipelines/crawler/isp/flows.py
pipelines/crawler/me_cnpj/flows.py
pipelines/crawler/me_rais/flows.py
pipelines/crawler/rf/flows.py
pipelines/crawler/rf_cnpj/flows.py
pipelines/crawler/tse_eleicoes/flows.py
pipelines/datasets/au_abs_cpi/flows.py
pipelines/datasets/au_abs_labour_force/flows.py
pipelines/datasets/au_ato_abr/flows.py
pipelines/datasets/au_ato_taxation_statistics/flows.py
pipelines/datasets/au_geoscape_gnaf/flows.py
pipelines/datasets/au_rba_statistical_tables/flows.py
pipelines/datasets/br_anp_precos_combustiveis/flows.py
pipelines/datasets/br_ans_beneficiario/flows.py
pipelines/datasets/br_bcb_agencia/flows.py
pipelines/datasets/br_bcb_estban/flows.py
pipelines/datasets/br_bcb_ifdata/flows.py
pipelines/datasets/br_bcb_sicor/flows.py
pipelines/datasets/br_bcb_taxa_cambio/flows.py
pipelines/datasets/br_bcb_taxa_selic/flows.py
pipelines/datasets/br_bd_indicadores/flows.py
pipelines/datasets/br_cgu_emendas_parlamentares/flows.py
pipelines/datasets/br_cgu_pessoal_executivo_federal/flows.py
pipelines/datasets/br_cgu_sancoes/flows.py
pipelines/datasets/br_cnj_improbidade_administrativa/flows.py
pipelines/datasets/br_cvm_oferta_publica_distribuicao/flows.py
pipelines/datasets/br_denatran_frota/flows.py
pipelines/datasets/br_ibge_pnadc/flows.py
pipelines/datasets/br_inmet_bdmep/flows.py
pipelines/datasets/br_me_caged/flows.py
pipelines/datasets/br_me_comex_stat/flows.py
pipelines/datasets/br_me_siconfi/flows.py
pipelines/datasets/br_mf_divida_ativa/flows.py
pipelines/datasets/br_mp_pep/flows.py
pipelines/datasets/br_poder360_pesquisas/flows.py
pipelines/datasets/br_rf_cafir/flows.py
pipelines/datasets/br_sedec_desastres/flows.py
pipelines/datasets/br_senado_dados_abertos/flows.py
pipelines/datasets/br_sfb_sicar/flows.py
pipelines/datasets/br_stf_corte_aberta/flows.py
pipelines/datasets/fundacao_lemann/flows.py
pipelines/datasets/mx_sesnsp_incidencia_delictiva/flows.py
pipelines/datasets/us_bea/flows.py
pipelines/datasets/us_bls_cpi/flows.py
pipelines/datasets/us_bls_oes/flows.py
pipelines/datasets/us_bls_qcew/flows.py
pipelines/datasets/us_cfpb_hmda/flows.py
pipelines/datasets/us_fec_campaign_finance/flows.py
pipelines/datasets/us_fed_fred/flows.py
pipelines/datasets/us_sec_edgar/flows.py
pipelines/datasets/world_cricsheet/flows.py
pipelines/datasets/world_wb_wdi/flows.py
```
Sugestão de abordagem: um codemod (ex. `sed`/script pequeno) que troca o
padrão em todos os arquivos de uma vez, com sua própria PR e CI passando —
mudança mecânica e de baixo risco (o pior caso de regressão é o rename
continuar não acontecendo, que já é o estado atual), mas grande demais em
diff pra ir junto de qualquer outra PR.
## Achado durante
basedosdados/pipelines#1867 (pipeline orientado a eventos com automações
Prefect 3) — documentado em detalhe em `staging-multi-ambiente.md`/
`issue-1867-pipeline-eventos.md` (docs locais da sessão que investigou).
Contributor guide
Research direction
Start with rename_flow_run_dataset_table in pipelines/utils/tasks.py and inspect the remaining call sites listed in the issue, especially their surrounding synchronous flows. Check Prefect's run_coro_as_sync utility, then search for the unused-coroutine pattern across the 63 files. Done means all 73 remaining call sites are handled consistently and CI passes.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- python
- Domain
- data-engineering
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Active
- Clarity
- Clearly specified
- Newbie friendliness
- 55/100