diff --git a/.claude/rules/prefect-pipeline-conventions.md b/.claude/rules/prefect-pipeline-conventions.md index eda694bd91..d8067f69bb 100644 --- a/.claude/rules/prefect-pipeline-conventions.md +++ b/.claude/rules/prefect-pipeline-conventions.md @@ -32,6 +32,21 @@ picks up `Flow` objects whose function is defined in that file (`obj.fn.__code__.co_filename` check). A factory that returns an inner `@flow` (see `br_ibge_ipca`) is fine because the inner fn is still defined in that file. +### `@flow` comes from `pipelines.utils.flow`, not from `prefect` + +```python +from pipelines.utils.flow import flow # NOT `from prefect import flow` +``` + +The deploy script reads two attributes off the flow object — `deploy_schedules` +and `job_variables` — that `prefect.Flow` does not declare, so setting them on a +plain Prefect flow is a Pyrefly `missing-attribute` error. `pipelines/utils/flow.py` +declares both on a `prefect.Flow` subclass and exports a `flow` decorator that +builds it; it takes the same arguments as `prefect.flow` and the object stays a +`prefect.Flow` for every `isinstance` check. Both attributes default to empty +(no schedule, work-pool default infrastructure). `deploy_schedules` is a list of +`prefect.schedules.Cron` — `Cron("0 16 10 * *", timezone="America/Sao_Paulo")`. + ## DRY with the onboarding code The cleaning transform lives in **one place** and is shared: @@ -275,8 +290,17 @@ on the **deployed Prefect worker** (its pod SA has access) — the local Schedule inline on the flow object (do NOT register storage/run-config by hand): ```python +from prefect.schedules import Cron + +from pipelines.utils.flow import flow + + +@flow(name="my_flow", log_prints=True) +def my_flow() -> None: ... + + my_flow.deploy_schedules = [ - {"cron": "0 16 10,11,12,13 * *", "timezone": "America/Sao_Paulo"} + Cron("0 16 10,11,12,13 * *", timezone="America/Sao_Paulo") ] my_flow.job_variables = { "memory": "8Gi" @@ -290,7 +314,7 @@ Deploy is CI, via `.github/scripts/deploy_flows.py`: looks fine. The workflow triggers on `labeled` and `synchronize`, so adding the label is itself enough. Schedules are **stripped** — manual runs only. - **Prod pool** (`cd-prefect3.yaml`, `--pool basedosdados --all`, on merge to main): - schedules become `Cron` objects; deployed **`paused=True`**. + `deploy_schedules` is passed straight to the deployment; deployed **`paused=True`**. - Cron in `America/Sao_Paulo`; see crontab.guru. For a monthly source, poll across a few release-window days — the source-poll guard no-ops until a new period lands. diff --git a/.github/scripts/deploy_flows.py b/.github/scripts/deploy_flows.py index 2b1af55d78..6bd2854f67 100644 --- a/.github/scripts/deploy_flows.py +++ b/.github/scripts/deploy_flows.py @@ -15,9 +15,9 @@ import sys from pathlib import Path -from prefect import Flow from prefect.runner.storage import GitRepository -from prefect.schedules import Cron + +from pipelines.utils.flow import Flow REPO_URL = "https://github.com/basedosdados/pipelines.git" @@ -69,23 +69,15 @@ def deploy_flow( entrypoint = f"{file_path}:{flow_name}" is_dev = "dev" in pool_name - schedules = getattr(flow, "deploy_schedules", None) - if is_dev: - schedules = None # flows em dev não têm schedule - elif schedules: - # Convert dict {"cron": "...", "timezone": "..."} to Cron schedule objects - schedules = [ - Cron(s["cron"], timezone=s.get("timezone", "UTC")) - if isinstance(s, dict) - else s - for s in schedules - ] + # flows em dev não têm schedule + schedules = None if is_dev else flow.deploy_schedules - job_variables = getattr(flow, "job_variables", None) + job_variables = flow.job_variables print(f" Registrando {flow_name} → {entrypoint}") try: + # pyrefly: ignore [missing-attribute] flow.from_source( source=GitRepository( url=REPO_URL, diff --git a/AGENTS.md b/AGENTS.md index 6372994909..395114b1ce 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -64,18 +64,25 @@ uv run manage.py add-pipeline ### File conventions -- `flows.py`: Define flows with `@flow`. Flows **must be defined at module level in this file** — `deploy_flows.py` only collects `Flow` objects whose function is defined there (an `obj.fn.__code__.co_filename` check). +- `flows.py`: Define flows with `@flow` from **`pipelines.utils.flow`**, never `prefect.flow` — the repo's decorator returns a `prefect.Flow` subclass that declares the deploy attributes (`deploy_schedules`, `job_variables`), which the Prefect class does not, so setting them on a plain `prefect.Flow` is a Pyrefly `missing-attribute` error. Flows **must be defined at module level in this file** — `deploy_flows.py` only collects `Flow` objects whose function is defined there (an `obj.fn.__code__.co_filename` check). - `tasks.py`: Define tasks with `@task`. - `constants.py`: Use a `constants` enum or plain constants — no hardcoded values elsewhere. - `utils.py`: Pure helper functions with no Prefect decorators. -There is no `schedules.py`. Attach the schedule to the flow object in `flows.py`; CI turns -these dicts into `Cron` objects at deploy time: +There is no `schedules.py`. Attach the schedule to the flow object in `flows.py`, as +`Cron` objects from `prefect.schedules` (the `timezone` is an argument of `Cron`): ```python -my_flow.deploy_schedules = [ - {"cron": "0 16 10 * *", "timezone": "America/Sao_Paulo"} -] +from prefect.schedules import Cron + +from pipelines.utils.flow import flow + + +@flow(name="my_flow", log_prints=True) +def my_flow() -> None: ... + + +my_flow.deploy_schedules = [Cron("0 16 10 * *", timezone="America/Sao_Paulo")] my_flow.job_variables = { "memory": "8Gi" } # optional; size to the flow's peak RAM diff --git a/pipelines/crawler/ibge_inflacao/flows.py b/pipelines/crawler/ibge_inflacao/flows.py index 273f33ea9c..7cbda8ebac 100644 --- a/pipelines/crawler/ibge_inflacao/flows.py +++ b/pipelines/crawler/ibge_inflacao/flows.py @@ -3,13 +3,12 @@ Prefect 3 — use os flows dos datasets (br_ibge_ipca, br_ibge_inpc) para deploy. """ -from prefect import flow - from pipelines.crawler.ibge_inflacao.tasks import ( check_for_updates, collect_data_utils, json_to_csv, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import DateFormat, PartBdpro, YearMonth from pipelines.utils.metadata.tasks import ( commit_source_update_task, @@ -136,5 +135,4 @@ def ibge_inflacao_flow( ) -# pyrefly: ignore [missing-attribute] ibge_inflacao_flow.deploy_schedules = [] diff --git a/pipelines/datasets/au_abs_cpi/flows.py b/pipelines/datasets/au_abs_cpi/flows.py index 8f6329f750..8d137f3296 100644 --- a/pipelines/datasets/au_abs_cpi/flows.py +++ b/pipelines/datasets/au_abs_cpi/flows.py @@ -15,10 +15,11 @@ import shutil import tempfile -from prefect import flow +from prefect.schedules import Cron from pipelines.datasets.au_abs_cpi.constants import constants from pipelines.datasets.au_abs_cpi.tasks import clean_cpi, download_cpi +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( AllFree, DateFormat, @@ -173,7 +174,6 @@ def au_abs_cpi_flow( # ABS publishes the monthly CPI in the last week of each month (moving to the # 4th Wednesday from Feb 2027). Poll across the last week at 16:00 BRT; the # source-poll guard no-ops until a new month lands. -# pyrefly: ignore [missing-attribute] au_abs_cpi_flow.deploy_schedules = [ - {"cron": "0 16 22,23,24,25,26,27,28 * *", "timezone": "America/Sao_Paulo"} + Cron("0 16 22,23,24,25,26,27,28 * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/au_abs_labour_force/flows.py b/pipelines/datasets/au_abs_labour_force/flows.py index ac7c428d40..144d8d5b35 100644 --- a/pipelines/datasets/au_abs_labour_force/flows.py +++ b/pipelines/datasets/au_abs_labour_force/flows.py @@ -19,7 +19,7 @@ import shutil import tempfile -from prefect import flow +from prefect.schedules import Cron from pipelines.datasets.au_abs_labour_force.constants import constants from pipelines.datasets.au_abs_labour_force.tasks import ( @@ -28,6 +28,7 @@ download_sdmx_task, latest_month_task, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, FreeLag, @@ -191,10 +192,8 @@ def au_abs_labour_force_flow( # ABS releases Labour Force monthly, on a Thursday roughly the 3rd-4th week, at # 11:30 Canberra time. Poll daily across that window at 06:00 BRT (= evening AEST, # after the morning release); the source-poll guard no-ops until a new month lands. -# pyrefly: ignore [missing-attribute] au_abs_labour_force_flow.deploy_schedules = [ - {"cron": "0 6 14-27 * *", "timezone": "America/Sao_Paulo"} + Cron("0 6 14-27 * *", timezone="America/Sao_Paulo") ] # openpyxl reads the ~38 MB SEM1 pivot; give the worker headroom. -# pyrefly: ignore [missing-attribute] au_abs_labour_force_flow.job_variables = {"memory": "6Gi"} diff --git a/pipelines/datasets/br_anatel_banda_larga_fixa/flows.py b/pipelines/datasets/br_anatel_banda_larga_fixa/flows.py index 33d3795dc1..8138213c4b 100644 --- a/pipelines/datasets/br_anatel_banda_larga_fixa/flows.py +++ b/pipelines/datasets/br_anatel_banda_larga_fixa/flows.py @@ -5,11 +5,12 @@ deste diretório. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.anatel.banda_larga_fixa.flows import ( _run_anatel_banda_larga_fixa, ) +from pipelines.utils.flow import flow def _anatel_blf_flow(table_id: str, cron: str | None): @@ -39,9 +40,8 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] _flow.deploy_schedules = ( - [{"cron": cron, "timezone": "America/Sao_Paulo"}] if cron else [] + [Cron(cron, timezone="America/Sao_Paulo")] if cron else [] ) return _flow diff --git a/pipelines/datasets/br_anatel_telefonia_movel/flows.py b/pipelines/datasets/br_anatel_telefonia_movel/flows.py index 7811fa7dc9..d0750206f5 100644 --- a/pipelines/datasets/br_anatel_telefonia_movel/flows.py +++ b/pipelines/datasets/br_anatel_telefonia_movel/flows.py @@ -1,10 +1,11 @@ """Flows for br_anatel_telefonia_movel — Prefect 3.""" -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.anatel.telefonia_movel.flows import ( _run_anatel_telefonia_movel, ) +from pipelines.utils.flow import flow def _anatel_tm_flow(table_id: str, cron: str): @@ -35,9 +36,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] - # pyrefly: ignore [missing-attribute] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] _flow.job_variables = {"memory_limit": "8Gi", "memory_request": "2Gi"} return _flow diff --git a/pipelines/datasets/br_anp_precos_combustiveis/flows.py b/pipelines/datasets/br_anp_precos_combustiveis/flows.py index 9b82f3ddf3..35d227526b 100644 --- a/pipelines/datasets/br_anp_precos_combustiveis/flows.py +++ b/pipelines/datasets/br_anp_precos_combustiveis/flows.py @@ -2,13 +2,14 @@ Flow br_anp_precos_combustiveis__microdados — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.anp_precos_combustiveis.tasks import ( download_and_transform, get_data_source_anp_max_date, make_partitions, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, DateOnly, @@ -120,7 +121,6 @@ def br_anp_precos_combustiveis__microdados( ) -# pyrefly: ignore [missing-attribute] br_anp_precos_combustiveis__microdados.deploy_schedules = [ - {"cron": "0 10 * * *", "timezone": "America/Sao_Paulo"} + Cron("0 10 * * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_ans_beneficiario/flows.py b/pipelines/datasets/br_ans_beneficiario/flows.py index 978ceb7c09..79ef9571f4 100644 --- a/pipelines/datasets/br_ans_beneficiario/flows.py +++ b/pipelines/datasets/br_ans_beneficiario/flows.py @@ -2,7 +2,7 @@ Flow br_ans_beneficiario__informacao_consolidada — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.ans_beneficiario.tasks import ( crawler_ans, @@ -10,6 +10,7 @@ files_to_download, get_file_max_date, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, PartBdpro, @@ -133,7 +134,6 @@ def br_ans_beneficiario__informacao_consolidada( ) -# pyrefly: ignore [missing-attribute] br_ans_beneficiario__informacao_consolidada.deploy_schedules = [ - {"cron": "0 21 * * *", "timezone": "America/Sao_Paulo"} + Cron("0 21 * * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_bcb_agencia/flows.py b/pipelines/datasets/br_bcb_agencia/flows.py index 7777c39c5c..4dd66c2130 100644 --- a/pipelines/datasets/br_bcb_agencia/flows.py +++ b/pipelines/datasets/br_bcb_agencia/flows.py @@ -2,7 +2,7 @@ Flow br_bcb_agencia__agencia — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.bcb_agencia.tasks import ( clean_data, @@ -11,6 +11,7 @@ get_documents_metadata, get_latest_file, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, PartBdpro, @@ -146,7 +147,6 @@ def br_bcb_agencia__agencia( ) -# pyrefly: ignore [missing-attribute] br_bcb_agencia__agencia.deploy_schedules = [ - {"cron": "0 22 25-31 * *", "timezone": "America/Sao_Paulo"} + Cron("0 22 25-31 * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_bcb_estban/flows.py b/pipelines/datasets/br_bcb_estban/flows.py index 5069b82474..4b047307cb 100644 --- a/pipelines/datasets/br_bcb_estban/flows.py +++ b/pipelines/datasets/br_bcb_estban/flows.py @@ -2,7 +2,7 @@ Flows para br_bcb_estban — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.bcb_estban.tasks import ( cleaning_data, @@ -12,6 +12,7 @@ get_id_municipio, get_latest_file, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, PartBdpro, @@ -167,8 +168,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_bcb_sicor/flows.py b/pipelines/datasets/br_bcb_sicor/flows.py index c62ec69d99..219ff15661 100644 --- a/pipelines/datasets/br_bcb_sicor/flows.py +++ b/pipelines/datasets/br_bcb_sicor/flows.py @@ -2,10 +2,11 @@ Flows para br_bcb_sicor — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.bcb.flows import _run_bcb_sicor from pipelines.crawler.bcb.tasks import create_load_dictionary +from pipelines.utils.flow import flow from pipelines.utils.tasks import ( rename_flow_run_dataset_table, run_dbt, @@ -52,8 +53,7 @@ def _flow( local_redis_execution=local_redis_execution, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_bcb_taxa_cambio/flows.py b/pipelines/datasets/br_bcb_taxa_cambio/flows.py index f8694f21dc..8834410576 100644 --- a/pipelines/datasets/br_bcb_taxa_cambio/flows.py +++ b/pipelines/datasets/br_bcb_taxa_cambio/flows.py @@ -2,12 +2,13 @@ Flow br_bcb_taxa_cambio — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.bcb_taxa_cambio.tasks import ( get_data_taxa_cambio, treat_data_taxa_cambio, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( AllBdpro, DateFormat, @@ -94,7 +95,6 @@ def br_bcb_taxa_cambio__taxa_cambio( ) -# pyrefly: ignore [missing-attribute] br_bcb_taxa_cambio__taxa_cambio.deploy_schedules = [ - {"cron": "0 8 * * *", "timezone": "America/Sao_Paulo"} + Cron("0 8 * * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_bcb_taxa_selic/flows.py b/pipelines/datasets/br_bcb_taxa_selic/flows.py index c9921f1fe7..3eedc0a440 100644 --- a/pipelines/datasets/br_bcb_taxa_selic/flows.py +++ b/pipelines/datasets/br_bcb_taxa_selic/flows.py @@ -6,8 +6,10 @@ import pandas as pd import requests -from prefect import flow, task +from prefect import task +from prefect.schedules import Cron +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( AllBdpro, DateFormat, @@ -154,7 +156,6 @@ def br_bcb_taxa_selic__taxa_selic( ) -# pyrefly: ignore [missing-attribute] br_bcb_taxa_selic__taxa_selic.deploy_schedules = [ - {"cron": "0 8 * * *", "timezone": "America/Sao_Paulo"} + Cron("0 8 * * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_bd_indicadores/flows.py b/pipelines/datasets/br_bd_indicadores/flows.py index 0cddc97877..224d70c455 100755 --- a/pipelines/datasets/br_bd_indicadores/flows.py +++ b/pipelines/datasets/br_bd_indicadores/flows.py @@ -2,8 +2,6 @@ Flows for br_bd_indicadores — Prefect 3. """ -from prefect import flow - from pipelines.crawler.bd_indicadores.tasks import ( crawler_metricas, crawler_real_time, @@ -14,6 +12,7 @@ has_new_tweets, save_data_to_csv, ) +from pipelines.utils.flow import flow from pipelines.utils.tasks import ( download_data_to_gcs, rename_flow_run_dataset_table, @@ -309,19 +308,11 @@ def br_bd_indicadores__pessoas( # Schedules — apenas contabilidade e receitas tinham schedule no Prefect 0 -# pyrefly: ignore [missing-attribute] br_bd_indicadores__contabilidade.deploy_schedules = [] -# pyrefly: ignore [missing-attribute] br_bd_indicadores__receitas_planejadas.deploy_schedules = [] -# pyrefly: ignore [missing-attribute] br_bd_indicadores__twitter_metrics.deploy_schedules = [] -# pyrefly: ignore [missing-attribute] br_bd_indicadores__twitter_metrics_agg.deploy_schedules = [] -# pyrefly: ignore [missing-attribute] br_bd_indicadores__page_views.deploy_schedules = [] -# pyrefly: ignore [missing-attribute] br_bd_indicadores__website_user.deploy_schedules = [] -# pyrefly: ignore [missing-attribute] br_bd_indicadores__equipes.deploy_schedules = [] -# pyrefly: ignore [missing-attribute] br_bd_indicadores__pessoas.deploy_schedules = [] diff --git a/pipelines/datasets/br_bd_siga_o_dinheiro/flows.py b/pipelines/datasets/br_bd_siga_o_dinheiro/flows.py index fb65cc41ed..d257f8865d 100644 --- a/pipelines/datasets/br_bd_siga_o_dinheiro/flows.py +++ b/pipelines/datasets/br_bd_siga_o_dinheiro/flows.py @@ -2,9 +2,8 @@ Flow br_bd_siga_o_dinheiro — Prefect 3. """ -from prefect import flow - from pipelines.crawler.bd_siga_o_dinheiro.tasks import get_table_ids +from pipelines.utils.flow import flow from pipelines.utils.tasks import download_data_to_gcs, run_dbt @@ -30,5 +29,4 @@ def br_bd_siga_o_dinheiro( download_data_to_gcs(dataset_id=dataset_id, table_id=table_id) -# pyrefly: ignore [missing-attribute] br_bd_siga_o_dinheiro.deploy_schedules = [] diff --git a/pipelines/datasets/br_bndes_operacoes_contratadas/flows.py b/pipelines/datasets/br_bndes_operacoes_contratadas/flows.py index a9d0eadd7b..cc7f71f3e0 100644 --- a/pipelines/datasets/br_bndes_operacoes_contratadas/flows.py +++ b/pipelines/datasets/br_bndes_operacoes_contratadas/flows.py @@ -5,12 +5,13 @@ orquestracao (poll deferido) vive em pipelines/crawler/bndes/flows.py. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.bndes.flows import ( _run_operacoes, _run_operacoes_administracao_publica, ) +from pipelines.utils.flow import flow @flow( @@ -41,9 +42,8 @@ def br_bndes_operacoes_contratadas__operacoes_indiretas_automaticas( ) -# pyrefly: ignore [missing-attribute] br_bndes_operacoes_contratadas__operacoes_indiretas_automaticas.deploy_schedules = [ - {"cron": "0 6 * * 1", "timezone": "America/Sao_Paulo"} + Cron("0 6 * * 1", timezone="America/Sao_Paulo") ] @@ -75,9 +75,8 @@ def br_bndes_operacoes_contratadas__operacoes_nao_automaticas( ) -# pyrefly: ignore [missing-attribute] br_bndes_operacoes_contratadas__operacoes_nao_automaticas.deploy_schedules = [ - {"cron": "0 6 * * 1", "timezone": "America/Sao_Paulo"} + Cron("0 6 * * 1", timezone="America/Sao_Paulo") ] @@ -111,7 +110,6 @@ def br_bndes_operacoes_contratadas__operacoes_administracao_publica( # cron semanal (segunda 06h BRT), igual a irma; a fonte atualiza mensal e o poll # deferido no-opa quando nao ha novidade. Ajuste se quiser outra janela. -# pyrefly: ignore [missing-attribute] br_bndes_operacoes_contratadas__operacoes_administracao_publica.deploy_schedules = [ - {"cron": "0 6 * * 1", "timezone": "America/Sao_Paulo"} + Cron("0 6 * * 1", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_camara_dados_abertos/flows.py b/pipelines/datasets/br_camara_dados_abertos/flows.py index 62239e5203..7d88ca079b 100644 --- a/pipelines/datasets/br_camara_dados_abertos/flows.py +++ b/pipelines/datasets/br_camara_dados_abertos/flows.py @@ -2,11 +2,12 @@ Flows para br_camara_dados_abertos — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.camara_dados_abertos.flows import ( _run_camara_dados_abertos, ) +from pipelines.utils.flow import flow def _camara_flow(table_id: str, cron: str): @@ -33,8 +34,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_cgu_beneficios_cidadao/flows.py b/pipelines/datasets/br_cgu_beneficios_cidadao/flows.py index f194ac6db3..f5113774fc 100644 --- a/pipelines/datasets/br_cgu_beneficios_cidadao/flows.py +++ b/pipelines/datasets/br_cgu_beneficios_cidadao/flows.py @@ -6,9 +6,10 @@ em pipelines/crawler/cgu/utils.py e pipelines/utils/utils.py. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.cgu.flows import _run_cgu_beneficios_cidadao +from pipelines.utils.flow import flow def _flow_factory(table_id: str, cron: str): @@ -37,8 +38,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_cgu_cartao_pagamento/flows.py b/pipelines/datasets/br_cgu_cartao_pagamento/flows.py index 40f99b31d6..76924317cb 100644 --- a/pipelines/datasets/br_cgu_cartao_pagamento/flows.py +++ b/pipelines/datasets/br_cgu_cartao_pagamento/flows.py @@ -2,9 +2,10 @@ Flows para br_cgu_cartao_pagamento — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.cgu.flows import _run_cgu_cartao_pagamento +from pipelines.utils.flow import flow def _flow_factory(table_id: str, cron: str): @@ -33,8 +34,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_cgu_emendas_parlamentares/flows.py b/pipelines/datasets/br_cgu_emendas_parlamentares/flows.py index 62905a438c..2529a4d9e0 100644 --- a/pipelines/datasets/br_cgu_emendas_parlamentares/flows.py +++ b/pipelines/datasets/br_cgu_emendas_parlamentares/flows.py @@ -2,12 +2,13 @@ Flows para br_cgu_emendas_parlamentares — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.cgu_emendas_parlamentares.tasks import ( convert_str_to_float, get_last_modified_time, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, FreeLag, @@ -119,7 +120,6 @@ def br_cgu_emendas_parlamentares__microdados( ) -# pyrefly: ignore [missing-attribute] br_cgu_emendas_parlamentares__microdados.deploy_schedules = [ - {"cron": "30 19 * * *", "timezone": "America/Sao_Paulo"} + Cron("30 19 * * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_cgu_licitacao_contrato/flows.py b/pipelines/datasets/br_cgu_licitacao_contrato/flows.py index a3c031164f..05fad06df1 100644 --- a/pipelines/datasets/br_cgu_licitacao_contrato/flows.py +++ b/pipelines/datasets/br_cgu_licitacao_contrato/flows.py @@ -2,9 +2,10 @@ Flows para br_cgu_licitacao_contrato — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.cgu.flows import _run_cgu_licitacao_contrato +from pipelines.utils.flow import flow def _flow_factory(table_id: str, cron: str | None): @@ -34,10 +35,7 @@ def _flow( ) if cron: - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [ - {"cron": cron, "timezone": "America/Sao_Paulo"} - ] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_cgu_pessoal_executivo_federal/flows.py b/pipelines/datasets/br_cgu_pessoal_executivo_federal/flows.py index d3962bf117..a106e1dcc8 100644 --- a/pipelines/datasets/br_cgu_pessoal_executivo_federal/flows.py +++ b/pipelines/datasets/br_cgu_pessoal_executivo_federal/flows.py @@ -1,11 +1,12 @@ """Flows para br_cgu_pessoal_executivo_federal — Prefect 3.""" -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.cgu_pessoal_executivo_federal.tasks import ( clean_save_table, crawl, ) +from pipelines.utils.flow import flow from pipelines.utils.tasks import ( rename_flow_run_dataset_table, run_dbt, @@ -75,7 +76,6 @@ def br_cgu_pessoal_executivo_federal__terceirizados( ) -# pyrefly: ignore [missing-attribute] br_cgu_pessoal_executivo_federal__terceirizados.deploy_schedules = [ - {"cron": "0 0 28 2/4 *", "timezone": "America/Sao_Paulo"} + Cron("0 0 28 2/4 *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_cgu_servidores_executivo_federal/flows.py b/pipelines/datasets/br_cgu_servidores_executivo_federal/flows.py index 56160a3b21..bfc9939874 100644 --- a/pipelines/datasets/br_cgu_servidores_executivo_federal/flows.py +++ b/pipelines/datasets/br_cgu_servidores_executivo_federal/flows.py @@ -17,9 +17,10 @@ enquanto o ZIP é gerado de forma assíncrona). """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.cgu.flows import _run_cgu_servidores_publicos +from pipelines.utils.flow import flow def _flow_factory(table_id: str, cron: str): @@ -48,8 +49,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_cnj_improbidade_administrativa/flows.py b/pipelines/datasets/br_cnj_improbidade_administrativa/flows.py index 8c880a3ef7..da0aca3369 100644 --- a/pipelines/datasets/br_cnj_improbidade_administrativa/flows.py +++ b/pipelines/datasets/br_cnj_improbidade_administrativa/flows.py @@ -2,7 +2,7 @@ Flow br_cnj_improbidade_administrativa — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.cnj_improbidade_administrativa.tasks import ( get_max_date, @@ -10,6 +10,7 @@ main_task, write_csv_file, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, DateOnly, @@ -99,7 +100,6 @@ def br_cnj_improbidade_administrativa__condenacao( ) -# pyrefly: ignore [missing-attribute] br_cnj_improbidade_administrativa__condenacao.deploy_schedules = [ - {"cron": "0 7 * * 1", "timezone": "America/Sao_Paulo"} + Cron("0 7 * * 1", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_cvm_administradores_carteira/flows.py b/pipelines/datasets/br_cvm_administradores_carteira/flows.py index 941da42c00..7ea108d12d 100644 --- a/pipelines/datasets/br_cvm_administradores_carteira/flows.py +++ b/pipelines/datasets/br_cvm_administradores_carteira/flows.py @@ -2,11 +2,12 @@ Flows for br_cvm_administradores_carteira — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.cvm_administradores_carteira.flows import ( _run_cvm_administradores_carteira, ) +from pipelines.utils.flow import flow def _adm_cart_flow(table_id: str, cron: str): @@ -33,8 +34,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_cvm_fi/flows.py b/pipelines/datasets/br_cvm_fi/flows.py index a66a1821fa..e169dc7e34 100644 --- a/pipelines/datasets/br_cvm_fi/flows.py +++ b/pipelines/datasets/br_cvm_fi/flows.py @@ -2,9 +2,10 @@ Flows for br_cvm_fi — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.cvm.flows import _run_cvm_fi +from pipelines.utils.flow import flow def _cvm_fi_flow(table_id: str, cron: str, date_column_name: dict): @@ -34,8 +35,7 @@ def _flow( url=url, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_cvm_oferta_publica_distribuicao/flows.py b/pipelines/datasets/br_cvm_oferta_publica_distribuicao/flows.py index 3a0d7376ce..1961da5360 100644 --- a/pipelines/datasets/br_cvm_oferta_publica_distribuicao/flows.py +++ b/pipelines/datasets/br_cvm_oferta_publica_distribuicao/flows.py @@ -2,12 +2,13 @@ Flows for br_cvm_oferta_publica_distribuicao — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.cvm_oferta_publica_distribuicao.tasks import ( clean_table_oferta_distribuicao, crawl, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, DateOnly, @@ -95,7 +96,6 @@ def br_cvm_oferta_publica_distribuicao__dia( ) -# pyrefly: ignore [missing-attribute] br_cvm_oferta_publica_distribuicao__dia.deploy_schedules = [ - {"cron": "45 6 * * 1-5", "timezone": "America/Sao_Paulo"} + Cron("45 6 * * 1-5", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_denatran_frota/flows.py b/pipelines/datasets/br_denatran_frota/flows.py index f42ea5bd75..7ad0f2ca91 100644 --- a/pipelines/datasets/br_denatran_frota/flows.py +++ b/pipelines/datasets/br_denatran_frota/flows.py @@ -2,7 +2,7 @@ Flows para br_denatran_frota — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.denatran_frota.constants import ( constants as denatran_constants, @@ -14,6 +14,7 @@ treat_municipio_tipo_task, treat_uf_tipo_task, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, PartBdpro, @@ -197,11 +198,9 @@ def br_denatran_frota__municipio_tipo( ) -# pyrefly: ignore [missing-attribute] br_denatran_frota__uf_tipo.deploy_schedules = [ - {"cron": "0 21 10-30 * *", "timezone": "America/Sao_Paulo"} + Cron("0 21 10-30 * *", timezone="America/Sao_Paulo") ] -# pyrefly: ignore [missing-attribute] br_denatran_frota__municipio_tipo.deploy_schedules = [ - {"cron": "20 21 10-30 * *", "timezone": "America/Sao_Paulo"} + Cron("20 21 10-30 * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_fgv_igp/flows.py b/pipelines/datasets/br_fgv_igp/flows.py index 43b59ea067..876626f4f4 100644 --- a/pipelines/datasets/br_fgv_igp/flows.py +++ b/pipelines/datasets/br_fgv_igp/flows.py @@ -2,9 +2,10 @@ Flows for br_fgv_igp — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.fgv_igp.flows import _run_fgv_igp +from pipelines.utils.flow import flow def _igp_flow(table_id: str, indice: str, periodo: str, cron: str | None): @@ -33,9 +34,8 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] _flow.deploy_schedules = ( - [{"cron": cron, "timezone": "America/Sao_Paulo"}] if cron else [] + [Cron(cron, timezone="America/Sao_Paulo")] if cron else [] ) return _flow diff --git a/pipelines/datasets/br_ibge_inpc/flows.py b/pipelines/datasets/br_ibge_inpc/flows.py index b7f351d69b..a388da1f80 100644 --- a/pipelines/datasets/br_ibge_inpc/flows.py +++ b/pipelines/datasets/br_ibge_inpc/flows.py @@ -2,9 +2,10 @@ Flows para br_ibge_inpc — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.ibge_inflacao.flows import _run_ibge_inflacao +from pipelines.utils.flow import flow def _inpc_flow(table_id: str, cron: str): @@ -30,8 +31,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_ibge_ipca/flows.py b/pipelines/datasets/br_ibge_ipca/flows.py index ca4b5fff11..5dc9d050a2 100644 --- a/pipelines/datasets/br_ibge_ipca/flows.py +++ b/pipelines/datasets/br_ibge_ipca/flows.py @@ -5,9 +5,10 @@ (guards contra bloco vazio da API do IBGE) em utils.py. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.ibge_inflacao.flows import _run_ibge_inflacao +from pipelines.utils.flow import flow def _ipca_flow(table_id: str, cron: str): @@ -33,8 +34,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_ibge_ipca15/flows.py b/pipelines/datasets/br_ibge_ipca15/flows.py index efad5a4d08..ecc22f3413 100644 --- a/pipelines/datasets/br_ibge_ipca15/flows.py +++ b/pipelines/datasets/br_ibge_ipca15/flows.py @@ -2,9 +2,10 @@ Flows para br_ibge_ipca15 — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.ibge_inflacao.flows import _run_ibge_inflacao +from pipelines.utils.flow import flow def _ipca15_flow(table_id: str, cron: str): @@ -33,8 +34,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_ibge_pnadc/flows.py b/pipelines/datasets/br_ibge_pnadc/flows.py index 81dd45d004..a08435faee 100644 --- a/pipelines/datasets/br_ibge_pnadc/flows.py +++ b/pipelines/datasets/br_ibge_pnadc/flows.py @@ -2,7 +2,7 @@ Flow br_ibge_pnadc — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.ibge_pnadc.tasks import ( build_partitions, @@ -10,6 +10,7 @@ get_data_source_date_and_url, ) from pipelines.datasets.br_ibge_pnadc.tasks import build_dicionario_task +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, PartBdpro, @@ -121,9 +122,8 @@ def br_ibge_pnadc__microdados( ) -# pyrefly: ignore [missing-attribute] br_ibge_pnadc__microdados.deploy_schedules = [ - {"cron": "0 5 15-31 2,5,8,11 *", "timezone": "America/Sao_Paulo"} + Cron("0 5 15-31 2,5,8,11 *", timezone="America/Sao_Paulo") ] @@ -196,7 +196,6 @@ def br_ibge_pnadc__dicionario( ) -# pyrefly: ignore [missing-attribute] br_ibge_pnadc__dicionario.deploy_schedules = [ - {"cron": "0 5 1,15 * *", "timezone": "America/Sao_Paulo"} + Cron("0 5 1,15 * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_inmet_bdmep/flows.py b/pipelines/datasets/br_inmet_bdmep/flows.py index 2e0703b7b6..ad1fd2f754 100644 --- a/pipelines/datasets/br_inmet_bdmep/flows.py +++ b/pipelines/datasets/br_inmet_bdmep/flows.py @@ -2,12 +2,13 @@ Flows for br_inmet_bdmep — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.inmet_bdmep.tasks import ( extract_last_date_from_source, get_base_inmet, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, DateOnly, @@ -115,7 +116,6 @@ def br_inmet_bdmep__microdados( ) -# pyrefly: ignore [missing-attribute] br_inmet_bdmep__microdados.deploy_schedules = [ - {"cron": "0 22 * * 1-5", "timezone": "America/Sao_Paulo"}, + Cron("0 22 * * 1-5", timezone="America/Sao_Paulo"), ] diff --git a/pipelines/datasets/br_me_caged/flows.py b/pipelines/datasets/br_me_caged/flows.py index 1adb52087a..ff07cd9cba 100644 --- a/pipelines/datasets/br_me_caged/flows.py +++ b/pipelines/datasets/br_me_caged/flows.py @@ -2,7 +2,7 @@ Flows para br_me_caged — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.me_caged.tasks import ( build_partitions, @@ -12,6 +12,7 @@ get_source_last_date, get_table_last_date, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, PartBdpro, @@ -149,8 +150,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_me_cnpj/flows.py b/pipelines/datasets/br_me_cnpj/flows.py index 1c84aa3f29..80f0fb3e6c 100644 --- a/pipelines/datasets/br_me_cnpj/flows.py +++ b/pipelines/datasets/br_me_cnpj/flows.py @@ -2,9 +2,10 @@ Flows for br_me_cnpj — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.me_cnpj.flows import _run_me_cnpj +from pipelines.utils.flow import flow def _me_cnpj_flow(table_id: str, cron: str): @@ -31,20 +32,17 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow br_me_cnpj__empresas = _me_cnpj_flow(table_id="empresas", cron="0 6 * * *") -# pyrefly: ignore [missing-attribute] br_me_cnpj__empresas.job_variables = { "memory_limit": "5Gi", "memory_request": "2Gi", } br_me_cnpj__socios = _me_cnpj_flow(table_id="socios", cron="0 7 * * *") -# pyrefly: ignore [missing-attribute] br_me_cnpj__socios.job_variables = { "memory_limit": "5Gi", "memory_request": "2Gi", @@ -55,7 +53,6 @@ def _flow( br_me_cnpj__estabelecimentos = _me_cnpj_flow( table_id="estabelecimentos", cron="0 9 * * *" ) -# pyrefly: ignore [missing-attribute] br_me_cnpj__estabelecimentos.job_variables = { "memory_limit": "5Gi", "memory_request": "2Gi", diff --git a/pipelines/datasets/br_me_comex_stat/flows.py b/pipelines/datasets/br_me_comex_stat/flows.py index bf91b41afe..ce198933a5 100644 --- a/pipelines/datasets/br_me_comex_stat/flows.py +++ b/pipelines/datasets/br_me_comex_stat/flows.py @@ -2,7 +2,7 @@ Flows for br_me_comex_stat — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.me_comex_stat.constants import ( constants as comex_constants, @@ -12,6 +12,7 @@ download_br_me_comex_stat, parse_last_date, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, PartBdpro, @@ -130,8 +131,7 @@ def _flow( date_format="%Y-%m", ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_me_rais/flows.py b/pipelines/datasets/br_me_rais/flows.py index c87c132e6f..57d7455c45 100644 --- a/pipelines/datasets/br_me_rais/flows.py +++ b/pipelines/datasets/br_me_rais/flows.py @@ -2,9 +2,8 @@ Flows for br_me_rais — Prefect 3. """ -from prefect import flow - from pipelines.crawler.me_rais.flows import _run_rais +from pipelines.utils.flow import flow @flow( @@ -61,7 +60,5 @@ def br_me_rais__microdados_vinculos( ) -# pyrefly: ignore [missing-attribute] br_me_rais__microdados_estabelecimentos.deploy_schedules = [] -# pyrefly: ignore [missing-attribute] br_me_rais__microdados_vinculos.deploy_schedules = [] diff --git a/pipelines/datasets/br_me_siconfi/flows.py b/pipelines/datasets/br_me_siconfi/flows.py index de8b88e2f5..b17a3252e5 100644 --- a/pipelines/datasets/br_me_siconfi/flows.py +++ b/pipelines/datasets/br_me_siconfi/flows.py @@ -23,10 +23,11 @@ import tempfile from datetime import datetime -from prefect import flow +from prefect.schedules import Cron from pipelines.datasets.br_me_siconfi import tasks, utils from pipelines.datasets.br_me_siconfi.constants import constants +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import AllFree, DateFormat, YearOnly from pipelines.utils.metadata.tasks import ( commit_source_update_task, @@ -213,12 +214,10 @@ def br_me_siconfi_flow( # SICONFI is annual but revised retroactively; rebuild once a month (1st at # 16:00 BRT). Each run rebuilds fully — there is no source-poll no-op here. -# pyrefly: ignore [missing-attribute] br_me_siconfi_flow.deploy_schedules = [ - {"cron": "0 16 1 * *", "timezone": "America/Sao_Paulo"} + Cron("0 16 1 * *", timezone="America/Sao_Paulo") ] # The município window build holds a full year of data in pandas at a time. -# pyrefly: ignore [missing-attribute] br_me_siconfi_flow.job_variables = {"memory": "16Gi"} @@ -299,5 +298,4 @@ def br_me_siconfi_seed_flow( # The legacy build holds one year of município Excel in pandas at a time. -# pyrefly: ignore [missing-attribute] br_me_siconfi_seed_flow.job_variables = {"memory": "8Gi"} diff --git a/pipelines/datasets/br_mp_pep/flows.py b/pipelines/datasets/br_mp_pep/flows.py index 020e4aecdf..d46148e60d 100644 --- a/pipelines/datasets/br_mp_pep/flows.py +++ b/pipelines/datasets/br_mp_pep/flows.py @@ -4,7 +4,7 @@ import datetime -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.mp_pep.tasks import ( clean_data, @@ -14,6 +14,7 @@ scraper, setup_web_driver, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, PartBdpro, @@ -112,7 +113,6 @@ def br_mp_pep__cargos_funcoes( ) -# pyrefly: ignore [missing-attribute] br_mp_pep__cargos_funcoes.deploy_schedules = [ - {"cron": "0 14 * * 3", "timezone": "America/Sao_Paulo"} + Cron("0 14 * * 3", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_ms_cnes/flows.py b/pipelines/datasets/br_ms_cnes/flows.py index d78adc729e..90828086d9 100644 --- a/pipelines/datasets/br_ms_cnes/flows.py +++ b/pipelines/datasets/br_ms_cnes/flows.py @@ -13,9 +13,10 @@ este arquivo, o deploy sai "0 registrados, N pulados" e passa. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.datasus.flows import _run_cnes +from pipelines.utils.flow import flow def _cnes_flow(table_id: str, cron: str | None): @@ -45,10 +46,7 @@ def _flow( ) if cron: - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [ - {"cron": cron, "timezone": "America/Sao_Paulo"} - ] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_ms_sia/flows.py b/pipelines/datasets/br_ms_sia/flows.py index d33651e882..5ac0ed0b9e 100644 --- a/pipelines/datasets/br_ms_sia/flows.py +++ b/pipelines/datasets/br_ms_sia/flows.py @@ -2,9 +2,10 @@ Flows for br_ms_sia — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.datasus.flows import _run_siasus +from pipelines.utils.flow import flow def _sia_flow(table_id: str, cron: str): @@ -33,8 +34,7 @@ def _flow( year_month_to_extract=year_month_to_extract, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_ms_sih/flows.py b/pipelines/datasets/br_ms_sih/flows.py index b1a72dbf2f..d6d0332956 100644 --- a/pipelines/datasets/br_ms_sih/flows.py +++ b/pipelines/datasets/br_ms_sih/flows.py @@ -2,9 +2,10 @@ Flows for br_ms_sih — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.datasus.flows import _run_sihsus +from pipelines.utils.flow import flow def _sih_flow(table_id: str, cron: str): @@ -33,8 +34,7 @@ def _flow( year_month_to_extract=year_month_to_extract, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_ms_sinan/flows.py b/pipelines/datasets/br_ms_sinan/flows.py index 1ecade2593..6c010104a8 100644 --- a/pipelines/datasets/br_ms_sinan/flows.py +++ b/pipelines/datasets/br_ms_sinan/flows.py @@ -2,9 +2,8 @@ Flows for br_ms_sinan — Prefect 3. """ -from prefect import flow - from pipelines.crawler.datasus.flows import _run_sinan +from pipelines.utils.flow import flow @flow( diff --git a/pipelines/datasets/br_poder360_pesquisas/flows.py b/pipelines/datasets/br_poder360_pesquisas/flows.py index 73fe414b75..3b01c7ef24 100644 --- a/pipelines/datasets/br_poder360_pesquisas/flows.py +++ b/pipelines/datasets/br_poder360_pesquisas/flows.py @@ -2,9 +2,10 @@ Flow br_poder360_pesquisas — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.poder360_pesquisas.tasks import crawler +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, DateOnly, @@ -88,7 +89,6 @@ def br_poder360_pesquisas__microdados( ) -# pyrefly: ignore [missing-attribute] br_poder360_pesquisas__microdados.deploy_schedules = [ - {"cron": "42 3 * * *", "timezone": "America/Sao_Paulo"} + Cron("42 3 * * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_rf_cafir/flows.py b/pipelines/datasets/br_rf_cafir/flows.py index 350239fad4..c1813ffd6a 100644 --- a/pipelines/datasets/br_rf_cafir/flows.py +++ b/pipelines/datasets/br_rf_cafir/flows.py @@ -2,7 +2,7 @@ Flows for br_rf_cafir — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.rf_cafir.constants import ( constants as br_rf_cafir_constants, @@ -13,6 +13,7 @@ task_get_last_update_date, task_parse_api_metadata, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, DateOnly, @@ -134,7 +135,6 @@ def br_rf_cafir__imoveis_rurais( ) -# pyrefly: ignore [missing-attribute] br_rf_cafir__imoveis_rurais.deploy_schedules = [ - {"cron": "0 0 * * *", "timezone": "America/Sao_Paulo"} + Cron("0 0 * * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_rf_cno/flows.py b/pipelines/datasets/br_rf_cno/flows.py index 5d0c5306af..ba3324a170 100644 --- a/pipelines/datasets/br_rf_cno/flows.py +++ b/pipelines/datasets/br_rf_cno/flows.py @@ -9,9 +9,10 @@ `safe_cast(data as date)` do model virava NULL e o filtro incremental nunca inseria. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.rf.flows import _run_rf +from pipelines.utils.flow import flow def _cno_flow(table_id: str, cron: str): @@ -40,8 +41,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_rj_isp_estatisticas_seguranca/flows.py b/pipelines/datasets/br_rj_isp_estatisticas_seguranca/flows.py index ea6b602a94..50bf628782 100644 --- a/pipelines/datasets/br_rj_isp_estatisticas_seguranca/flows.py +++ b/pipelines/datasets/br_rj_isp_estatisticas_seguranca/flows.py @@ -2,9 +2,10 @@ Flows for br_rj_isp_estatisticas_seguranca — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.isp.flows import _run_isp +from pipelines.utils.flow import flow def _isp_flow(table_id: str, cron: str): @@ -31,8 +32,7 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] - _flow.deploy_schedules = [{"cron": cron, "timezone": "America/Sao_Paulo"}] + _flow.deploy_schedules = [Cron(cron, timezone="America/Sao_Paulo")] return _flow diff --git a/pipelines/datasets/br_senado_dados_abertos/flows.py b/pipelines/datasets/br_senado_dados_abertos/flows.py index 4ecc255cb7..e140139a88 100644 --- a/pipelines/datasets/br_senado_dados_abertos/flows.py +++ b/pipelines/datasets/br_senado_dados_abertos/flows.py @@ -16,10 +16,11 @@ import shutil import tempfile -from prefect import flow +from prefect.schedules import Cron from pipelines.datasets.br_senado_dados_abertos.constants import constants from pipelines.datasets.br_senado_dados_abertos.tasks import extract_clean +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, DateOnly, @@ -165,9 +166,7 @@ def br_senado_dados_abertos_flow( # Legislative activity updates on business days; refresh every morning (BRT). -# pyrefly: ignore [missing-attribute] br_senado_dados_abertos_flow.deploy_schedules = [ - {"cron": "0 8 * * *", "timezone": "America/Sao_Paulo"} + Cron("0 8 * * *", timezone="America/Sao_Paulo") ] -# pyrefly: ignore [missing-attribute] br_senado_dados_abertos_flow.job_variables = {"memory": "4Gi"} diff --git a/pipelines/datasets/br_sfb_sicar/flows.py b/pipelines/datasets/br_sfb_sicar/flows.py index b3dc147e4c..16ab1a0046 100644 --- a/pipelines/datasets/br_sfb_sicar/flows.py +++ b/pipelines/datasets/br_sfb_sicar/flows.py @@ -2,7 +2,7 @@ Flow br_sfb_sicar — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.sfb_sicar.constants import Constants from pipelines.crawler.sfb_sicar.tasks import ( @@ -10,6 +10,7 @@ get_each_uf_release_date, unzip_to_parquet, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, DateOnly, @@ -113,7 +114,6 @@ def br_sfb_sicar__area_imovel( ) -# pyrefly: ignore [missing-attribute] br_sfb_sicar__area_imovel.deploy_schedules = [ - {"cron": "15 21 15 * *", "timezone": "America/Sao_Paulo"} + Cron("15 21 15 * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_stf_corte_aberta/flows.py b/pipelines/datasets/br_stf_corte_aberta/flows.py index cc10eb4e45..162c8df559 100644 --- a/pipelines/datasets/br_stf_corte_aberta/flows.py +++ b/pipelines/datasets/br_stf_corte_aberta/flows.py @@ -2,13 +2,14 @@ Flow br_stf_corte_aberta — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.stf_corte_aberta.tasks import ( download_and_transform, get_data_source_stf_max_date, make_partitions, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( DateFormat, DateOnly, @@ -119,7 +120,6 @@ def br_stf_corte_aberta__decisoes( ) -# pyrefly: ignore [missing-attribute] br_stf_corte_aberta__decisoes.deploy_schedules = [ - {"cron": "0 12 * * *", "timezone": "America/Sao_Paulo"} + Cron("0 12 * * *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/br_tse_eleicoes/flows.py b/pipelines/datasets/br_tse_eleicoes/flows.py index 1512902939..340e982ad3 100644 --- a/pipelines/datasets/br_tse_eleicoes/flows.py +++ b/pipelines/datasets/br_tse_eleicoes/flows.py @@ -2,9 +2,10 @@ Flows for br_tse_eleicoes — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron from pipelines.crawler.tse_eleicoes.flows import _run_tse_eleicoes +from pipelines.utils.flow import flow def _tse_flow(table_id: str, cron: str | None): @@ -31,9 +32,8 @@ def _flow( force_run=force_run, ) - # pyrefly: ignore [missing-attribute] _flow.deploy_schedules = ( - [{"cron": cron, "timezone": "America/Sao_Paulo"}] if cron else [] + [Cron(cron, timezone="America/Sao_Paulo")] if cron else [] ) return _flow diff --git a/pipelines/datasets/fundacao_lemann/flows.py b/pipelines/datasets/fundacao_lemann/flows.py index ffa8044fcf..f2b64ab202 100644 --- a/pipelines/datasets/fundacao_lemann/flows.py +++ b/pipelines/datasets/fundacao_lemann/flows.py @@ -2,8 +2,9 @@ Flow fundacao_lemann — Prefect 3. """ -from prefect import flow +from prefect.schedules import Cron +from pipelines.utils.flow import flow from pipelines.utils.tasks import ( download_data_to_gcs, rename_flow_run_dataset_table, @@ -40,7 +41,6 @@ def fundacao_lemann__ano_escola_serie_educacao_aprendizagem_adequada( download_data_to_gcs(dataset_id=dataset_id, table_id=table_id) -# pyrefly: ignore [missing-attribute] fundacao_lemann__ano_escola_serie_educacao_aprendizagem_adequada.deploy_schedules = [ - {"cron": "0 9 1 1 *", "timezone": "America/Sao_Paulo"} + Cron("0 9 1 1 *", timezone="America/Sao_Paulo") ] diff --git a/pipelines/datasets/test_dataset/flows.py b/pipelines/datasets/test_dataset/flows.py index 3d0740fefb..c2ea2101f2 100644 --- a/pipelines/datasets/test_dataset/flows.py +++ b/pipelines/datasets/test_dataset/flows.py @@ -2,8 +2,7 @@ Flows de teste para a função download_data_to_gcs. """ -from prefect import flow - +from pipelines.utils.flow import flow from pipelines.utils.tasks import run_dbt DATASET_ID = "test_dataset" diff --git a/pipelines/datasets/us_bls_cpi/flows.py b/pipelines/datasets/us_bls_cpi/flows.py index 04c57b0616..437b5392e2 100644 --- a/pipelines/datasets/us_bls_cpi/flows.py +++ b/pipelines/datasets/us_bls_cpi/flows.py @@ -13,10 +13,11 @@ import shutil import tempfile -from prefect import flow +from prefect.schedules import Cron from pipelines.datasets.us_bls_cpi.constants import constants from pipelines.datasets.us_bls_cpi.tasks import clean_cpi, download_cpi +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( AllFree, DateFormat, @@ -177,10 +178,8 @@ def us_bls_cpi_flow( # BLS releases CPI monthly, ~2nd week, on a US business day. Poll across a few # mid-month days at 16:00 BRT; the source-poll guard no-ops until a new month # actually appears. -# pyrefly: ignore [missing-attribute] us_bls_cpi_flow.deploy_schedules = [ - {"cron": "0 16 10,11,12,13,14,15 * *", "timezone": "America/Sao_Paulo"} + Cron("0 16 10,11,12,13,14,15 * *", timezone="America/Sao_Paulo") ] # The clean step holds ~4M rows in pandas; give the worker headroom. -# pyrefly: ignore [missing-attribute] us_bls_cpi_flow.job_variables = {"memory": "8Gi"} diff --git a/pipelines/datasets/us_bls_qcew/flows.py b/pipelines/datasets/us_bls_qcew/flows.py index 3c1dc69f2d..b13224d306 100644 --- a/pipelines/datasets/us_bls_qcew/flows.py +++ b/pipelines/datasets/us_bls_qcew/flows.py @@ -19,13 +19,14 @@ import shutil import tempfile -from prefect import flow +from prefect.schedules import Cron from pipelines.datasets.us_bls_qcew.constants import constants from pipelines.datasets.us_bls_qcew.tasks import ( clean_qcew, latest_source_period, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( AllFree, DateFormat, @@ -188,11 +189,9 @@ def us_bls_qcew_flow( # and Wages news releases land in early March, June, September, and December. # Poll across the first ~10 days of those months at 16:00 BRT; the source-poll # guard no-ops until a new quarter actually appears in the singlefiles. -# pyrefly: ignore [missing-attribute] us_bls_qcew_flow.deploy_schedules = [ - {"cron": "0 16 1-10 3,6,9,12 *", "timezone": "America/Sao_Paulo"} + Cron("0 16 1-10 3,6,9,12 *", timezone="America/Sao_Paulo") ] # The clean step streams ~15M-row singlefiles one chunk at a time (peak ~1.75GB # in pandas); give the worker headroom above that. -# pyrefly: ignore [missing-attribute] us_bls_qcew_flow.job_variables = {"memory": "8Gi"} diff --git a/pipelines/datasets/world_cricsheet/flows.py b/pipelines/datasets/world_cricsheet/flows.py index 56da898510..dffb41c83b 100644 --- a/pipelines/datasets/world_cricsheet/flows.py +++ b/pipelines/datasets/world_cricsheet/flows.py @@ -16,13 +16,14 @@ import shutil import tempfile -from prefect import flow +from prefect.schedules import Cron from pipelines.datasets.world_cricsheet.constants import constants from pipelines.datasets.world_cricsheet.tasks import ( clean_cricsheet, download_cricsheet, ) +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import ( AllFree, DateFormat, @@ -212,11 +213,9 @@ def world_cricsheet_flow( # WEEKLY instead — Monday 06:00 BRT — which cuts the cost ~4x while keeping # freshness fine. The source-poll guard still no-ops between real releases, and # the full-replace dump means overlapping windows never duplicate. -# pyrefly: ignore [missing-attribute] world_cricsheet_flow.deploy_schedules = [ - {"cron": "0 6 * * 1", "timezone": "America/Sao_Paulo"} + Cron("0 6 * * 1", timezone="America/Sao_Paulo") ] # The deliveries build streams 11.4M rows and the bundle extracts to several GB; # give the worker headroom. -# pyrefly: ignore [missing-attribute] world_cricsheet_flow.job_variables = {"memory": "12Gi"} diff --git a/pipelines/utils/execute_dbt_model/flows.py b/pipelines/utils/execute_dbt_model/flows.py index 70c4c57931..b8388c9d08 100644 --- a/pipelines/utils/execute_dbt_model/flows.py +++ b/pipelines/utils/execute_dbt_model/flows.py @@ -9,7 +9,9 @@ from typing import Any from dbt.cli.main import dbtRunner -from prefect import flow, task +from prefect import task + +from pipelines.utils.flow import flow @task @@ -115,5 +117,4 @@ def run_dbt_model_flow( dbt_done.result() -# pyrefly: ignore [missing-attribute] run_dbt_model_flow.deploy_schedules = [] diff --git a/pipelines/utils/flow.py b/pipelines/utils/flow.py new file mode 100644 index 0000000000..bc5a6309ae --- /dev/null +++ b/pipelines/utils/flow.py @@ -0,0 +1,120 @@ +""" +`Flow` da Base dos Dados — o `Flow` do Prefect 3 mais os atributos de deploy. + +`.github/scripts/deploy_flows.py` lê dois atributos do objeto flow que o +`prefect.Flow` não declara: + +- `deploy_schedules`: lista de agendamentos (`prefect.schedules.Cron`), usada no + deploy de produção (no pool de dev os schedules são descartados); +- `job_variables`: overrides da configuração de infraestrutura do work pool + (memória, CPU...). + +Atribuí-los a um `prefect.Flow` funciona em runtime, mas o Pyrefly acusa +`missing-attribute`, já que a classe do Prefect não os declara. Este módulo +declara os dois numa subclasse de `prefect.Flow` e expõe um decorator `flow` +que instancia essa subclasse. Use sempre este `flow` em `flows.py`, nunca o +`prefect.flow`: + +```python +from prefect.schedules import Cron + +from pipelines.utils.flow import flow + + +@flow(name="meu_dataset", log_prints=True) +def meu_dataset_flow() -> None: ... + + +meu_dataset_flow.deploy_schedules = [ + Cron("0 16 10 * *", timezone="America/Sao_Paulo") +] +meu_dataset_flow.job_variables = {"memory": "8Gi"} +``` + +Como `Flow` herda de `prefect.Flow`, as checagens `isinstance(obj, Flow)` do +script de deploy e do próprio Prefect continuam valendo. +""" + +from collections.abc import Callable +from typing import Any, ParamSpec, TypeVar, overload + +from prefect import Flow as PrefectFlow +from prefect.futures import PrefectFuture +from prefect.schedules import Schedule +from prefect.task_runners import TaskRunner + +P = ParamSpec("P") +R = TypeVar("R") + + +class Flow(PrefectFlow[P, R]): + """Flow do Prefect 3 com os atributos que o deploy da BD lê. + + Attributes: + deploy_schedules: Agendamentos do flow, construídos com + `prefect.schedules.Cron` (que já recebe `timezone`). Vazio (o + padrão) significa deployment sem schedule — execução apenas manual. + job_variables: Overrides da configuração de infraestrutura do work + pool, por exemplo `{"memory": "8Gi"}`. Vazio usa o padrão do pool. + """ + + deploy_schedules: list[Schedule] | None + job_variables: dict[str, Any] | None + + def __init__(self, *args: Any, **kwargs: Any) -> None: + super().__init__(*args, **kwargs) + self.deploy_schedules = None + self.job_variables = None + + +# Uso sem parênteses (`@flow`). A assinatura da função decorada não é +# preservada aqui — usar `Callable[P, R]` tornaria esta sobrecarga genérica, o +# que exigiria a sintaxe de type parameters do Python 3.12 (ruff UP047). +@overload +def flow(fn: Callable[..., Any], /) -> Flow[..., Any]: ... + + +@overload +def flow( + fn: None = None, + /, + *, + name: str | None = None, + version: str | None = None, + flow_run_name: Callable[[], str] | str | None = None, + retries: int | None = None, + retry_delay_seconds: int | float | None = None, + task_runner: TaskRunner[PrefectFuture[Any]] | None = None, + description: str | None = None, + timeout_seconds: int | float | None = None, + validate_parameters: bool = True, + persist_result: bool | None = None, + cache_result_in_memory: bool = True, + log_prints: bool | None = None, + **kwargs: Any, +) -> Callable[[Callable[P, R]], Flow[P, R]]: ... + + +def flow( + fn: Callable[..., Any] | None = None, /, **kwargs: Any +) -> Flow[..., Any] | Callable[[Callable[..., Any]], Flow[..., Any]]: + """Decorator equivalente ao `prefect.flow`, mas devolvendo um `Flow` da BD. + + Aceita os mesmos argumentos nomeados do `prefect.flow` (`name`, + `log_prints`, `retries`, `flow_run_name`, ...), repassados ao construtor do + `prefect.Flow`. + + Args: + fn: A função decorada, quando o decorator é usado sem parênteses. + **kwargs: Argumentos do `prefect.flow`. + + Returns: + O `Flow` construído, ou o decorator que o constrói. + """ + if fn is not None: + return Flow(fn=fn, **kwargs) + + def decorator(fn: Callable[..., Any]) -> Flow[..., Any]: + return Flow(fn=fn, **kwargs) + + return decorator diff --git a/pipelines/utils/materialize_prod/flows.py b/pipelines/utils/materialize_prod/flows.py index fdc7a17d1b..f53eb8519d 100644 --- a/pipelines/utils/materialize_prod/flows.py +++ b/pipelines/utils/materialize_prod/flows.py @@ -5,8 +5,7 @@ from __future__ import annotations -from prefect import flow - +from pipelines.utils.flow import flow from pipelines.utils.materialize_prod.tasks import ( download_files_from_bucket_folders, ) @@ -170,5 +169,4 @@ def transfer_files_to_prod_flow( ) -# pyrefly: ignore [missing-attribute] transfer_files_to_prod_flow.deploy_schedules = [] # utilitário, disparo manual diff --git a/pipelines/utils/metadata/flows.py b/pipelines/utils/metadata/flows.py index 76508f2ccd..2d6a9ebf9c 100644 --- a/pipelines/utils/metadata/flows.py +++ b/pipelines/utils/metadata/flows.py @@ -9,8 +9,7 @@ Pydantic no parsing do parâmetro do flow. """ -from prefect import flow - +from pipelines.utils.flow import flow from pipelines.utils.metadata.domain import CoverageSpec from pipelines.utils.metadata.tasks import ( register_table_materialization_task, @@ -37,5 +36,4 @@ def update_temporal_coverage( ) -# pyrefly: ignore [missing-attribute] update_temporal_coverage.deploy_schedules = [] diff --git a/pipelines/utils/tests/test_flow.py b/pipelines/utils/tests/test_flow.py new file mode 100644 index 0000000000..9c109ba503 --- /dev/null +++ b/pipelines/utils/tests/test_flow.py @@ -0,0 +1,59 @@ +"""Testes do decorator `flow` da BD (`pipelines.utils.flow`).""" + +from prefect import Flow as PrefectFlow +from prefect.schedules import Cron + +from pipelines.utils.flow import Flow, flow + + +@flow(name="flow_de_teste", log_prints=True) +def _flow_com_opcoes(x: int = 1) -> int: + return x + + +@flow +def _flow_sem_parenteses() -> str: + return "ok" + + +def test_e_um_flow_do_prefect(): + """O deploy e o próprio Prefect fazem `isinstance(obj, prefect.Flow)`.""" + assert isinstance(_flow_com_opcoes, Flow) + assert isinstance(_flow_com_opcoes, PrefectFlow) + assert isinstance(_flow_sem_parenteses, PrefectFlow) + + +def test_repassa_as_opcoes_do_prefect(): + assert _flow_com_opcoes.name == "flow_de_teste" + assert _flow_com_opcoes.log_prints is True + # Sem `name=`, o Prefect infere o nome a partir da função. + assert _flow_sem_parenteses.name == "-flow-sem-parenteses" + + +def test_atributos_de_deploy_comecam_vazios(): + assert _flow_com_opcoes.deploy_schedules == [] + assert _flow_com_opcoes.job_variables == {} + + +def test_atributos_de_deploy_nao_sao_compartilhados(): + """Cada flow tem os seus — nada de default mutável de classe.""" + assert ( + _flow_com_opcoes.deploy_schedules + is not _flow_sem_parenteses.deploy_schedules + ) + assert ( + _flow_com_opcoes.job_variables + is not _flow_sem_parenteses.job_variables + ) + + +def test_aceita_atribuicao_dos_atributos_de_deploy(): + schedules = [Cron("0 16 10 * *", timezone="America/Sao_Paulo")] + + _flow_com_opcoes.deploy_schedules = schedules + _flow_com_opcoes.job_variables = {"memory": "8Gi"} + + assert _flow_com_opcoes.deploy_schedules == schedules + assert _flow_com_opcoes.deploy_schedules[0].cron == "0 16 10 * *" + assert _flow_com_opcoes.deploy_schedules[0].timezone == "America/Sao_Paulo" + assert _flow_com_opcoes.job_variables == {"memory": "8Gi"} diff --git a/pipelines/{{cookiecutter.pipeline_name}}/flows.py b/pipelines/{{cookiecutter.pipeline_name}}/flows.py index da266dc06b..e04346b4a0 100644 --- a/pipelines/{{cookiecutter.pipeline_name}}/flows.py +++ b/pipelines/{{cookiecutter.pipeline_name}}/flows.py @@ -6,16 +6,18 @@ # # Aqui é onde devem ser definidos os flows da pipeline (Prefect 3). # -# Cada flow é uma função decorada com `@flow`. Os passos são chamadas de -# tasks (funções decoradas com `@task`, definidas em `tasks.py`) executadas -# na ordem do corpo da função. +# Cada flow é uma função decorada com `@flow` — o decorator da BD, em +# `pipelines.utils.flow`, e não o do Prefect: ele devolve um `Flow` que declara +# os atributos de deploy (`deploy_schedules`, `job_variables`). Os passos são +# chamadas de tasks (funções decoradas com `@task`, definidas em `tasks.py`) +# executadas na ordem do corpo da função. # # O deploy é feito por `.github/scripts/deploy_flows.py`, que descobre os # objetos `Flow` deste arquivo automaticamente — não é preciso registrar # storage nem run_config manualmente. # # O agendamento é definido inline no próprio objeto do flow, via -# `.deploy_schedules`, uma lista de dicts `{"cron": ..., "timezone": ...}`. +# `.deploy_schedules`, uma lista de `Cron` (`prefect.schedules`). # Em `dev` os schedules são ignorados; em `prod` são ativados na sincronização # com o backend. Veja https://crontab.guru/ para montar a expressão cron. # @@ -25,9 +27,10 @@ # ############################################################################### -from prefect import flow +from prefect.schedules import Cron from pipelines.datasets.{{cookiecutter.pipeline_name}}.tasks import say_hello +from pipelines.utils.flow import flow @flow(name="{{cookiecutter.pipeline_name}}", log_prints=True) @@ -37,5 +40,5 @@ def {{cookiecutter.pipeline_name}}_flow(name: str = "World") -> None: # Agendamento (opcional). Remova se o flow não deve rodar em intervalos fixos. {{cookiecutter.pipeline_name}}_flow.deploy_schedules = [ - {"cron": "0 14 * * 1", "timezone": "America/Sao_Paulo"} # segundas, 14:00 + Cron("0 14 * * 1", timezone="America/Sao_Paulo") # segundas, 14:00 ]