Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
122 changes: 119 additions & 3 deletions pipelines/crawler/cgu/flows.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,28 @@
upload_to_gcs,
)

# Formato em que cada tabela do br_cgu_beneficios_cidadao é gravada em disco por
Comment thread
DaviMacielCavalcante marked this conversation as resolved.
# `read_and_partition_beneficios_cidadao`. Precisa casar com o que o
# `upload_to_gcs` procura: quando a staging já existe, ele chama `dump_header`,
# que estoura FileNotFoundError se não achar arquivo no formato declarado.
_SOURCE_FORMAT_BENEFICIOS_CIDADAO = {
Comment thread
DaviMacielCavalcante marked this conversation as resolved.
"novo_bolsa_familia": "parquet",
"garantia_safra": "parquet",
"bpc": "csv",
}


def _part_bdpro_year_month(year: str, month: str) -> PartBdpro:
"""Cobertura padrão dos flows CGU: part_bdpro mensal (ano/mes), free_lag 6m."""
"""Monta a cobertura padrão dos flows CGU: part_bdpro mensal.

Args:
year: nome da coluna de ano na tabela (ex. `ano`, `ano_extrato`).
month: nome da coluna de mês na tabela (ex. `mes`, `mes_extrato`).

Returns:
Cobertura part_bdpro em granularidade ano/mês, com o free_lag
padrão de 6 meses definido em `PartBdpro`.
"""
return PartBdpro(
date_column=YearMonth(year=year, month=month),
date_format=DateFormat.YEAR_MONTH,
Expand All @@ -46,14 +65,41 @@ def _materialize_and_metadata(
update_metadata: bool,
coverage: CoverageSpec,
source_max_date=None,
source_format: str = "csv",
) -> None:
"""Sobe os dados particionados, materializa no dbt e atualiza metadados.

Sempre roda a metade de dev (upload no bucket `basedosdados-dev` +
`dbt run/test` no target dev). A metade de prod só roda com
`materialize_after_dump`, e a escrita de metadados só dentro dela.

Args:
filepath: caminho do diretório particionado gerado pelo passo de
particionamento.
dataset_id: id do dataset no BigQuery (ex. `br_cgu_servidores_publicos`).
table_id: id da tabela dentro do dataset.
dbt_alias: repassado ao `run_dbt`; indica se o modelo usa alias.
target: target do dbt na etapa de prod. A etapa de dev é sempre
`dev`, independente deste valor.
materialize_after_dump: se falso, para depois do dbt em dev e não
toca em prod nem em metadados.
update_metadata: se verdadeiro, registra a materialização da tabela
e grava o Update da fonte original — sempre em `env="prod"`.
source_max_date: data máxima de cobertura da fonte, no formato
`%Y-%m`. Quando `None`, o Update da fonte não é gravado.
coverage: cobertura da tabela, usada no registro da materialização.
source_format: formato em que o passo de particionamento gravou os
arquivos em disco (`csv` ou `parquet`). Vale para os dois
uploads; se não casar com o que está no disco, o `dump_header`
chamado pelo `upload_to_gcs` estoura `FileNotFoundError`.
"""
upload_to_gcs(
data_path=filepath,
dataset_id=dataset_id,
table_id=table_id,
bucket_name="basedosdados-dev",
dump_mode="append",
source_format="csv",
source_format=source_format,
Comment thread
coderabbitai[bot] marked this conversation as resolved.
)
run_dbt(
dataset_id=dataset_id,
Expand All @@ -72,7 +118,7 @@ def _materialize_and_metadata(
table_id=table_id,
bucket_name="basedosdados",
dump_mode="append",
source_format="csv",
source_format=source_format,
)
run_dbt(
dataset_id=dataset_id,
Expand Down Expand Up @@ -111,6 +157,23 @@ def _run_cgu_cartao_pagamento(
target: str,
force_run: bool,
) -> None:
"""Roda uma tabela do br_cgu_cartao_pagamento de ponta a ponta.

Baixa o mês seguinte ao último registrado nos metadados, para se a fonte
não tiver nada novo, e então particiona, sobe e materializa.

Args:
dataset_id: id do dataset no BigQuery.
table_id: id da tabela dentro do dataset.
relative_month: quantos meses somar à última data registrada nos
metadados para chegar no período a baixar.
materialize_after_dump: se verdadeiro, materializa também em prod.
dbt_alias: repassado ao `run_dbt`.
update_metadata: se verdadeiro, atualiza cobertura da tabela e o
Update da fonte original.
target: target do dbt na etapa de prod.
force_run: ignora o poll e roda mesmo sem dado novo na fonte.
"""
# pyrefly: ignore [unused-coroutine]
rename_flow_run_dataset_table(
prefix="Dump: ", dataset_id=dataset_id, table_id=table_id
Expand Down Expand Up @@ -155,6 +218,25 @@ def _run_cgu_servidores_publicos(
target: str,
force_run: bool,
) -> None:
"""Roda uma tabela do br_cgu_servidores_publicos de ponta a ponta.

Diferente dos outros flows CGU, checa antes se todas as URLs do mês
existem — a fonte publica os arquivos de servidores em lotes, e baixar
com parte deles no ar geraria mês incompleto.

Args:
dataset_id: id do dataset no BigQuery.
table_id: id da tabela dentro do dataset.
relative_month: quantos meses somar à última data registrada nos
metadados para chegar no período a baixar.
materialize_after_dump: se verdadeiro, materializa também em prod.
dbt_alias: repassado ao `run_dbt`.
update_metadata: se verdadeiro, atualiza cobertura da tabela e o
Update da fonte original.
target: target do dbt na etapa de prod.
force_run: ignora a checagem de URLs e o poll, rodando mesmo sem
dado novo na fonte.
"""
# pyrefly: ignore [unused-coroutine]
rename_flow_run_dataset_table(
prefix="Dump: ", dataset_id=dataset_id, table_id=table_id
Expand Down Expand Up @@ -206,6 +288,20 @@ def _run_cgu_licitacao_contrato(
target: str,
force_run: bool,
) -> None:
"""Roda uma tabela do br_cgu_licitacao_contrato de ponta a ponta.

Args:
dataset_id: id do dataset no BigQuery.
table_id: id da tabela dentro do dataset.
relative_month: quantos meses somar à última data registrada nos
metadados para chegar no período a baixar.
materialize_after_dump: se verdadeiro, materializa também em prod.
dbt_alias: repassado ao `run_dbt`.
update_metadata: se verdadeiro, atualiza cobertura da tabela e o
Update da fonte original.
target: target do dbt na etapa de prod.
force_run: ignora o poll e roda mesmo sem dado novo na fonte.
"""
# pyrefly: ignore [unused-coroutine]
rename_flow_run_dataset_table(
prefix="Dump: ", dataset_id=dataset_id, table_id=table_id
Expand Down Expand Up @@ -250,6 +346,25 @@ def _run_cgu_beneficios_cidadao(
target: str,
force_run: bool,
) -> None:
"""Roda uma tabela do br_cgu_beneficios_cidadao de ponta a ponta.

Único flow CGU que declara `source_format` por tabela: o
`novo_bolsa_familia` e a `garantia_safra` são gravados em parquet, o
`bpc` em csv (ver `_SOURCE_FORMAT_BENEFICIOS_CIDADAO`).

Args:
dataset_id: id do dataset no BigQuery.
table_id: id da tabela dentro do dataset. Precisa estar em
`_SOURCE_FORMAT_BENEFICIOS_CIDADAO`.
relative_month: quantos meses somar à última data registrada nos
metadados para chegar no período a baixar.
materialize_after_dump: se verdadeiro, materializa também em prod.
dbt_alias: repassado ao `run_dbt`.
update_metadata: se verdadeiro, atualiza cobertura da tabela e o
Update da fonte original.
target: target do dbt na etapa de prod.
force_run: ignora o poll e roda mesmo sem dado novo na fonte.
"""
# pyrefly: ignore [unused-coroutine]
rename_flow_run_dataset_table(
prefix="Dump: ", dataset_id=dataset_id, table_id=table_id
Expand Down Expand Up @@ -281,4 +396,5 @@ def _run_cgu_beneficios_cidadao(
update_metadata=update_metadata,
source_max_date=data_source_max_date,
coverage=_part_bdpro_year_month(**dict_for_table(table_id)),
source_format=_SOURCE_FORMAT_BENEFICIOS_CIDADAO[table_id],
)
7 changes: 4 additions & 3 deletions pipelines/datasets/br_cgu_beneficios_cidadao/flows.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
"""
Flows para br_cgu_beneficios_cidadao — Prefect 3.

Redeploy para propagar o fix de detecção/download no portal da CGU
(User-Agent de browser + tratamento do 202 assíncrono + URL sem barra final)
em pipelines/crawler/cgu/utils.py e pipelines/utils/utils.py.
Redeploy para propagar o source_format por tabela, em
pipelines/crawler/cgu/flows.py: o novo_bolsa_familia e a garantia_safra são
gravados em parquet, mas o upload declarava csv para as três tabelas. Desde a
#1677 isso derruba a run em dump_header, antes de qualquer escrita.
"""

from prefect import flow
Expand Down
Loading