diff --git a/models/br_rf_cnpj/br_rf_cnpj__empresas.sql b/models/br_rf_cnpj/br_rf_cnpj__empresas.sql index 210f7d68a1..0367871614 100644 --- a/models/br_rf_cnpj/br_rf_cnpj__empresas.sql +++ b/models/br_rf_cnpj/br_rf_cnpj__empresas.sql @@ -30,8 +30,8 @@ with where porte != "porte" {% if is_incremental() %} - and safe.parse_date('%Y-%m', data_referencia) - > (select max(data_referencia) from {{ this }}) + and data_referencia + > format_date('%Y-%m', (select max(data_referencia) from {{ this }})) {% else %} -- Dados históricos até 2023-04-30 foram migrados do modelo -- br_me_cnpj.estabelecimentos diff --git a/models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql b/models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql index d8f000778c..fa61323b18 100644 --- a/models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql +++ b/models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql @@ -55,8 +55,8 @@ with from {{ set_datalake_project("br_rf_cnpj_staging.estabelecimentos") }} {% if is_incremental() %} where - safe.parse_date('%Y-%m', data_referencia) - > (select max(data_referencia) from {{ this }}) + data_referencia + > format_date('%Y-%m', (select max(data_referencia) from {{ this }})) -- Dados históricos até 2023-04-30 foram migrados do modelo -- br_me_cnpj.estabelecimentos {% else %} diff --git a/models/br_rf_cnpj/br_rf_cnpj__socios.sql b/models/br_rf_cnpj/br_rf_cnpj__socios.sql index 6432fe5717..c72d415b51 100644 --- a/models/br_rf_cnpj/br_rf_cnpj__socios.sql +++ b/models/br_rf_cnpj/br_rf_cnpj__socios.sql @@ -33,8 +33,8 @@ with where safe_cast(qualificacao as string) != "qualificacao" {% if is_incremental() %} - and safe.parse_date('%Y-%m', data_referencia) - > (select max(data_referencia) from {{ this }}) + an data_referencia + > format_date('%Y-%m', (select max(data_referencia) from {{ this }})) {% else %} -- Dados históricos até 2023-04-30 foram migrados do modelo -- br_me_cnpj.socios diff --git a/pipelines/crawler/rf_cafir/tasks.py b/pipelines/crawler/rf_cafir/tasks.py deleted file mode 100644 index fb5be2e0be..0000000000 --- a/pipelines/crawler/rf_cafir/tasks.py +++ /dev/null @@ -1,190 +0,0 @@ -""" -Tasks for br_ms_cnes -""" - -import datetime -import os - -import pandas as pd -from prefect import task - -from pipelines.constants import constants -from pipelines.crawler.rf_cafir.constants import ( - constants as br_rf_cafir_constants, -) -from pipelines.crawler.rf_cafir.utils import ( - decide_files_to_download, - download_csv_files, - get_last_update_date, - parse_api_metadata, - preserve_zeros, - remove_ascii_zero_from_df, - strip_string, -) -from pipelines.utils.utils import log - - -@task( - retries=2, - retry_delay_seconds=constants.TASK_RETRY_DELAY.value, -) -def task_parse_api_metadata(url: str) -> pd.DataFrame: - return parse_api_metadata(url=url) - - -@task( - retries=2, - retry_delay_seconds=constants.TASK_RETRY_DELAY.value, -) -def task_get_last_update_date(url: str) -> str: - """Data de modificação mais recente entre os arquivos publicados na origem. - - Usada para nomear a partição `data=` no Storage. Ver issue #1696 — esta - não é a data de referência do dado, e a origem dela deve mudar. - - Returns: - str: Data no formato YYYY-MM-DD - """ - return get_last_update_date(url=url) - - -@task( - retries=2, - retry_delay_seconds=constants.TASK_RETRY_DELAY.value, -) -def task_decide_files_to_download( - df: pd.DataFrame, - last_update_date: str, - data_especifica: datetime.date | None = None, - data_maxima: bool = True, -) -> tuple[list[str], str | datetime.date]: - """Decide quais arquivos baixar a partir dos metadados da fonte. - - Args: - df: Metadados dos arquivos publicados (nome e data de atualização). - last_update_date: Data de modificação mais recente na origem, usada para - nomear a partição no Storage. - data_especifica: Se informada, filtra os arquivos por essa data. - data_maxima: Se True, seleciona os arquivos da data mais recente. - - Returns: - tuple[list[str], str | datetime.date]: Lista de nomes de arquivos e a - data correspondente (string quando data_maxima; date quando - data_especifica). - """ - # pyrefly: ignore [bad-return] - return decide_files_to_download( - df=df, - last_update_date=last_update_date, - data_especifica=data_especifica, - data_maxima=data_maxima, - ) - - -@task( - retries=3, - retry_delay_seconds=constants.TASK_RETRY_DELAY.value, -) -def task_download_files( - url: str, - file_list: list[str], - data_atualizacao: list[datetime.date], - last_update_date: str, -) -> str: - """Essa task faz o download dos arquivos do FTP, faz o parse dos dados e salva os arquivos em um diretório temporário. - - Returns: - str: Caminho do diretório temporário - """ - - date = data_atualizacao - - log(f"------ Extraindo dados para data: {date}") - log( - f"------ A data máxima extraida da API da Receita Federal que será utilizada para gerar partições no Storage: {last_update_date}" - ) - - files_list = file_list - log( - f"------ Os seguintes arquivos foram selecionados para download: {files_list}" - ) - - for file in files_list: - log(f"Baixando arquivo: {file} de {url}") - - # monta url - complete_url = url + file - - # baixa arquivo - download_csv_files( - file_name=file, - url=complete_url, - download_directory=br_rf_cafir_constants.PATH.value[0], - ) - - # constroi caminho do arquivo - file_path = br_rf_cafir_constants.PATH.value[0] + "/" + file - log(f"Lendo arquivo: {file} de : {file_path}") - - # Le o arquivo txt - df = pd.read_fwf( - file_path, - widths=br_rf_cafir_constants.WIDTHS.value, - names=br_rf_cafir_constants.COLUMN_NAMES.value, - dtype=br_rf_cafir_constants.DTYPE.value, - converters={ - col: preserve_zeros - for col in br_rf_cafir_constants.COLUMN_NAMES.value - }, - encoding="ISO-8859-1", - ) - - # Remove ascii /x00 (zero) - crasha tabela na materialização no BQ - df = remove_ascii_zero_from_df(df) - - # tira os espacos em branco - # pyrefly: ignore [not-callable] - df = df.applymap(strip_string) - - log(f"Salvando arquivo: {file}") - - # constroi diretório - os.makedirs( - br_rf_cafir_constants.PATH.value[1] - + f"/imoveis_rurais/data={last_update_date}/", - exist_ok=True, - ) - - # NOTE: Com modificação do formato de divulgação do FTP os arquivos passaram a ser divulgados csvs particionados por UF - # A partir de 2025, a nomenclaruta dos no Storage arquivos mudou para: "imoveis_rurais_uf_numero.csv" no lugar de "imoveris_rurais_numero.csv" - - save_path = ( - br_rf_cafir_constants.PATH.value[1] - + f"/imoveis_rurais/data={last_update_date}/" - + "imoveis_rurais_" - # extrai uf e numeração do nome do arquivo - + file.split(".")[-2] - + ".csv" - ) - - df.to_csv( - save_path, - index=False, - sep=",", - na_rep="", - encoding="utf-8", - escapechar="\\", - ) - - log(f"Arquivo salvo: {save_path.split('/')[-1]}") - - del df - - log( - f"----- Removendo o arquivo: {os.listdir(br_rf_cafir_constants.PATH.value[0])} do diretório de input" - ) - - # remove o arquivo de input - os.remove(os.path.join(br_rf_cafir_constants.PATH.value[0], file)) - - return br_rf_cafir_constants.PATH.value[1] + "/imoveis_rurais" diff --git a/pipelines/crawler/rf_cafir/utils.py b/pipelines/crawler/rf_cafir/utils.py deleted file mode 100644 index 0b6870ddc5..0000000000 --- a/pipelines/crawler/rf_cafir/utils.py +++ /dev/null @@ -1,213 +0,0 @@ -""" -General purpose functions for the br_ms_cnes project -""" - -import datetime -import os - -import pandas as pd -import requests -from bs4 import BeautifulSoup - -from pipelines.utils.utils import log - - -def strip_string(x: pd.DataFrame) -> pd.DataFrame: - """Aplica o strip em uma coluna de um dataframe, caso seja string. - ps: usar com applymap - - Args: - x (pd.Dataframe): Dataframe - - Returns: - pd.Dataframe: Dataframe com valores de linha sem espaços no início e no final das strings - """ - if isinstance(x, str): - # pyrefly: ignore [bad-return] - return x.strip() - return x - - -def remove_ascii_zero_from_df(df: pd.DataFrame) -> pd.DataFrame: - """Remove ASCII 0 (NULL) de colunas tipadas como string de um DataFrame. - Returns: - pd.DataFrame: DataFrame sem ascii 0 (\x00). - """ - # pyrefly: ignore [not-callable] - return df.applymap( - lambda x: x.replace("\x00", "") if isinstance(x, str) else x - ) - - -def requests_url(url: str) -> requests.Response: - xml_body = """ - - - - """ - - headers = { - "Depth": "1", - "Content-Type": "application/xml", - "Accept": "application/xml", - "User-Agent": "Mozilla/5.0", - } - try: - response = requests.request( - method="PROPFIND", - url=url, - headers=headers, - data=xml_body, - timeout=30, - ) - - response.raise_for_status() - - except requests.exceptions.RequestException as e: - log(f"Erro durante a requisição: {e}") - raise - - return response - - -def parse_api_metadata(url: str | None = None) -> pd.DataFrame: - """ - Faz uma requisição para a URL fornecida e extrai metadados de arquivos CSV. - Args: - url (str): A URL da API para fazer a requisição. - headers (dict, opcional): Cabeçalhos HTTP para incluir na requisição. Padrão é None. - Returns: - pd.DataFrame: Um DataFrame contendo os nomes dos arquivos e suas respectivas datas de atualização. - Raises: - ValueError: Se a quantidade de arquivos extraídos for diferente da quantidade de datas de atualização. - """ - - # pyrefly: ignore [bad-argument-type] - soup = BeautifulSoup(requests_url(url).text, "lxml") - - csvs_com_data = [] - - for x in soup.find_all("d:href"): - href = x.text - if not href.endswith(".csv"): - continue - - data_raw = href.split(".")[-3] - data_raw = data_raw.replace("D", "202") - - data = datetime.datetime.strptime(data_raw, "%Y%m%d").strftime( - "%Y-%m-%d" - ) - - csvs_com_data.append( - {"nome_arquivo": href.split("/")[-1], "data_atualizacao": data} - ) - - return pd.DataFrame(csvs_com_data) - - -def get_last_update_date(url: str) -> str: - soup = BeautifulSoup(requests_url(url).text, "lxml") - - return str( - max( - datetime.datetime.strptime( - p.find("d:getlastmodified").text, "%a, %d %b %Y %H:%M:%S GMT" - ) - for p in soup.find_all("d:prop") - if p.find("d:getlastmodified") - ).date() - ) - - -def decide_files_to_download( - last_update_date: str, - df: pd.DataFrame, - data_especifica: datetime.date | None = None, - data_maxima: bool = True, -) -> tuple[list[str], list[datetime.datetime]]: - """ - Decide quais arquivos baixar a depender da necessidade de atualização - - Parâmetros: - df (pd.DataFrame): DataFrame contendo informações dos arquivos, incluindo a data de atualização e o nome do arquivo. - data_especifica (datetime.date, opcional): Data específica para filtrar os arquivos. O Padrão é "%yyyy-%mm-%dd". - data_maxima (bool): Se True, retorna os arquivos com a data de atualização mais recente. Padrão é True. - - Retorna: - tuple: Uma tupla contendo uma lista de nomes de arquivos que atendem aos critérios fornecidos e a data correspondente. - - Levanta: - ValueError: Se não houver arquivos disponíveis para a data específica fornecida. - """ - - if data_maxima: - max_date = df["data_atualizacao"].max() - log( - f"A data máxima extraida da API da Receita Federal que será utilizada para comparar com os metadados da BD: {max_date}" - ) - - log( - f"A data máxima extraida da API da Receita Federal que será utilizada para gerar partições no Storage: {last_update_date}" - ) - - return df[df["data_atualizacao"] == max_date][ - "nome_arquivo" - ].tolist(), max_date - - elif data_especifica: - filtered_df = df[df["data_atualizacao"] == data_especifica] - if filtered_df.empty: - raise ValueError( - f"Não há arquivos disponíveis para a data {data_especifica}. Verifique o FTP da Receita Federal." - ) - # pyrefly: ignore [bad-return] - return filtered_df["nome_arquivo"].tolist(), data_especifica - - else: - raise ValueError( - "Critérios inválidos: deve-se selecionar pelo menos um dos parâmetros: 'data_maxima' ou 'data_especifica'." - ) - - -def download_csv_files( - url: str, file_name: str, download_directory: str -) -> None: - """ - Faz o download de um arquivo CSV a partir de uma URL e salva em um diretório especificado. - - Args: - url (str): A URL do arquivo CSV a ser baixado. - file_name (str): O nome do arquivo a ser salvo. - download_directory (str): O diretório onde o arquivo será salvo. - headers (dict): Cabeçalhos HTTP a serem enviados com a requisição. - - Returns: - None - """ - # cria diretório - os.makedirs(download_directory, exist_ok=True) - - log(f"Downloading--------- {url}") - # Extrai links de download - - # Setta path - file_path = os.path.join(download_directory, file_name) - - # faz request - response = requests.get(url) - - if response.status_code == 200: - # Salva no diretório especificado - with open(file_path, "wb") as f: - f.write(response.content) - log(f"Downloaded {file_name}") - else: - log( - f"Failed to download {file_name}. Status code: {response.status_code}" - ) - - -def preserve_zeros(x): - """Preserva os zeros a esquerda de um número""" - return x.strip() diff --git a/pipelines/crawler/rf_cnpj/flows.py b/pipelines/crawler/rf_cnpj/flows.py index cd82c5693c..0f1ce9eb73 100644 --- a/pipelines/crawler/rf_cnpj/flows.py +++ b/pipelines/crawler/rf_cnpj/flows.py @@ -118,9 +118,9 @@ def _run_rf_cnpj( commit_source_update_task( dataset_id=dataset_id, table_id=table_id, - source_max_date=folder_date, + source_max_date=last_modified_date, env="prod", - date_format="%Y-%m", + date_format="%Y-%m-%d", update_metadata=update_metadata, materialize_after_dump=materialize_after_dump, ) diff --git a/pipelines/datasets/br_rf_cafir/README.md b/pipelines/datasets/br_rf_cafir/README.md index 6a6e71ce37..0ca153e783 100644 --- a/pipelines/datasets/br_rf_cafir/README.md +++ b/pipelines/datasets/br_rf_cafir/README.md @@ -19,15 +19,25 @@ e é atualizada com frequência **diária**. --- -## Particionamento no Storage — atenção +## Particionamento no Storage -A partição do Storage (`data=YYYY-MM-DD/`) e a coluna `data_referencia` que -dela deriva **não** guardam a data de referência do dado. Guardam a **data de -modificação do arquivo no servidor** — o `max()` de `getlastmodified` sobre -*todos* os arquivos da pasta (`crawler/rf_cafir/utils.py::get_last_update_date`). +A partição legado do Storage (`data=YYYY-MM-DD/`) e a coluna `data_referencia` que +dela deriva **não** guardavam a data de referência do dado, mas som a **data de +modificação do arquivo no servidor**. -O poll e o commit de metadados, por outro lado, usam a data do **nome do -arquivo** (`K34313UF.D60701...` → `2026-07-01`). As duas datas divergem. +Isso foi modificado e usa-se como `data_referencia` a data do arquivo, presente em seu nome como `YMMDD` (`K34313UF.D60701...` → `2026-07-01`). + +A coluna `data_referencia` passa a armazenar de fato a data de referência do arquivo e foi incluída a `data_modificacao`, que armazena a data da última modificação do arquivo no servidor. + +## Download e processamento em paralelo + +- `download_files`/`process_files` (que faziam loop interno sequencial) viraram três tasks: +* `extract_file_records`: converte o DataFrame filtrado em list[dict] de escalares. +* `download_file`: baixa um único arquivo. +* `process_file`: processa um único arquivo (chama process_csv_file). +flows.py + +O flow agora dispara download_file.submit(...) para cada registro (download paralelo real via Prefect), e process_file.submit(..., wait_for=[download_futures[...]]) — cada processamento aguarda apenas o seu próprio download, não o lote inteiro, então arquivos diferentes correm em paralelo em todo o pipeline (download₁ pode estar em processamento enquanto download₂ ainda baixa). **Risco:** `data_referencia` é a chave do modelo incremental. Se a Receita mexer em qualquer arquivo da pasta — inclusive de anos antigos — a data de diff --git a/pipelines/crawler/rf_cafir/__init__.py b/pipelines/datasets/br_rf_cafir/__init__.py similarity index 100% rename from pipelines/crawler/rf_cafir/__init__.py rename to pipelines/datasets/br_rf_cafir/__init__.py diff --git a/pipelines/crawler/rf_cafir/constants.py b/pipelines/datasets/br_rf_cafir/constants.py similarity index 94% rename from pipelines/crawler/rf_cafir/constants.py rename to pipelines/datasets/br_rf_cafir/constants.py index 4fd748382d..2842b17f66 100644 --- a/pipelines/crawler/rf_cafir/constants.py +++ b/pipelines/datasets/br_rf_cafir/constants.py @@ -21,11 +21,6 @@ class constants(Enum): "Priority": "u=0, i", } - PATH = [ - "/tmp/input/br_rf_cafir", - "/tmp/output/br_rf_cafir", - ] - TABLE = ["imoveis_rurais"] COLUMN_NAMES = [ diff --git a/pipelines/datasets/br_rf_cafir/flows.py b/pipelines/datasets/br_rf_cafir/flows.py index 8bc1054b9e..f75232823c 100644 --- a/pipelines/datasets/br_rf_cafir/flows.py +++ b/pipelines/datasets/br_rf_cafir/flows.py @@ -2,16 +2,21 @@ Flows for br_rf_cafir — Prefect 3. """ +import datetime + from prefect import flow -from pipelines.crawler.rf_cafir.constants import ( +from pipelines.datasets.br_rf_cafir.constants import ( constants as br_rf_cafir_constants, ) -from pipelines.crawler.rf_cafir.tasks import ( - task_decide_files_to_download, - task_download_files, - task_get_last_update_date, - task_parse_api_metadata, +from pipelines.datasets.br_rf_cafir.tasks import ( + build_paths, + decide_files_to_download, + download_file, + extract_file_records, + get_api_metadata, + get_last_reference_date, + process_file, ) from pipelines.utils.metadata.domain import ( DateFormat, @@ -28,6 +33,7 @@ run_dbt, upload_to_gcs, ) +from pipelines.utils.utils import log @flow( @@ -42,27 +48,28 @@ def br_rf_cafir__imoveis_rurais( update_metadata: bool = True, target: str = "prod", force_run: bool = False, + data_referencia: str | None = None, ) -> None: # pyrefly: ignore [unused-coroutine] rename_flow_run_dataset_table( prefix="Dump: ", dataset_id=dataset_id, table_id=table_id ) - df_metadata = task_parse_api_metadata(url=br_rf_cafir_constants.URL.value) - - last_update = task_get_last_update_date( - url=br_rf_cafir_constants.URL.value - ) + input_folder, output_folder = build_paths() + df_metadata = get_api_metadata(url=br_rf_cafir_constants.URL.value) - arquivos, data_atualizacao = task_decide_files_to_download( - df=df_metadata, last_update_date=last_update - ) + if data_referencia is None: + reference_date = get_last_reference_date(df_metadata) + else: + reference_date = datetime.datetime.strptime( + data_referencia, "%Y-%m-%d" + ).date() if not force_run: has_new_data = poll_source_for_update_task( dataset_id=dataset_id, table_id=table_id, - source_max_date=data_atualizacao, + source_max_date=reference_date, env="prod", date_format="%Y-%m-%d", compare_against="coverage", @@ -76,23 +83,57 @@ def br_rf_cafir__imoveis_rurais( commit_source_update_task( dataset_id=dataset_id, table_id=table_id, - source_max_date=data_atualizacao, + source_max_date=reference_date, env="prod", date_format="%Y-%m-%d", update_metadata=update_metadata, materialize_after_dump=materialize_after_dump, ) - # pyrefly: ignore [no-matching-overload] - file_path = task_download_files( - url=br_rf_cafir_constants.URL.value, - file_list=arquivos, - data_atualizacao=data_atualizacao, - last_update_date=last_update, + filtered_df = decide_files_to_download( + df_metadata=df_metadata, reference_date=reference_date + ) + + # Extrai os registros (file_name, reference_date: cada task de + # download/processamento paralela recebe apenas o seu próprio registro + file_records = extract_file_records(df_metadata=filtered_df) + + log( + f"------ Os seguintes arquivos foram selecionados para download: " + f"{[r['nome_arquivo'] for r in file_records]}" ) + # Download por arquivo em paralelo. + download_futures = { + record["nome_arquivo"]: download_file.submit( + file_name=record["nome_arquivo"], + url=br_rf_cafir_constants.URL.value, + input_folder=input_folder, + ) + for record in file_records + } + + # Processamento roda em paralelo, mas só começa depois do seu + # próprio download (wait_for). + process_futures = [ + # pyrefly: ignore [no-matching-overload] + process_file.submit( + file_name=record["nome_arquivo"], + reference_date=record["data_referencia"], + input_folder=input_folder, + output_folder=output_folder, + wait_for=[download_futures[record["nome_arquivo"]]], + ) + for record in file_records + ] + + for future in process_futures: + future.result() + + output_path = output_folder / "imoveis_rurais" + upload_to_gcs( - data_path=file_path, + data_path=output_path, dataset_id=dataset_id, table_id=table_id, bucket_name="basedosdados-dev", @@ -111,7 +152,7 @@ def br_rf_cafir__imoveis_rurais( return upload_to_gcs( - data_path=file_path, + data_path=output_path, dataset_id=dataset_id, table_id=table_id, bucket_name="basedosdados", diff --git a/pipelines/datasets/br_rf_cafir/tasks.py b/pipelines/datasets/br_rf_cafir/tasks.py new file mode 100644 index 0000000000..aca3db9464 --- /dev/null +++ b/pipelines/datasets/br_rf_cafir/tasks.py @@ -0,0 +1,153 @@ +""" +Tasks for br_ms_cnes +""" + +import datetime +from pathlib import Path + +import pandas as pd +from prefect import task + +from pipelines.constants import constants +from pipelines.datasets.br_rf_cafir.utils import ( + download_csv_file, + parse_api_metadata, + process_csv_file, + requests_url, +) +from pipelines.utils.utils import log + + +@task +def build_paths() -> tuple[Path, Path]: + tmp_folder = Path("tmp") + input_folder = tmp_folder / "input" / "br_rf_cafir" + output_folder = tmp_folder / "output" / "br_rf_cafir" + + input_folder.mkdir(parents=True, exist_ok=True) + output_folder.mkdir(parents=True, exist_ok=True) + + return (input_folder, output_folder) + + +@task( + retries=2, + retry_delay_seconds=constants.TASK_RETRY_DELAY.value, +) +def get_api_metadata(url: str | None = None) -> pd.DataFrame: + """ + Faz uma requisição para a URL fornecida e extrai metadados de arquivos CSV. + Args: + url (str): A URL da API para fazer a requisição. + Returns: + pd.DataFrame: Um DataFrame contendo os nomes dos arquivos e suas respectivas datas de atualização. + Raises: + ValueError: Se a quantidade de arquivos extraídos for diferente da quantidade de datas de atualização. + """ + # pyrefly: ignore [bad-argument-type] + response = requests_url(url) + df_metadata = parse_api_metadata(response_text=response.text) + + return df_metadata + + +@task +def get_last_reference_date(df_metadata: pd.DataFrame) -> datetime.date: + max_date = df_metadata["data_referencia"].max() + return max_date.date() + + +@task( + retries=2, + retry_delay_seconds=constants.TASK_RETRY_DELAY.value, +) +def decide_files_to_download( + df_metadata: pd.DataFrame, + reference_date: datetime.date | None = None, +) -> pd.DataFrame: + """ + Decide quais arquivos baixar a depender da necessidade de atualização + + Parâmetros: + df_metadata (pd.DataFrame): DataFrame contendo informações dos arquivos, incluindo a data de atualização e o nome do arquivo. + reference_date (datetime.date, opcional): Data de referência específica para filtrar os arquivos. O Padrão é "%yyyy-%mm-%dd". + + Retorna: + pd.DataFrame: Dataframe filtrado com lista de nomes de arquivos que atendem aos critérios fornecidos e as datas correspondente. + + Levanta: + ValueError: Se não houver arquivos disponíveis para a data específica fornecida. + """ + if reference_date is None: + max_date = df_metadata["data_referencia"].max().date() + log( + f"A data máxima extraida da API da Receita Federal que será utilizada para comparar com os metadados da BD: {max_date}" + ) + + log( + f"A data máxima extraida da API da Receita Federal que será utilizada para gerar partições no Storage: {max_date}" + ) + + return df_metadata[df_metadata["data_referencia"].dt.date == max_date] + + else: + filtered_df = df_metadata[ + df_metadata["data_referencia"].dt.date == reference_date + ] + if filtered_df.empty: + raise ValueError( + f"Não há arquivos disponíveis para a data {reference_date}. Verifique o FTP da Receita Federal." + ) + # pyrefly: ignore [bad-return] + return filtered_df + + +@task +def extract_file_records( + df_metadata: pd.DataFrame, + filename_col: str = "nome_arquivo", + reference_date_col: str = "data_referencia", +) -> list[dict]: + """ + Converte as linhas do DataFrame de metadados em uma lista de dicts + (valores escalares) ANTES do fan-out em tasks paralelas — cada task de + download/processamento passa a receber apenas seu próprio file_name e + reference_date, nunca o DataFrame inteiro compartilhado. + """ + return df_metadata[[filename_col, reference_date_col]].to_dict("records") + + +@task( + retries=3, + retry_delay_seconds=constants.TASK_RETRY_DELAY.value, +) +def download_file( + file_name: str, + url: str, + input_folder: Path, +) -> str: + log(f"Baixando arquivo: {file_name} de {url}") + download_csv_file( + file_name=file_name, + url=url + file_name, + input_folder=input_folder, + ) + return file_name + + +@task( + retries=3, + retry_delay_seconds=constants.TASK_RETRY_DELAY.value, +) +def process_file( + file_name: str, + reference_date: datetime.date, + input_folder: Path, + output_folder: Path, +) -> Path: + file_path = input_folder / file_name + return process_csv_file( + file_path=file_path, + reference_date=reference_date, + output_folder=output_folder, + ) diff --git a/pipelines/datasets/br_rf_cafir/utils.py b/pipelines/datasets/br_rf_cafir/utils.py new file mode 100644 index 0000000000..57baaa7fee --- /dev/null +++ b/pipelines/datasets/br_rf_cafir/utils.py @@ -0,0 +1,214 @@ +""" +General purpose functions for the br_ms_cnes project +""" + +import datetime +import os +from pathlib import Path + +import pandas as pd +import requests +from bs4 import BeautifulSoup + +from pipelines.datasets.br_rf_cafir.constants import ( + constants as br_rf_cafir_constants, +) +from pipelines.utils.utils import log + + +def strip_string(x: pd.DataFrame) -> pd.DataFrame: + """Aplica o strip em uma coluna de um dataframe, caso seja string. + ps: usar com applymap + + Args: + x (pd.Dataframe): Dataframe + + Returns: + pd.Dataframe: Dataframe com valores de linha sem espaços no início e no final das strings + """ + if isinstance(x, str): + # pyrefly: ignore [bad-return] + return x.strip() + return x + + +def remove_ascii_zero_from_df(df: pd.DataFrame) -> pd.DataFrame: + """Remove ASCII 0 (NULL) de colunas tipadas como string de um DataFrame. + Returns: + pd.DataFrame: DataFrame sem ascii 0 (\x00). + """ + # pyrefly: ignore [not-callable] + return df.applymap( + lambda x: x.replace("\x00", "") if isinstance(x, str) else x + ) + + +def requests_url(url: str) -> requests.Response: + xml_body = """ + + + + """ + + headers = { + "Depth": "1", + "Content-Type": "application/xml", + "Accept": "application/xml", + "User-Agent": "Mozilla/5.0", + } + try: + response = requests.request( + method="PROPFIND", + url=url, + headers=headers, + data=xml_body, + timeout=30, + ) + + response.raise_for_status() + + except requests.exceptions.RequestException as e: + log(f"Erro durante a requisição: {e}") + raise + + return response + + +def parse_api_metadata(response_text: str) -> pd.DataFrame: + """ + Extrai metadados de arquivos CSV a partir da resposta da API. + Args: + response_text (str): Texto da resposta da API a ser parseado. + Returns: + pd.DataFrame: Um DataFrame contendo os nomes dos arquivos e suas respectivas datas de atualização. + Raises: + ValueError: Se a quantidade de arquivos extraídos for diferente da quantidade de datas de atualização. + """ + soup = BeautifulSoup(response_text, "lxml") + + items_urls = soup.find_all("d:href") + items_dates = soup.find_all("d:prop") + + files_metadata = [] + for index, item in enumerate(items_urls): + href = item.text + print(href) + if href.endswith(".csv"): + reference_date_str = href.split(".")[-3].replace("D", "202") + reference_date = datetime.datetime.strptime( + reference_date_str, "%Y%m%d" + ) + # .strftime("%Y-%m-%d") + update_date = datetime.datetime.strptime( + items_dates[index].find("d:getlastmodified").text, + "%a, %d %b %Y %H:%M:%S GMT", + ) + + files_metadata.append( + { + "nome_arquivo": href.split("/")[-1], + "data_referencia": reference_date, + "data_modificacao": update_date, + } + ) + return pd.DataFrame(files_metadata) + + +def get_last_date(df_metadata: pd.DataFrame, date_column: str) -> str: + max_date = df_metadata[date_column].max() + return str(max_date.date()) + + +def download_csv_file(url: str, file_name: str, input_folder: Path) -> None: + """ + Faz o download de um arquivo CSV a partir de uma URL e salva em um diretório especificado. + + Args: + url (str): A URL do arquivo CSV a ser baixado. + file_name (str): O nome do arquivo a ser salvo. + input_folder (Path): O diretório onde o arquivo será salvo. + headers (dict): Cabeçalhos HTTP a serem enviados com a requisição. + + Returns: + None + """ + log(f"Downloading--------- {url}") + file_path = input_folder / file_name + response = requests.get(url) + + if response.status_code == 200: + with open(file_path, "wb") as f: + f.write(response.content) + log(f"Downloaded {file_name}") + else: + log( + f"Failed to download {file_name}. Status code: {response.status_code}" + ) + + +def preserve_zeros(x): + """Preserva os zeros a esquerda de um número""" + return x.strip() + + +def process_csv_file( + file_path: Path, + reference_date: datetime.date, + output_folder: Path, +) -> Path: + """ + Lê, limpa e particiona um único arquivo csv de largura fixa. + Args: + file_path (Path): O caminho do arquivo a ser processado. + reference_date (date): Data de referência do arquivo a ser processado + output_folder (Path): Pasta em que o arquivo pocessado terá salvas suas partições + """ + file_name = file_path.name + log(f"Lendo arquivo: {file_name} de : {file_path}") + + df = pd.read_fwf( + file_path, + widths=br_rf_cafir_constants.WIDTHS.value, + names=br_rf_cafir_constants.COLUMN_NAMES.value, + dtype=br_rf_cafir_constants.DTYPE.value, + converters={ + col: preserve_zeros + for col in br_rf_cafir_constants.COLUMN_NAMES.value + }, + encoding="ISO-8859-1", + ) + + # Remove ascii /x00 (zero) - pois dá erro na materialização no BQ + df = remove_ascii_zero_from_df(df) + + # Tira os espacos em branco + # pyrefly: ignore [not-callable] + df = df.applymap(strip_string) + + log(f"Salvando arquivo: {file_name}") + partitions_path = ( + output_folder / "imoveis_rurais" / f"data={reference_date}" + ) + partitions_path.mkdir(parents=True, exist_ok=True) + + # NOTE: Com modificação do formato de divulgação do FTP os arquivos passaram a ser divulgados csvs particionados por UF + # A partir de 2025, a nomenclaruta dos no Storage arquivos mudou para: "imoveis_rurais_uf_numero.csv" no lugar de "imoveris_rurais_numero.csv" + save_path = partitions_path / ( + "imoveis_rurais_" + file_name.split(".")[-2] + ".csv" + ) + + df.to_csv( + save_path, + index=False, + sep=",", + na_rep="", + encoding="utf-8", + escapechar="\\", + ) + + log(f"Arquivo salvo: {save_path.as_posix().split('/')[-1]}") + + del df + os.remove(file_path) + + return save_path diff --git a/pipelines/datasets/br_rf_cnpj/README.md b/pipelines/datasets/br_rf_cnpj/README.md new file mode 100644 index 0000000000..fe91441e3e --- /dev/null +++ b/pipelines/datasets/br_rf_cnpj/README.md @@ -0,0 +1,108 @@ +# Documentação do Conjunto de Dados: CNPJ (Cadastro de Pessoa Jurídica) + +Pipeline de dados de CNPJ (empresas, estabelecimentos, sócios e simples) da +Receita Federal. Este README documenta as mudanças estruturais feitas na +migração desta pipeline (antes `br_me_cnpj`), o motivo de cada uma e seus +impactos no reprocessamento histórico. + +## 1. Correção de encoding na leitura dos CSVs + +Os arquivos CSV de origem da Receita Federal são publicados em `latin1` +(`ISO-8859-1`), não em `utf-8`. A pipeline anterior lia esses arquivos com o +encoding incorreto, o que inseria caracteres inválidos em colunas de +texto livre — por exemplo, `razao_social`, `nome_fantasia` e demais campos com +acentuação. + +A correção lê o CSV de origem como `latin1` e regrava o CSV intermediário como +`utf-8`, preservando a acentuação corretamente (ver `process_csv_*` em +`pipelines/crawler/rf_cnpj/utils.py`). + +## 2. Necessidade de reprocessamento desde 2023 + +Como o bug de encoding afeta todo dado ingerido pela pipeline antiga, os dados +de referência precisaram ser reprocessados do zero para corrigir os caracteres já gravados incorretamente em BigQuery. + +**Data Range:** 2023-06-10 a 2026-05-10 + +**Colunas afetadas:** +- br_me_cnpj.estabelecimentos : bairro, complemento, email, logradouro, nome_fantasia, numero, tipo_logradouro + +- br_me_cnpj.empresas: razao_social + +- br_me_cnpj.socios: nome + +## 3. Migração de `br_me_cnpj` para `br_rf_cnpj` + +O dataset estava indexado sob a organização ME (Ministério da Economia), que +não é a fonte dos dados. A fonte real é a Receita Federal (RF), então o dataset +foi renomeado/migrado para `br_rf_cnpj` para refletir corretamente a +organização de origem. Toda referência à organização e aos metadados do +dataset (incluindo joins em `br_bd_diretorios_brasil__empresa`) foi atualizada +de `br_me_cnpj` para `br_rf_cnpj`. + +## 4. Particionamento por data de referência (`folder_date`, não `last_modified_date`) + +O particionamento passou a ser feito pela **data de referência do arquivo na +fonte** (`folder_date` — o mês/competência a que os dados dizem respeito), e +não mais pela `last_modified_date` (data em que o arquivo foi modificado/gerado +pela Receita Federal). Isso alinha o particionamento ao período que os dados de +fato representam. + + +## 5. Pipeline recorrente de dicionário + +A tabela `br_rf_cnpj__dicionario` traduz os códigos usados nas demais tabelas +(natureza jurídica, qualificação do responsável, motivo da situação +cadastral, município, país, CNAE, etc.) para seus valores legíveis +(`chave` - `valor`). Ela é materializada como `table`, não incremental, +apenas o snapshot mais atual dos códigos. + +O conteúdo vem de duas origens, unidas no modelo: + +- **Arquivos de dicionário publicados pela própria Receita Federal** junto + com a publicação mensal: arquivo de `Cnaes`, `Naturezas`, `Qualificacoes`, + `Municipios`, `Paises`, `Motivos`. Esses arquivos são baixados + e processados por `process_csv_dicionario` + (`pipelines/crawler/rf_cnpj/utils.py`), que lê cada CSV `chave;valor` e adiciona `id_tabela`/`nome_coluna` para identificar a qual tabela e coluna cada código pertence. +- **Entradas manuais** (`dicionario_not_found`, no modelo dbt), cobrindo + chaves que aparecem nos dados reais mas não constam nos arquivos oficiais + da Receita Federal (ex.: códigos `36`, `994`/`393`, `8`/`9`/`32` sem + correspondência na fonte). Sem essas entradas, os valores ficam sem + tradução no dicionário. + +Como `simples` e `dicionario` não têm cobertura temporal por competência +(`NonHistorical`), o polling dessas tabelas compara contra `Table.Update` +(quando rodamos pela última vez), e não contra `Coverage`, ao decidir se há +dado novo a processar. + +## 6. Tabelas legado (`*_legado`) e materialização full-refresh + +Foram criadas tabelas `_legado` (`empresas_legado`, `estabelecimentos_legado`, +`socios_legado`) em staging, contendo os dados históricos migrados do +`br_me_cnpj`. Os modelos principais (`br_rf_cnpj__empresas`, +`br_rf_cnpj__estabelecimentos`, `br_rf_cnpj__socios`) são incrementais e, na +materialização **full-refresh** (primeira execução / quando não há +`is_incremental()`), fazem `union all` com os dados das tabelas `_legado` para +reconstruir o histórico completo (dados até 2023-04-30). Em execuções +incrementais normais, apenas os dados novos vindos da staging atual +(`empresas`, `estabelecimentos`, `socios`) são adicionados. + + +## 7. Alerta: mudança na fonte na representação do documento de sócio pessoa jurídica (14 → 8 caracteres) + +Na tabela de sócios, para registros de sócio **pessoa jurídica** (`tipo = "1"`), +a coluna `documento` armazenava historicamente o **CNPJ completo (14 +caracteres)** do sócio PJ. Em **agosto de 2026** a própria Receita Federal passou a +publicar esse campo com **apenas 8 caracteres** (o `cnpj_basico`, sem +filial/dígito verificador), reduzindo a granularidade da identificação do +sócio PJ na fonte. + +Essa é uma mudança **na fonte**, não introduzida por esta pipeline. + +## Resumo +* Correção de encoding: os CSVs de origem da Receita Federal eram lidos incorretamente, gerando caracteres inválidos em campos de texto (ex.: razão social, nome fantasia). +* Reprocessamento histórico: dados de referência a partir de 2023 foram reprocessados para corrigir os registros afetados pelo erro de encoding. +* Correção da organização de origem: dataset migrado de br_me_cnpj (Ministério da Economia) para br_rf_cnpj, refletindo corretamente a Receita Federal como fonte dos dados. +* Particionamento por data de referência: a partição das tabelas passou a ser baseada na data de referência do arquivo na fonte (competência dos dados), em vez da data de modificação do arquivo. +* Preservação do histórico: dados anteriores a maio/2023, hoje indisponíveis na fonte original, foram preservados em tabelas legado e continuam incluídos na base a cada atualização completa da tabela. +* Alerta de fonte: o documento do sócio pessoa jurídica (tabela `socios`, `tipo = "1"`) passou a ser publicado pela Receita Federal com 8 caracteres em vez dos 14 caracteres do CNPJ completo, a partir da competência 2026-08. \ No newline at end of file