basedosdados / basedosdados/pipelines

rename_flow_run_dataset_table nunca renomeia o flow run (async task chamada sem await)

Open
#1,940 0 comments 0 reactions 0 assignees View on GitHub
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

Open the contributing 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

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.