Skip to content

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

Description

@Winzen

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. "<prefix><dataset_id>.<table_id>"). 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):

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:

# pyrefly: ignore [unused-coroutine]
rename_flow_run_dataset_table(
    prefix="...", dataset_id=dataset_id, table_id=table_id
)

por:

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

#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).

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions