Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
36 commits
Select commit Hold shift + click to select a range
9f4ff9b
feat(test-event-pipeline): scaffold pilot flows for event-driven auto…
Winzen Aug 30, 2026
b7e1f0e
fix(test-event-pipeline): pass download_params as JSON, compute refer…
Winzen Aug 31, 2026
dbcf4c1
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Aug 31, 2026
0187ab4
feat(test-event-pipeline): add mat_test_flow and Automação 2 (flow_do…
Winzen Aug 31, 2026
5e779da
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Aug 31, 2026
164649f
fix(test-event-pipeline): revert UTC import to timezone.utc for py3.1…
Winzen Aug 31, 2026
852b029
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Aug 31, 2026
b93696e
fix(test-event-pipeline): fix pyrefly type errors from PR #1932 CI
Winzen Aug 31, 2026
dd8b572
feat(automations): match_related on the etapa tag, not just the event…
Winzen Aug 31, 2026
b9a7235
feat(test-event-pipeline): real check_update against backend coverage…
Winzen Sep 1, 2026
734e0b3
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Sep 1, 2026
3d90322
feat(materialize-prod): wire generic mat_test_flow to real dev->prod …
Winzen Sep 1, 2026
a14d8e4
fix(materialize-prod): make dbt_command a parameter, default "run"
Winzen Sep 1, 2026
b7f8e11
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Sep 1, 2026
d4d2af5
fix(materialize-prod): default download_billing_project to basedosdados
Winzen Sep 1, 2026
0ae30da
test(test-event-pipeline): exercise the real dev->prod path
Winzen Sep 1, 2026
180be33
refactor(mat_test): promote dataset_id/table_id to real flow parameters
Winzen Sep 1, 2026
cb06d8a
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Sep 1, 2026
6f8d490
fix(mat_test): coerce coverage dict to CoverageSpec and run rename sy…
Winzen Sep 1, 2026
bdd8273
fix(materialize_prod): run rename_flow_run_dataset_table synchronously
Winzen Sep 1, 2026
baeaddd
docs(mat_test): update mat_test_flow docstring to match current code
Winzen Sep 1, 2026
5c6412a
refactor(1867): replace Automation/emit_event chain with run_deployme…
Winzen Sep 1, 2026
5119fa9
refactor(1867): clean up leftover automation-era vocabulary
Winzen Sep 1, 2026
d5f4124
refactor(1867): rename backend_* params/constants to dataset_id/table_id
Winzen Sep 1, 2026
8472abb
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Sep 1, 2026
2e9bc38
refactor(1867): encapsula check_update/flow_download em CheckThenDown…
Winzen Sep 3, 2026
14cfaa3
feat(1867): adiciona piloto particionado (ano=/mes=) pra testar dev->…
Winzen Sep 3, 2026
8ece717
fix(1867): não repetir ano/mes dentro do CSV particionado
Winzen Sep 3, 2026
b9f0ac5
refactor(1867): consolida pilotos event_pipeline em test_dataset/
Winzen Sep 3, 2026
e54a55b
refactor(1867): move upload_to_gcs pra capsula e simplifica nomenclatura
Winzen Sep 5, 2026
90cfd08
feat(1867): migra 11 datasets reais pro pipeline orientado a eventos
Winzen Sep 9, 2026
a1895b2
refactor(1867): pipeline_factory() e make_pipeline sempre em tasks.py
Winzen Sep 13, 2026
2533e97
revert(1867): desfaz migração de br_denatran_frota
Winzen Sep 14, 2026
9fdc4bc
refactor(1867): remove sufixo _flow do nome dos flows migrados
Winzen Sep 17, 2026
36929f9
revert(1867): desfaz migração de br_me_cnpj e br_sfb_sicar
Winzen Sep 17, 2026
a958e99
Merge remote-tracking branch 'origin/main' into feat/event-pipeline-a…
Winzen Sep 18, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions .github/scripts/deploy_flows.py
Original file line number Diff line number Diff line change
Expand Up @@ -193,6 +193,9 @@ def deploy_flow(

job_variables = getattr(flow, "job_variables", None)

extra_tags = getattr(flow, "deploy_tags", None) or []
tags = ["automated-deploy", *extra_tags]

try:
flow.from_source(
source=GitRepository(
Expand All @@ -203,7 +206,7 @@ def deploy_flow(
).deploy(
name=deployment_name,
work_pool_name=pool_name,
tags=["automated-deploy"],
tags=tags,
schedules=schedules,
job_variables=job_variables,
build=False,
Expand All @@ -214,7 +217,7 @@ def deploy_flow(
if not schedules
else f"com schedules: {schedules}"
)
return True, f" ✓ {deployment_name} registrado {status}"
return True, f" ✓ {deployment_name} registrado {status}, tags={tags}"
except Exception as e:
return False, f" ✗ Falha ao registrar {deployment_name}: {e}"

Expand Down
13 changes: 13 additions & 0 deletions models/test_dataset/test_dataset__test_event_pipeline.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
{{
config(
schema="test_dataset",
alias="test_event_pipeline",
materialized="table",
)
}}

-- Piloto da issue #1867: lê o CSV que flow_download_flow sobe pro staging
-- (upload_to_gcs) — cada rodada substitui a única linha pela data do
-- download mais recente.
select safe_cast(reference_date as date) as reference_date
from {{ set_datalake_project("test_dataset_staging.test_event_pipeline") }} as t
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
{{
config(
schema="test_dataset",
alias="test_event_pipeline_partitioned",
materialized="table",
)
}}

-- Variante particionada (ano=/mes=) do piloto da issue #1867 — testa
-- DownloadResult.partition_folders/transfer_files_to_prod_flow(folders=...)
-- promovendo só a fatia nova, não o staging inteiro. As colunas ano/mes
-- vêm explícitas no CSV (não extraídas do caminho Hive), igual ao dado
-- real de cada arquivo.
select
safe_cast(ano as int64) as ano,
safe_cast(mes as int64) as mes,
safe_cast(reference_date as date) as reference_date
from
{{ set_datalake_project("test_dataset_staging.test_event_pipeline_partitioned") }}
as t
8 changes: 8 additions & 0 deletions pipelines/datasets/br_ans_beneficiario/constants.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
"""
Constant values for br_ans_beneficiario.
"""

DATASET_ID = "br_ans_beneficiario"

INFORMACAO_CONSOLIDADA_TABLE_ID = "informacao_consolidada"
INFORMACAO_CONSOLIDADA_URL = "https://dadosabertos.ans.gov.br/FTP/PDA/informacoes_consolidadas_de_beneficiarios-024/"
178 changes: 51 additions & 127 deletions pipelines/datasets/br_ans_beneficiario/flows.py
Original file line number Diff line number Diff line change
@@ -1,150 +1,74 @@
"""
Flow br_ans_beneficiario__informacao_consolidada — Prefect 3.
Flows para br_ans_beneficiario — Prefect 3.

Migrado por completo pro pipeline orientado a eventos (issue #1867):
check_update -> download -> mat_test. Lógica específica do dataset mora em
`tasks.py`, constantes em `constants.py` — aqui só a fiação
(`CheckThenDownloadPipeline` + `@flow`).
"""

from prefect import flow

from pipelines.crawler.ans_beneficiario.tasks import (
crawler_ans,
extract_links_and_dates,
files_to_download,
get_file_max_date,
from pipelines.datasets.br_ans_beneficiario.constants import (
DATASET_ID,
INFORMACAO_CONSOLIDADA_TABLE_ID,
)
from pipelines.utils.metadata.domain import (
DateFormat,
PartBdpro,
YearMonth,
from pipelines.datasets.br_ans_beneficiario.tasks import (
br_ans_beneficiario_check_for_update,
br_ans_beneficiario_download,
)
from pipelines.utils.metadata.tasks import (
commit_source_update_task,
poll_source_for_update_task,
register_table_materialization_task,
from pipelines.utils.stage_dispatch import (
CheckThenDownloadPipeline,
Etapa,
deploy_tags,
)
from pipelines.utils.tasks import (
rename_flow_run_dataset_table,
run_dbt,
upload_to_gcs,

_informacao_consolidada_pipeline = CheckThenDownloadPipeline(
dataset_id=DATASET_ID,
table_id=INFORMACAO_CONSOLIDADA_TABLE_ID,
check_for_update=br_ans_beneficiario_check_for_update,
download_data=br_ans_beneficiario_download,
# Mesma granularidade do flow antigo: dado é mensal, sem dia.
date_format="%Y-%m",
)


@flow(
name="br_ans_beneficiario__informacao_consolidada",
name=_informacao_consolidada_pipeline.check_update_flow_name,
log_prints=True,
)
def br_ans_beneficiario__informacao_consolidada(
dataset_id: str = "br_ans_beneficiario",
table_id: str = "informacao_consolidada",
url: str = "https://dadosabertos.ans.gov.br/FTP/PDA/informacoes_consolidadas_de_beneficiarios-024/",
year: str | None = None,
materialize_after_dump: bool = True,
update_metadata: bool = True,
target: str = "prod",
force_run: bool = False,
) -> None:
# pyrefly: ignore [unused-coroutine]
rename_flow_run_dataset_table(
prefix="Dump: ", dataset_id=dataset_id, table_id=table_id
)

links_and_dates = extract_links_and_dates(url=url)
file_last_date = get_file_max_date(df=links_and_dates)

if force_run:
files = files_to_download(df=links_and_dates, year=year)
else:
has_new_data = poll_source_for_update_task(
dataset_id=dataset_id,
table_id=table_id,
source_max_date=file_last_date,
env="prod",
date_format="%Y-%m",
compare_against="coverage",
)
if not has_new_data:
print(f"Não há atualizações para a tabela {table_id}!")
return
files = files_to_download(df=links_and_dates, year=None)

# Comita o Update da fonte já aqui, antes de baixar/materializar: se o
# flow falhar no meio, o metadado da fonte ainda reflete que havia dado
# novo publicado, mesmo que a tabela não tenha sido atualizada.
commit_source_update_task(
dataset_id=dataset_id,
table_id=table_id,
source_max_date=file_last_date,
env="prod",
date_format="%Y-%m",
update_metadata=update_metadata,
materialize_after_dump=materialize_after_dump,
)

if not files:
print("Nenhum arquivo para baixar.")
return

output_filepath = crawler_ans(files=files)

# crawler_ans -> parquet_partition grava .parquet. Sem declarar o
# formato, o dump_header chamado por upload_to_gcs procura .csv (default)
# e não encontra nada.
upload_to_gcs(
data_path=output_filepath,
dataset_id=dataset_id,
table_id=table_id,
bucket_name="basedosdados-dev",
dump_mode="append",
source_format="parquet",
)
def br_ans_beneficiario_informacao_consolidada_check_update() -> None:
_informacao_consolidada_pipeline.run_check_update()

run_dbt(
dataset_id=dataset_id,
table_id=table_id,
dbt_command="run/test",
target="dev",
)

if not materialize_after_dump:
return

upload_to_gcs(
data_path=output_filepath,
dataset_id=dataset_id,
table_id=table_id,
bucket_name="basedosdados",
dump_mode="append",
source_format="parquet",
)
# pyrefly: ignore [missing-attribute]
br_ans_beneficiario_informacao_consolidada_check_update.deploy_tags = (
deploy_tags(DATASET_ID, Etapa.CHECK_UPDATE)
)

run_dbt(
dataset_id=dataset_id,
table_id=table_id,
dbt_command="run/test",
target=target,
)

if update_metadata:
# R4-normalização: o legado usava date_format="%Y-%m-%d" com colunas
# ano/mes (emitia um endDay/startDay=1 espúrio). O domínio refatorado
# força YearMonth↔YEAR_MONTH, descartando o dia — granularidade correta
# para dado mensal. Valores ano/mês idênticos ao legado; free_lag default
# (months=6) reproduz time_delta={"months":6}.
register_table_materialization_task(
dataset_id=dataset_id,
table_id=table_id,
coverage=PartBdpro(
date_column=YearMonth(year="ano", month="mes"),
date_format=DateFormat.YEAR_MONTH,
),
env="prod",
bq_project="basedosdados",
)
@flow(
name=_informacao_consolidada_pipeline.download_flow_name,
log_prints=True,
)
def br_ans_beneficiario_informacao_consolidada_download(
download_params: dict,
) -> None:
_informacao_consolidada_pipeline.run_download(download_params)


# pyrefly: ignore [missing-attribute]
br_ans_beneficiario__informacao_consolidada.deploy_schedules = [
{"cron": "0 21 * * *", "timezone": "America/Sao_Paulo"}
]
br_ans_beneficiario_informacao_consolidada_download.deploy_tags = deploy_tags(
DATASET_ID, Etapa.DOWNLOAD
)
# Pico medido em produção após otimizar parquet_partition (category dtype +
# del/gc.collect() por estado): ~1.78Gi. ~1.7x de margem sobre esse valor.
# del/gc.collect() por estado): ~1.78Gi. ~1.7x de margem sobre esse valor —
# mesmo tier do flow antigo, já que o download pesado (crawler_ans) continua
# acontecendo aqui.
# pyrefly: ignore [missing-attribute]
br_ans_beneficiario__informacao_consolidada.job_variables = {"memory": "3Gi"}
br_ans_beneficiario_informacao_consolidada_download.job_variables = {
"memory": "3Gi"
}
_informacao_consolidada_pipeline.download_deployment = (
br_ans_beneficiario_informacao_consolidada_download.fn.__name__
)
60 changes: 60 additions & 0 deletions pipelines/datasets/br_ans_beneficiario/tasks.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
"""
Tasks for br_ans_beneficiario.
"""

from datetime import date

from pipelines.crawler.ans_beneficiario.tasks import (
crawler_ans,
extract_links_and_dates,
files_to_download,
get_file_max_date,
)
from pipelines.datasets.br_ans_beneficiario.constants import (
INFORMACAO_CONSOLIDADA_URL,
)
from pipelines.utils.metadata.domain import DateFormat, PartBdpro, YearMonth
from pipelines.utils.stage_dispatch import CheckResult, DownloadResult

# ──────────────────────────────────────────────────────────────────────────────
# informacao_consolidada (issue #1867) — ver constants.py
#
# Diferente de br_ibge_ipca: aqui o check É leve de verdade e independente —
# `extract_links_and_dates` só faz um GET na página de listagem (HTML da
# pasta FTP-like da ANS) e lê datas de "última atualização" já presentes no
# próprio HTML, sem baixar nenhum arquivo de dado. `download_data` refaz essa
# mesma listagem (idem barato) só pra descobrir de novo quais arquivos
# baixar — não dá pra repassar o DataFrame entre pods, mas a chamada em si
# é a mesma leve de sempre, não o download pesado (que só acontece depois,
# em `crawler_ans`).
# ──────────────────────────────────────────────────────────────────────────────


def br_ans_beneficiario_check_for_update() -> CheckResult:
links_and_dates = extract_links_and_dates(url=INFORMACAO_CONSOLIDADA_URL)
file_last_date = get_file_max_date(df=links_and_dates) # "YYYY-MM"
reference_date = date.fromisoformat(f"{file_last_date}-01")
return CheckResult(reference_date=reference_date)


def br_ans_beneficiario_download(download_params: dict) -> DownloadResult:
links_and_dates = extract_links_and_dates(url=INFORMACAO_CONSOLIDADA_URL)
files = files_to_download(df=links_and_dates, year=None)

output_filepath = crawler_ans(files=files)

ref = date.fromisoformat(download_params["reference_date"])
return DownloadResult(
coverage=PartBdpro(
date_column=YearMonth(year="ano", month="mes"),
date_format=DateFormat.YEAR_MONTH,
).model_dump(),
data_path=output_filepath,
bq_project="basedosdados",
# crawler_ans -> parquet_partition grava .parquet, não .csv (default).
source_format="parquet",
# to_partitions particiona por ano/mes/sigla_uf/modalidade_operadora,
# mas só um ano/mes muda por execução — promover o nível ano/mes já
# carrega todas as sub-partições de uf/modalidade daquele mês.
partition_folders=[f"ano={ref.year}/mes={ref.month:02d}"],
)
14 changes: 14 additions & 0 deletions pipelines/datasets/br_ibge_ipca/constants.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
"""
Constant values for br_ibge_ipca.
"""

DATASET_ID = "br_ibge_ipca"

# As 4 tabelas do dataset — migradas pro pipeline orientado a eventos
# (issue #1867), substituindo o antigo `_ipca_flow`/`_run_ibge_inflacao`
# monolítico (que segue existindo em `crawler/ibge_inflacao/flows.py`,
# ainda usado por br_ibge_ipca15/br_ibge_inpc).
MES_BRASIL_TABLE_ID = "mes_brasil"
MES_CATEGORIA_BRASIL_TABLE_ID = "mes_categoria_brasil"
MES_CATEGORIA_RM_TABLE_ID = "mes_categoria_rm"
MES_CATEGORIA_MUNICIPIO_TABLE_ID = "mes_categoria_municipio"
Loading
Loading