diff --git a/pipelines/crawler/cgu/flows.py b/pipelines/crawler/cgu/flows.py index 1b844a13a4..58dfd928c6 100644 --- a/pipelines/crawler/cgu/flows.py +++ b/pipelines/crawler/cgu/flows.py @@ -26,9 +26,28 @@ upload_to_gcs, ) +# Formato em que cada tabela do br_cgu_beneficios_cidadao é gravada em disco por +# `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 = { + "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, @@ -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, ) run_dbt( dataset_id=dataset_id, @@ -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, @@ -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 @@ -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 @@ -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 @@ -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 @@ -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], ) diff --git a/pipelines/datasets/br_cgu_beneficios_cidadao/flows.py b/pipelines/datasets/br_cgu_beneficios_cidadao/flows.py index f194ac6db3..db4a53cbc0 100644 --- a/pipelines/datasets/br_cgu_beneficios_cidadao/flows.py +++ b/pipelines/datasets/br_cgu_beneficios_cidadao/flows.py @@ -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