From c46e6af7ca154043beb4035ea4434689e531c858 Mon Sep 17 00:00:00 2001 From: Davi Cavalcante <93160711+DaviMacielCavalcante@users.noreply.github.com> Date: Thu, 3 Sep 2026 10:55:51 -0300 Subject: [PATCH 1/9] =?UTF-8?q?feat:=20cria=C3=A7=C3=A3o=20de=20flow=20par?= =?UTF-8?q?a=20o=20conjunto=20br=5Fms=5Fsim?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- models/br_ms_sim/br_ms_sim__microdados.sql | 7 +- models/br_ms_sim/schema.yml | 4 + pipelines/datasets/br_ms_sim/README.md | 79 +++ pipelines/datasets/br_ms_sim/__init__.py | 0 pipelines/datasets/br_ms_sim/constants.py | 482 ++++++++++++++++++ pipelines/datasets/br_ms_sim/flows.py | 181 +++++++ pipelines/datasets/br_ms_sim/tasks.py | 64 +++ pipelines/datasets/br_ms_sim/utils.py | 536 +++++++++++++++++++++ 8 files changed, 1349 insertions(+), 4 deletions(-) create mode 100644 pipelines/datasets/br_ms_sim/README.md create mode 100644 pipelines/datasets/br_ms_sim/__init__.py create mode 100644 pipelines/datasets/br_ms_sim/constants.py create mode 100644 pipelines/datasets/br_ms_sim/flows.py create mode 100644 pipelines/datasets/br_ms_sim/tasks.py create mode 100644 pipelines/datasets/br_ms_sim/utils.py diff --git a/models/br_ms_sim/br_ms_sim__microdados.sql b/models/br_ms_sim/br_ms_sim__microdados.sql index 82c0698df6..ed444e882b 100644 --- a/models/br_ms_sim/br_ms_sim__microdados.sql +++ b/models/br_ms_sim/br_ms_sim__microdados.sql @@ -1,5 +1,3 @@ --- Microdados SIM: carga 2020-2024; particao end: 2024 (ver README). --- Staging reprocessado com cleaning.py (circunstancia_obito, estado_civil). {{ config( alias="microdados", @@ -8,7 +6,7 @@ partition_by={ "field": "ano", "data_type": "int64", - "range": {"start": 1996, "end": 2024, "interval": 1}, + "range": {"start": 1996, "end": 2031, "interval": 1}, }, cluster_by="sigla_uf", ) @@ -105,5 +103,6 @@ select safe_cast(tipo_nivel_investigador as string) tipo_nivel_investigador, safe_cast(numero_dias_informacao as int64) numero_dias_informacao, safe_cast(fontes_informacao as string) fontes_informacao, - safe_cast(alt_causa as string) alt_causa + safe_cast(alt_causa as string) alt_causa, + safe_cast(dado_preliminar as string) dado_preliminar from {{ set_datalake_project("br_ms_sim_staging.microdados") }} as t diff --git a/models/br_ms_sim/schema.yml b/models/br_ms_sim/schema.yml index 2f35c22881..7f49c8fd19 100644 --- a/models/br_ms_sim/schema.yml +++ b/models/br_ms_sim/schema.yml @@ -213,6 +213,10 @@ models: description: Fontes Informação - name: alt_causa description: Alt. Causa + - name: dado_preliminar + description: > + Indica se o registro vem da versão preliminar da fonte, ainda sujeita + a revisão (1), ou da versão definitiva (0) - name: br_ms_sim__dicionario description: Dicionário para tradução dos códigos do conjunto br_ms_sim. Para taduzir códigos compartilhados entre instituições, como id_municipio, buscar diff --git a/pipelines/datasets/br_ms_sim/README.md b/pipelines/datasets/br_ms_sim/README.md new file mode 100644 index 0000000000..12cc05b8b8 --- /dev/null +++ b/pipelines/datasets/br_ms_sim/README.md @@ -0,0 +1,79 @@ +# br_ms_sim — pipeline + +Carga dos microdados de óbitos não fetais (CID-10) do SIM/DATASUS, do FTP até a +materialização em `basedosdados.br_ms_sim.microdados`. + +Contexto da base, investigações de qualidade e decisões de tratamento estão em +[`models/br_ms_sim/README.md`](../../../models/br_ms_sim/README.md). + +## Sem schedule, por quê + +O DATASUS republica o SIM duas vezes por ano, sem data fixa. O flow é deployado +sem `deploy_schedules`: o deployment existe, aceita execução avulsa e não dispara +sozinho. Para armar depois, basta acrescentar a lista de crons em `flows.py`. + +## Definitivo e preliminar + +A fonte serve o mesmo ano em dois diretórios: + +| Diretório | Conteúdo | +|---|---| +| `SIM/CID10/DORES/` | definitivo, publicado cerca de um ano após o fechamento | +| `SIM/PRELIM/DORES/` | preliminar, ainda sujeito a revisão | + +`resolve_year_source` dá precedência ao definitivo. Reprocessar um ano depois do +fechamento troca o dado e muda `dado_preliminar` de `1` para `0` — não há passo +manual para a virada, só rodar o flow com aquele ano. + +## Parâmetros + +| Parâmetro | Padrão | Efeito | +|---|---|---| +| `ano` | vazio | Vazio pega o ano mais recente da fonte. Preenchido é backfill: o flow pula o poll e não mexe no metadado da fonte | +| `materialize_after_dump` | `True` | Sobe para prod e materializa lá | +| `update_metadata` | `True` | Registra a cobertura materializada | +| `force_run` | `False` | Materializa mesmo sem novidade na fonte | + +Execução de teste no pool de dev, sem tocar em produção: + +```json +{"materialize_after_dump": false, "update_metadata": false, "force_run": true} +``` + +Os padrões escrevem em **produção**, mesmo saindo do pool de teste. + +## Formato da staging + +A staging é CSV desde a carga original. O particionado sai em +`ano=/sigla_uf=/microdados.csv` e sobe com `dump_mode="append"`: os +caminhos são fixos, então reenviar um ano sobrescreve aquele ano e preserva o +resto da série. `overwrite` apagaria o prefixo inteiro, com ele 1996 em diante. + +Trocar para parquet exigiria recriar a tabela externa e, com ela, recarregar +toda a série. + +## Carga manual + +`utils.py` não importa Prefect, então a transformação roda fora do flow: + +```python +from pipelines.datasets.br_ms_sim import utils + +source = utils.resolve_year_source(2025) +utils.download_table("microdados", 2025, source) +utils.clean_table("microdados", 2025, source) +``` + +## Pontos de atenção + +- `dado_preliminar` é coluna nova no modelo. Em `dump_mode="append"` o schema da + tabela externa só é ampliado por `_sync_staging_schema` + (`pipelines/utils/tasks.py`), que abre o cliente do BigQuery sem credencial e + responde `403 bigquery.tables.update`. É o primeiro caso em que esse caminho é + de fato exercitado. +- O dicionário do conjunto vem de uma staging própria, alimentada por + `models/br_ms_sim/code/update_dicionario.py`. As linhas de `dado_preliminar` + precisam ser acrescentadas por lá. +- Os scripts em `models/br_ms_sim/code/microdados/` são a carga anterior e não + conhecem `dado_preliminar`: gravam CSV com uma coluna a menos do que o modelo + espera. Usar o flow. diff --git a/pipelines/datasets/br_ms_sim/__init__.py b/pipelines/datasets/br_ms_sim/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/pipelines/datasets/br_ms_sim/constants.py b/pipelines/datasets/br_ms_sim/constants.py new file mode 100644 index 0000000000..2a2b939196 --- /dev/null +++ b/pipelines/datasets/br_ms_sim/constants.py @@ -0,0 +1,482 @@ +""" +Constantes de br_ms_sim. +""" + +from enum import Enum + + +class constants(Enum): + """Constantes de br_ms_sim.""" + + FTP_FINAL = ( + "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/" + "DO{sigla_uf}{ano}.dbc" + ) + FTP_PRELIM = ( + "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/PRELIM/DORES/" + "DO{sigla_uf}{ano}.dbc" + ) + FTP_FINAL_DIR = ( + "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/" + ) + FTP_PRELIM_DIR = ( + "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/PRELIM/DORES/" + ) + + # Área de trabalho do pod. `input/` recebe os .dbc, `output/` o particionado + # que sobe para o GCS. + PATH = "/tmp/br_ms_sim/" + + SOURCE_FORMAT = "csv" + + UFS = [ + "AC", + "AL", + "AM", + "AP", + "BA", + "CE", + "DF", + "ES", + "GO", + "MA", + "MG", + "MS", + "MT", + "PA", + "PB", + "PE", + "PI", + "PR", + "RJ", + "RN", + "RO", + "RR", + "RS", + "SC", + "SE", + "SP", + "TO", + ] + + TABLES = { + "microdados": { + "file_prefix": "DO", + "partition_columns": ["ano", "sigla_uf"], + }, + } + + DATE_COLUMNS = [ + "data_obito", + "data_nascimento", + "data_atestado", + "data_investigacao", + "data_cadastro", + "data_recebimento", + "data_recebimento_original", + "data_recebimento_original_a", + "data_cadastro_informacao", + "data_cadastro_investigacao", + "data_conclusao_investigacao", + "data_conclusao_caso", + ] + + RENAME = { + "CONTADOR": "sequencial_obito", + "TIPOBITO": "tipo_obito", + "DTOBITO": "data_obito", + "HORAOBITO": "hora_obito", + "NATURAL": "naturalidade", + "CODMUNNATU": "id_municipio_6_naturalidade", + "DTNASC": "data_nascimento", + "IDADE": "idade_raw", + "SEXO": "sexo", + "RACACOR": "raca_cor", + "ESTCIV": "estado_civil", + "ESC": "escolaridade", + "ESC2010": "escolaridade_2010", + "SERIESCFAL": "serie_escolar_falecido", + "OCUP": "ocupacao", + "CODMUNRES": "id_municipio_6_resid", + "CODBAIRES": "codigo_bairro_residencia", + "LOCOCOR": "local_ocorrencia", + "CODESTAB": "codigo_estabelecimento", + "ESTABDESCR": "descricao_estabelecimento", + "CODMUNOCOR": "id_municipio_6_ocor", + "CODBAIOCOR": "codigo_bairro_ocorrencia", + "IDADEMAE": "idade_mae", + "ESCMAE": "escolaridade_mae", + "ESCMAE2010": "escolaridade_mae_2010", + "SERIESCMAE": "serie_escolar_mae", + "OCUPMAE": "ocupacao_mae", + "QTDFILVIVO": "quantidade_filhos_vivos", + "QTDFILMORT": "quantidade_filhos_mortos", + "GRAVIDEZ": "gravidez", + "SEMAGESTAC": "semanas_gestacao", + "GESTACAO": "gestacao", + "PARTO": "parto", + "OBITOPARTO": "obito_parto", + "MORTEPARTO": "morte_parto", + "PESO": "peso", + "TPMORTEOCO": "tipo_morte_ocorrencia", + "OBITOGRAV": "obito_gravidez", + "OBITOPUERP": "obito_puerperio", + "ASSISTMED": "assistencia_medica", + "EXAME": "exame", + "CIRURGIA": "cirurgia", + "NECROPSIA": "necropsia", + "LINHAA": "linha_a", + "LINHAB": "linha_b", + "LINHAC": "linha_c", + "LINHAD": "linha_d", + "LINHAII": "linha_ii", + "CAUSABAS": "causa_basica", + "CB_PRE": "causa_basica_pre", + "COMUNSVOIM": "id_municipio_6_svo_iml", + "DTATESTADO": "data_atestado", + "CIRCOBITO": "circunstancia_obito", + "ACIDTRAB": "acidente_trabalho", + "FONTE": "fonte", + "NUMEROLOTE": "numero_lote", + "TPPOS": "tipo_pos", + "DTINVESTIG": "data_investigacao", + "CAUSABAS_O": "causa_basica_original", + "DTCADASTRO": "data_cadastro", + "ATESTANTE": "atestante", + "STCODIFICA": "status_codificadora", + "CODIFICADO": "codificado", + "VERSAOSIST": "versao_sistema", + "VERSAOSCB": "versao_scb", + "FONTEINV": "fonte_investigacao", + "DTRECEBIM": "data_recebimento", + "ATESTADO": "atestado", + "DTRECORIG": "data_recebimento_original", + "DTRECORIGA": "data_recebimento_original_a", + "CAUSAMAT": "causa_materna", + "ESCMAEAGR1": "escolaridade_mae_2010_agr", + "ESCFALAGR1": "escolaridade_falecido_2010_agr", + "STDOEPIDEM": "status_do_epidem", + "STDONOVA": "status_do_nova", + "DIFDATA": "diferenca_data", + "NUDIASOBCO": "numero_dias_obito_investigacao", + "NUDIASOBIN": "numero_dias_obito_ficha", + "DTCADINV": "data_cadastro_investigacao", + "TPOBITOCOR": "tipo_obito_ocorrencia", + "DTCONINV": "data_conclusao_investigacao", + "FONTES": "fontes", + "TPRESGINFO": "tipo_resgate_informacao", + "TPNIVELINV": "tipo_nivel_investigador", + "NUDIASINF": "numero_dias_informacao", + "DTCADINF": "data_cadastro_informacao", + "DTCONCASO": "data_conclusao_caso", + "FONTESINF": "fontes_informacao", + "ALTCAUSA": "alt_causa", + "CRM": "crm", + } + + COLUMNS = [ + "ano", + "sigla_uf", + "sequencial_obito", + "tipo_obito", + "causa_basica", + "data_obito", + "hora_obito", + "naturalidade", + "data_nascimento", + "idade", + "sexo", + "raca_cor", + "estado_civil", + "escolaridade", + "ocupacao", + "codigo_bairro_residencia", + "id_municipio_residencia", + "local_ocorrencia", + "codigo_bairro_ocorrencia", + "id_municipio_ocorrencia", + "idade_mae", + "escolaridade_mae", + "ocupacao_mae", + "quantidade_filhos_vivos", + "quantidade_filhos_mortos", + "gravidez", + "gestacao", + "parto", + "obito_parto", + "morte_parto", + "peso", + "obito_gravidez", + "obito_puerperio", + "assistencia_medica", + "exame", + "cirurgia", + "necropsia", + "linha_a", + "linha_b", + "linha_c", + "linha_d", + "linha_ii", + "circunstancia_obito", + "acidente_trabalho", + "fonte", + "codigo_estabelecimento", + "atestante", + "data_atestado", + "tipo_pos", + "data_investigacao", + "causa_basica_original", + "data_cadastro", + "fonte_investigacao", + "data_recebimento", + "causa_basica_pre", + "tipo_obito_ocorrencia", + "tipo_morte_ocorrencia", + "data_cadastro_informacao", + "data_cadastro_investigacao", + "id_municipio_svo_iml", + "data_recebimento_original", + "data_recebimento_original_a", + "causa_materna", + "status_do_epidem", + "status_do_nova", + "serie_escolar_falecido", + "serie_escolar_mae", + "escolaridade_2010", + "escolaridade_mae_2010", + "escolaridade_falecido_2010_agr", + "escolaridade_mae_2010_agr", + "semanas_gestacao", + "diferenca_data", + "data_conclusao_investigacao", + "data_conclusao_caso", + "numero_dias_obito_investigacao", + "id_municipio_naturalidade", + "descricao_estabelecimento", + "crm", + "numero_lote", + "status_codificadora", + "codificado", + "versao_sistema", + "versao_scb", + "atestado", + "numero_dias_obito_ficha", + "fontes", + "tipo_resgate_informacao", + "tipo_nivel_investigador", + "numero_dias_informacao", + "fontes_informacao", + "alt_causa", + "dado_preliminar", + ] + + # Códigos que representam ausência de informação e viram NULL. + NULLIFY = { + "local_ocorrencia": ["0", "6", "7", "9"], + "sexo": ["0", "6", "7", "9"], + "raca_cor": ["0", "6", "7", "9"], + "estado_civil": ["0", "9"], + "escolaridade": ["0", "6", "7", "9", "A"], + "escolaridade_mae": ["0", "6", "7", "9", "A"], + "escolaridade_2010": ["9"], + "escolaridade_mae_2010": ["9"], + "gravidez": ["0", "9"], + "gestacao": ["0", "9"], + "parto": ["0", "3", "4", "5", "6", "7", "9"], + "obito_parto": ["0", "4", "5", "6", "7", "9"], + "morte_parto": ["0", "4", "5", "6", "7", "9"], + "obito_gravidez": ["0", "3", "4", "5", "6", "7", "9"], + "obito_puerperio": ["0", "4", "5", "6", "7", "9"], + "assistencia_medica": ["0", "4", "5", "6", "7", "9"], + "exame": ["0", "4", "5", "6", "7", "9"], + "cirurgia": ["0", "4", "5", "6", "7", "9"], + "necropsia": ["0", "4", "5", "6", "7", "9"], + "acidente_trabalho": ["0", "4", "5", "6", "7", "9"], + "circunstancia_obito": ["0", "5", "6", "7", "9"], + "fonte": ["0", "5", "6", "7", "9"], + "fonte_investigacao": ["0", "9"], + "tipo_morte_ocorrencia": ["9"], + } + + # Código -> rótulo. Aplicado depois do NULLIFY. + RECODE = { + "tipo_obito": {"1": "fetal", "2": "nao-fetal"}, + "sexo": {"1": "masculino", "2": "feminino"}, + "raca_cor": { + "1": "branca", + "2": "preta", + "3": "amarela", + "4": "parda", + "5": "indigena", + }, + "estado_civil": { + "1": "solteiro", + "2": "casado", + "3": "viuvo", + "4": "separado judicialmente/divorciado", + "5": "uniao consensual", + }, + "escolaridade": { + "1": "nenhuma", + "2": "1 a 3 anos", + "3": "4 a 7 anos", + "4": "8 a 11 anos", + "5": "12 e mais", + "8": "9 a 11 anos", + }, + "escolaridade_mae": { + "1": "nenhuma", + "2": "1 a 3 anos", + "3": "4 a 7 anos", + "4": "8 a 11 anos", + "5": "12 e mais", + "8": "9 a 11 anos", + }, + "escolaridade_2010": { + "0": "sem escolaridade", + "1": "fundamental I", + "2": "fundamental II", + "3": "medio", + "4": "superior incompleto", + "5": "superior completo", + }, + "escolaridade_mae_2010": { + "0": "sem escolaridade", + "1": "fundamental I", + "2": "fundamental II", + "3": "medio", + "4": "superior incompleto", + "5": "superior completo", + }, + "escolaridade_mae_2010_agr": { + "00": "sem escolaridade", + "01": "fundamental I incompleto", + "02": "fundamental I completo", + "03": "fundamental II incompleto", + "04": "fundamental II completo", + "05": "ensino medio incompleto", + "06": "ensino medio completo", + "07": "ensino superior incompleto", + "08": "ensino superior completo", + "09": "ignorado", + "10": "fundamental I incompleto ou inespecifico", + "11": "fundamental II incompleto ou inespecifico", + "12": "ensino medio incompleto ou inespecifico", + }, + "escolaridade_falecido_2010_agr": { + "00": "sem escolaridade", + "01": "fundamental I incompleto", + "02": "fundamental I completo", + "03": "fundamental II incompleto", + "04": "fundamental II completo", + "05": "ensino medio incompleto", + "06": "ensino medio completo", + "07": "ensino superior incompleto", + "08": "ensino superior completo", + "09": "ignorado", + "10": "fundamental I incompleto ou inespecifico", + "11": "fundamental II incompleto ou inespecifico", + "12": "ensino medio incompleto ou inespecifico", + }, + "local_ocorrencia": { + "1": "hospital", + "2": "outro estabelecimento de saude", + "3": "domicilio", + "4": "via publica", + "5": "outros", + "6": "aldeia indigena", + }, + "gravidez": {"1": "unica", "2": "dupla", "3": "tripla e mais"}, + "gestacao": { + "A": "21 a 27 semanas", + "1": "menos de 22 semanas", + "2": "22 a 27 semanas", + "3": "28 a 31 semanas", + "4": "32 a 36 semanas", + "5": "37 a 41 semanas", + "6": "42 semanas ou mais", + "7": "28 semanas ou mais", + "8": "28 a 36 semanas", + }, + "parto": {"1": "vaginal", "2": "cesareo"}, + "obito_parto": {"1": "antes", "2": "durante", "3": "depois"}, + "morte_parto": {"1": "antes", "2": "durante", "3": "depois"}, + "obito_gravidez": {"1": "sim", "2": "nao"}, + "obito_puerperio": { + "1": "0 a 42 dias", + "2": "43 dias a 1 ano", + "3": "nao", + }, + "assistencia_medica": {"1": "sim", "2": "nao"}, + "exame": {"1": "sim", "2": "nao"}, + "cirurgia": {"1": "sim", "2": "nao"}, + "necropsia": {"1": "sim", "2": "nao"}, + "circunstancia_obito": { + "1": "acidente", + "2": "suicidio", + "3": "homicidio", + "4": "outro", + }, + "acidente_trabalho": {"1": "sim", "2": "nao"}, + "fonte": { + "1": "boletim de ocorrencia", + "2": "hospital", + "3": "familia", + "4": "outro", + }, + "fonte_investigacao": { + "1": "comite de mortalidade materna e/ou infantil", + "2": "visita familiar / entrevista familia", + "3": "estabelecimento de saude / prontuario", + "4": "relacionamento com outros bancos de dados", + "5": "SVO", + "6": "IML", + "7": "outra fonte", + "8": "multiplas fontes", + }, + "status_do_epidem": {"1": "sim", "0": "nao"}, + "status_do_nova": {"1": "sim", "0": "nao"}, + "atestante": { + "1": "sim", + "2": "substituto", + "3": "IML", + "4": "SVO", + "5": "outros", + }, + "tipo_pos": {"S": "sim", "N": "nao"}, + "status_codificadora": {"S": "sim", "N": "nao"}, + "codificado": {"S": "sim", "N": "nao"}, + "tipo_obito_ocorrencia": { + "1": "durante a gestacao", + "2": "duranto abortamento", + "3": "apos abortamento", + "4": "no parto ou ate 1 hora apos o parto", + "5": "no puerperio (ate 42 dias do termino da gestacao)", + "6": "entre o 43º dia e ate um ano apos o termino da gestacao", + "7": "investigacao nao identificou o momento do obito", + "8": "mais de 1 ano apos o parto", + "9": "outras", + }, + "tipo_morte_ocorrencia": { + "1": "na gravidez", + "2": "no parto", + "3": "no aborto", + "4": "ate 42 dias apos o parto", + "5": "de 43 dias ate 1 ano apos o parto", + "8": "nao ocorreu nestes periodos", + }, + "tipo_resgate_informacao": { + "1": "nao acrescentou nem corrigiu informacao", + "2": "sim, permitiu o resgate de novas informacoes", + "3": "sim, permitiu a correcao de alguma das causas informadas originalmente", + }, + "tipo_nivel_investigador": { + "E": "estadual", + "R": "regional", + "M": "municipal", + }, + } + + # Colunas em que qualquer valor fora do RECODE também vira NULL: o DATASUS + # grava códigos que não constam do dicionário da própria fonte. + RECODE_STRICT = ["estado_civil", "circunstancia_obito"] diff --git a/pipelines/datasets/br_ms_sim/flows.py b/pipelines/datasets/br_ms_sim/flows.py new file mode 100644 index 0000000000..451e40c777 --- /dev/null +++ b/pipelines/datasets/br_ms_sim/flows.py @@ -0,0 +1,181 @@ +""" +Flows de br_ms_sim — Prefect 3. +""" + +from prefect import flow + +from pipelines.datasets.br_ms_sim.constants import constants +from pipelines.datasets.br_ms_sim.tasks import ( + clean_table, + download_table, + get_source_max_year, + resolve_year_source, +) +from pipelines.utils.metadata.domain import ( + AllFree, + CoverageSpec, + DateFormat, + YearOnly, +) +from pipelines.utils.metadata.tasks import ( + commit_source_update_task, + poll_source_for_update_task, + register_table_materialization_task, +) +from pipelines.utils.tasks import ( + rename_flow_run_dataset_table, + run_dbt, + upload_to_gcs, +) + +DATE_FORMAT = DateFormat.YEAR + + +def coverage(table_id: str) -> CoverageSpec: + """Devolve a cobertura da tabela. + + Raises: + ValueError: Se a tabela não constar de `constants.TABLES`. + """ + if table_id not in constants.TABLES.value: + raise ValueError(f"tabela sem cobertura definida: {table_id}") + return AllFree( + date_column=YearOnly(col="ano"), + date_format=DATE_FORMAT, + ) + + +def run_ms_sim( + *, + dataset_id: str, + table_id: str, + ano: int | None, + materialize_after_dump: bool, + update_metadata: bool, + target: str, + force_run: bool, + dump_mode: str, + source_format: str, +) -> None: + """Executa o ciclo baixar, limpar, subir, dbt e metadados de um ano.""" + # pyrefly: ignore [unused-coroutine] + rename_flow_run_dataset_table( + prefix="Dump: ", dataset_id=dataset_id, table_id=table_id + ) + + backfill = ano is not None + source_max_year = get_source_max_year() + ano = int(ano if backfill else source_max_year) + + if not force_run and not backfill: + has_new_data = poll_source_for_update_task( + dataset_id=dataset_id, + table_id=table_id, + source_max_date=source_max_year, + env="prod", + date_format=DATE_FORMAT, + compare_against="coverage", + ) + if not has_new_data: + print(f"Não há atualizações para a tabela {table_id}!") + return + + if not backfill: + commit_source_update_task( + dataset_id=dataset_id, + table_id=table_id, + source_max_date=source_max_year, + env="prod", + date_format=DATE_FORMAT, + update_metadata=update_metadata, + materialize_after_dump=materialize_after_dump, + ) + + source = resolve_year_source(ano) + print(f"Carregando {ano} a partir do diretório {source}") + + download_table(table_id=table_id, ano=ano, source=source) + filepath = clean_table(table_id=table_id, ano=ano, source=source) + + upload_to_gcs( + data_path=filepath, + dataset_id=dataset_id, + table_id=table_id, + bucket_name="basedosdados-dev", + dump_mode=dump_mode, + source_format=source_format, + ) + + 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=filepath, + dataset_id=dataset_id, + table_id=table_id, + bucket_name="basedosdados", + dump_mode=dump_mode, + source_format=source_format, + ) + + run_dbt( + dataset_id=dataset_id, + table_id=table_id, + dbt_command="run/test", + target=target, + ) + + if update_metadata: + register_table_materialization_task( + dataset_id=dataset_id, + table_id=table_id, + coverage=coverage(table_id), + env="prod", + bq_project="basedosdados", + ) + + +def ms_sim_flow( + table_id: str, + dump_mode: str = "append", + source_format: str = "csv", +): + """Carimba o flow de uma tabela.""" + + @flow( + name=f"br_ms_sim__{table_id}", + log_prints=True, + ) + def table_flow( + dataset_id: str = "br_ms_sim", + table_id: str = table_id, + ano: int | None = None, + materialize_after_dump: bool = True, + update_metadata: bool = True, + target: str = "prod", + force_run: bool = False, + ) -> None: + """Carrega um ano do SIM, do FTP do DATASUS até a materialização.""" + run_ms_sim( + dataset_id=dataset_id, + table_id=table_id, + ano=ano, + materialize_after_dump=materialize_after_dump, + update_metadata=update_metadata, + target=target, + force_run=force_run, + dump_mode=dump_mode, + source_format=source_format, + ) + + return table_flow + + +br_ms_sim__microdados = ms_sim_flow("microdados") diff --git a/pipelines/datasets/br_ms_sim/tasks.py b/pipelines/datasets/br_ms_sim/tasks.py new file mode 100644 index 0000000000..b1987afa3f --- /dev/null +++ b/pipelines/datasets/br_ms_sim/tasks.py @@ -0,0 +1,64 @@ +""" +Tasks de br_ms_sim. + +Cada task embrulha uma função de `utils.py`, onde fica a lógica. +""" + +from pathlib import Path + +from prefect import task + +from pipelines.datasets.br_ms_sim import utils + + +@task(retries=3, retry_delay_seconds=30) +def get_source_max_year() -> str: + """Lê na fonte até que ano ela publicou. + + Returns: + O ano mais recente, no formato `%Y`. + """ + return utils.get_source_max_year() + + +@task(retries=3, retry_delay_seconds=30) +def resolve_year_source(ano: int) -> str: + """Diz de qual diretório o ano deve ser baixado. + + Args: + ano: Ano a resolver. + + Returns: + `"definitivo"` ou `"preliminar"`. + """ + return utils.resolve_year_source(ano) + + +@task(retries=3, retry_delay_seconds=60) +def download_table(table_id: str, ano: int, source: str) -> Path: + """Baixa os arquivos das UFs do ano. + + Args: + table_id: Slug da tabela. + ano: Ano a baixar. + source: `"definitivo"` ou `"preliminar"`. + + Returns: + O diretório de entrada com os arquivos baixados. + """ + return utils.download_table(table_id=table_id, ano=ano, source=source) + + +@task +def clean_table(table_id: str, ano: int, source: str) -> Path: + """Limpa o ano já baixado e grava o particionado. + + Args: + table_id: Slug da tabela. + ano: Ano a limpar. + source: `"definitivo"` ou `"preliminar"`. + + Returns: + O diretório particionado, no formato esperado por `upload_to_gcs`. + """ + return utils.clean_table(table_id=table_id, ano=ano, source=source) diff --git a/pipelines/datasets/br_ms_sim/utils.py b/pipelines/datasets/br_ms_sim/utils.py new file mode 100644 index 0000000000..f581cb8462 --- /dev/null +++ b/pipelines/datasets/br_ms_sim/utils.py @@ -0,0 +1,536 @@ +""" +Download, limpeza e particionamento de br_ms_sim. + +O módulo não importa Prefect. As funções são chamadas pelas tasks de +`tasks.py` e também podem ser executadas diretamente. +""" + +import os +import shutil +import tempfile +import urllib.request +from pathlib import Path + +import basedosdados as bd +import pandas as pd +from datasus_dbc import decompress as dbc2dbf +from dbfread import DBF + +from pipelines.datasets.br_ms_sim.constants import constants + +FINAL, PRELIM = "definitivo", "preliminar" + + +def build_paths(table_id: str, ano: int) -> tuple[Path, Path]: + """Cria os diretórios de trabalho do ano e devolve os dois caminhos. + + O `output/` é apagado a cada chamada; o `input/` é preservado. O ano compõe + o caminho, de modo que execuções de anos diferentes não compartilham + diretório. + + Args: + table_id: Slug da tabela. + ano: Ano da carga. + + Returns: + Os caminhos de `input/` e de `output/`, nessa ordem. + """ + base = Path(constants.PATH.value) / table_id / str(ano) + input_dir, output_dir = base / "input", base / "output" + shutil.rmtree(output_dir, ignore_errors=True) + input_dir.mkdir(parents=True, exist_ok=True) + output_dir.mkdir(parents=True, exist_ok=True) + return input_dir, output_dir + + +def list_ftp_years(directory_url: str) -> set[int]: + """Lê a listagem de um diretório do FTP e devolve os anos com arquivo. + + Args: + directory_url: URL do diretório, terminada em barra. + + Returns: + Os anos extraídos dos nomes no padrão `DO.dbc`. + """ + with urllib.request.urlopen(directory_url, timeout=120) as response: + listing = response.read().decode("latin1") + + years = set() + for line in listing.splitlines(): + name = line.split()[-1] if line.strip() else "" + if not name.upper().startswith("DO") or not name.lower().endswith( + ".dbc" + ): + continue + stem = name[:-4] + if len(stem) >= 8 and stem[-4:].isdigit(): + years.add(int(stem[-4:])) + return years + + +def get_source_max_year() -> str: + """Devolve o ano mais recente publicado na fonte. + + Considera os dois diretórios, já que o ano corrente costuma existir apenas + no preliminar. O valor é a competência do dado, não a data da consulta. + + Returns: + O ano mais recente, no formato `%Y`. + + Raises: + RuntimeError: Se nenhum arquivo for encontrado nos dois diretórios. + """ + years = list_ftp_years(constants.FTP_FINAL_DIR.value) | list_ftp_years( + constants.FTP_PRELIM_DIR.value + ) + if not years: + raise RuntimeError( + "nenhum arquivo DO*.dbc encontrado no FTP do DATASUS — a fonte " + "mudou de layout ou está fora do ar" + ) + return str(max(years)) + + +def resolve_year_source(ano: int) -> str: + """Diz de qual diretório o ano deve ser baixado. + + O definitivo tem precedência sobre o preliminar. Um ano fechado pelo DATASUS + passa a existir nos dois diretórios, e reprocessá-lo substitui o dado + preliminar pelo definitivo. + + Args: + ano: Ano a resolver. + + Returns: + `"definitivo"` ou `"preliminar"`. + + Raises: + FileNotFoundError: Se o ano não existir em nenhum dos dois diretórios. + """ + if ano in list_ftp_years(constants.FTP_FINAL_DIR.value): + return FINAL + if ano in list_ftp_years(constants.FTP_PRELIM_DIR.value): + return PRELIM + raise FileNotFoundError( + f"ano {ano} não está no FTP do DATASUS, nem definitivo nem preliminar" + ) + + +def download_year(ano: int, source: str, input_dir: Path) -> Path: + """Baixa os arquivos `.dbc` das 27 UFs do ano. + + UF ausente na fonte é registrada no log e ignorada; a carga prossegue com as + demais. + + Args: + ano: Ano a baixar. + source: `"definitivo"` ou `"preliminar"`. + input_dir: Diretório de destino. + + Returns: + O diretório de destino. + + Raises: + RuntimeError: Se nenhuma das 27 UFs for baixada. + """ + template = ( + constants.FTP_FINAL.value + if source == FINAL + else constants.FTP_PRELIM.value + ) + + missing = [] + for sigla_uf in constants.UFS.value: + url = template.format(sigla_uf=sigla_uf, ano=ano) + destination = input_dir / f"DO{sigla_uf}{ano}.dbc" + try: + urllib.request.urlretrieve(url, filename=str(destination)) + print(f" {destination.name}: {destination.stat().st_size} bytes") + except Exception as error: + destination.unlink(missing_ok=True) + missing.append(sigla_uf) + print(f" {sigla_uf}: ausente na fonte — {error}") + + if len(missing) == len(constants.UFS.value): + raise RuntimeError( + f"nenhuma UF baixada para {ano} ({source}) — não há o que carregar" + ) + if missing: + print(f"UFs ausentes em {ano}: {', '.join(missing)}") + + return input_dir + + +def read_dbc(filepath: Path, encoding: str = "iso-8859-1") -> pd.DataFrame: + """Descompacta um arquivo `.dbc` e lê o `.dbf` resultante. + + Args: + filepath: Caminho do arquivo `.dbc`. + encoding: Codificação do `.dbf`. + + Returns: + O conteúdo do arquivo. + """ + file_descriptor, tmp_path = tempfile.mkstemp( + suffix=".dbf", dir=tempfile.gettempdir() + ) + os.close(file_descriptor) + try: + dbc2dbf(str(filepath), tmp_path) + table = DBF(tmp_path, encoding=encoding, load=True) + return pd.DataFrame(iter(table)) + finally: + Path(tmp_path).unlink(missing_ok=True) + + +def load_municipios() -> pd.DataFrame: + """Lê o de-para de código de município do diretório da Base dos Dados. + + Returns: + As colunas `id_municipio` e `id_municipio_6`, ambas como texto. + """ + return bd.read_sql( + "SELECT id_municipio, id_municipio_6 " + "FROM `basedosdados-dev.br_bd_diretorios_brasil.municipio`", + billing_project_id="basedosdados-dev", + from_file=True, + ).astype(str) + + +def convert_municipio_6_to_7( + dataframe: pd.DataFrame, + column_6: str, + column_7: str, + municipios: pd.DataFrame, +) -> pd.DataFrame: + """Substitui o código de município de 6 dígitos pelo de 7 do IBGE. + + Args: + dataframe: Dados a converter. + column_6: Coluna de origem, com o código de 6 dígitos. + column_7: Nome da coluna resultante. + municipios: De-para devolvido por `load_municipios`. + + Returns: + Os dados com a coluna convertida, ou inalterados se `column_6` não + existir. + """ + if column_6 not in dataframe.columns: + return dataframe + dataframe = dataframe.merge( + municipios[["id_municipio_6", "id_municipio"]], + how="left", + left_on=column_6, + right_on="id_municipio_6", + ) + dataframe = dataframe.drop(columns=[column_6, "id_municipio_6"]) + return dataframe.rename(columns={"id_municipio": column_7}) + + +def convert_municipio_resid_ocor( + dataframe: pd.DataFrame, ano: int, municipios: pd.DataFrame +) -> pd.DataFrame: + """Converte os municípios de residência e de ocorrência. + + Até 2005 a fonte grava esses campos com 7 dígitos, e as colunas são apenas + renomeadas. + + Args: + dataframe: Dados a converter. + ano: Ano do arquivo. + municipios: De-para devolvido por `load_municipios`. + + Returns: + Os dados com as duas colunas convertidas. + """ + if ano <= 2005: + renames = {} + if "id_municipio_6_resid" in dataframe.columns: + renames["id_municipio_6_resid"] = "id_municipio_residencia" + if "id_municipio_6_ocor" in dataframe.columns: + renames["id_municipio_6_ocor"] = "id_municipio_ocorrencia" + return dataframe.rename(columns=renames) + + dataframe = convert_municipio_6_to_7( + dataframe, + "id_municipio_6_resid", + "id_municipio_residencia", + municipios, + ) + return convert_municipio_6_to_7( + dataframe, "id_municipio_6_ocor", "id_municipio_ocorrencia", municipios + ) + + +def parse_date(value: object) -> str | None: + """Converte uma data no formato `DDMMAAAA`. + + Args: + value: Valor bruto do arquivo. + + Returns: + A data no formato `AAAA-MM-DD`, ou None se o valor não for uma data. + """ + if not value: + return None + text = str(value).strip() + if len(text) < 8 or text == "00000000": + return None + return f"{text[4:8]}-{text[2:4]}-{text[0:2]}" + + +def parse_hora(value: object) -> str | None: + """Converte um horário no formato `HHMM`. + + Args: + value: Valor bruto do arquivo. + + Returns: + O horário no formato `HH:MM:00`, ou None se o valor for mais curto que + quatro dígitos. + """ + if not value: + return None + text = str(value).strip() + if len(text) < 4: + return None + text = text.zfill(4) + return f"{text[0:2]}:{text[2:4]}:00" + + +def parse_idade(value: object) -> float | None: + """Converte a idade codificada do SIM em anos. + + O primeiro dígito é a unidade — 1 minuto, 2 hora, 3 mês, 4 ano, 5 ano acima + de 100 — e os demais são a quantidade. + + Args: + value: Valor bruto do arquivo. + + Returns: + A idade em anos, com duas casas decimais, ou None se a unidade for + desconhecida ou a quantidade não for numérica. + """ + if not value: + return None + text = str(value).strip() + if len(text) < 2: + return None + unit = text[0] + try: + amount = int(text[1:]) + except ValueError: + return None + if unit == "1": + idade = 0.0 + elif unit == "2": + idade = amount / 365 + elif unit == "3": + idade = amount / 12 + elif unit == "4": + idade = float(amount) + elif unit == "5": + idade = float(100 + amount) + else: + return None + return round(idade, 2) + + +def recode_columns(dataframe: pd.DataFrame) -> pd.DataFrame: + """Nulifica os códigos de ausência e traduz os demais para rótulos. + + Args: + dataframe: Dados a recodificar. + + Returns: + Uma cópia dos dados com as colunas de `NULLIFY` e `RECODE` tratadas. + """ + dataframe = dataframe.copy() + + for column, invalid in constants.NULLIFY.value.items(): + if column in dataframe.columns: + dataframe[column] = dataframe[column].replace(invalid, None) + + for column, mapping in constants.RECODE.value.items(): + if column in dataframe.columns: + dataframe[column] = dataframe[column].replace(mapping) + + # Código fora do dicionário da própria fonte atravessaria o replace intacto + # e viraria rótulo inválido na tabela. + for column in constants.RECODE_STRICT.value: + if column not in dataframe.columns: + continue + valid = set(constants.RECODE.value[column].values()) + unknown = dataframe[column].notna() & ~dataframe[column].isin(valid) + dataframe.loc[unknown, column] = None + + if "peso" in dataframe.columns: + dataframe["peso"] = dataframe["peso"].replace(["0"], None) + + return dataframe + + +def ensure_schema_columns(dataframe: pd.DataFrame) -> pd.DataFrame: + """Completa as colunas ausentes e aplica a ordem da arquitetura. + + Args: + dataframe: Dados de um ano, que pode não trazer todas as colunas. + + Returns: + Os dados com todas as colunas de `COLUMNS`, na ordem do modelo. + """ + for column in constants.COLUMNS.value: + if column not in dataframe.columns: + dataframe[column] = None + return dataframe[constants.COLUMNS.value] + + +def process_file( + filepath: Path, + ano: int, + sigla_uf: str, + municipios: pd.DataFrame, + is_prelim: bool, +) -> pd.DataFrame: + """Lê o arquivo de uma UF e devolve os dados no schema da arquitetura. + + Args: + filepath: Caminho do arquivo `.dbc`. + ano: Ano do arquivo. + sigla_uf: Sigla da unidade da federação. + municipios: De-para devolvido por `load_municipios`. + is_prelim: Se o arquivo vem do diretório preliminar, o que define o + valor de `dado_preliminar`. + + Returns: + Os dados renomeados, convertidos e recodificados. + """ + dataframe = read_dbc(filepath) + dataframe.columns = dataframe.columns.str.upper() + dataframe = dataframe.astype(str).replace( + {"None": None, "nan": None, "": None} + ) + dataframe = dataframe.replace("NA", None) + + dataframe = dataframe.rename(columns=constants.RENAME.value) + dataframe = dataframe.drop(columns=["ORIGEM", "UFINFORM"], errors="ignore") + + dataframe["ano"] = ano + dataframe["sigla_uf"] = sigla_uf + dataframe["dado_preliminar"] = "1" if is_prelim else "0" + + dataframe = convert_municipio_resid_ocor(dataframe, ano, municipios) + dataframe = convert_municipio_6_to_7( + dataframe, + "id_municipio_6_svo_iml", + "id_municipio_svo_iml", + municipios, + ) + dataframe = convert_municipio_6_to_7( + dataframe, + "id_municipio_6_naturalidade", + "id_municipio_naturalidade", + municipios, + ) + + for column in constants.DATE_COLUMNS.value: + if column in dataframe.columns: + dataframe[column] = dataframe[column].apply(parse_date) + + if "hora_obito" in dataframe.columns: + dataframe["hora_obito"] = dataframe["hora_obito"].apply(parse_hora) + + if "idade_raw" in dataframe.columns: + dataframe["idade"] = dataframe["idade_raw"].apply(parse_idade) + dataframe = dataframe.drop(columns=["idade_raw"]) + + dataframe = recode_columns(dataframe) + return ensure_schema_columns(dataframe) + + +def clean_year( + table_id: str, + ano: int, + source: str, + input_dir: Path, + output_dir: Path, +) -> Path: + """Limpa os arquivos do ano e grava o particionado em CSV. + + Args: + table_id: Slug da tabela, que define as colunas de partição. + ano: Ano processado. + source: `"definitivo"` ou `"preliminar"`. + input_dir: Diretório com os arquivos `.dbc`. + output_dir: Raiz do particionado. + + Returns: + O diretório particionado, no formato esperado por `upload_to_gcs`. + + Raises: + RuntimeError: Se nenhum arquivo do ano for encontrado em `input_dir`. + """ + table = constants.TABLES.value[table_id] + partition_columns = table["partition_columns"] + file_prefix = table["file_prefix"] + municipios = load_municipios() + is_prelim = source == PRELIM + total = 0 + + for filepath in sorted(input_dir.glob(f"{file_prefix}*{ano}.dbc")): + sigla_uf = filepath.stem[len(file_prefix) :][:2] + dataframe = process_file( + filepath, ano, sigla_uf, municipios, is_prelim + ) + + partition = output_dir / f"ano={ano}" / f"sigla_uf={sigla_uf}" + partition.mkdir(parents=True, exist_ok=True) + dataframe.drop(columns=partition_columns).to_csv( + partition / f"{table_id}.csv", index=False + ) + total += len(dataframe) + print(f" {sigla_uf}: {len(dataframe):,} linhas") + + if total == 0: + raise RuntimeError( + f"nenhum arquivo processado para {ano} — `input/` está vazio" + ) + + print(f"{ano} ({source}): {total:,} linhas") + return output_dir + + +def download_table(table_id: str, ano: int, source: str) -> Path: + """Prepara os diretórios e baixa o ano. + + Args: + table_id: Slug da tabela. + ano: Ano a baixar. + source: `"definitivo"` ou `"preliminar"`. + + Returns: + O diretório de entrada com os arquivos baixados. + """ + input_dir, _ = build_paths(table_id, ano) + return download_year(ano=ano, source=source, input_dir=input_dir) + + +def clean_table(table_id: str, ano: int, source: str) -> Path: + """Limpa o ano já baixado. + + Args: + table_id: Slug da tabela. + ano: Ano a limpar. + source: `"definitivo"` ou `"preliminar"`. + + Returns: + O diretório particionado, no formato esperado por `upload_to_gcs`. + """ + input_dir, output_dir = build_paths(table_id, ano) + return clean_year( + table_id=table_id, + ano=ano, + source=source, + input_dir=input_dir, + output_dir=output_dir, + ) From a3037ccec2b068d5c87be90bbd84052040d0bff8 Mon Sep 17 00:00:00 2001 From: Davi Cavalcante <93160711+DaviMacielCavalcante@users.noreply.github.com> Date: Thu, 3 Sep 2026 11:38:42 -0300 Subject: [PATCH 2/9] fix: ajustando para pedir * G de RAM --- pipelines/datasets/br_ms_sim/flows.py | 2 ++ pipelines/datasets/br_ms_sim/utils.py | 4 +++- 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/pipelines/datasets/br_ms_sim/flows.py b/pipelines/datasets/br_ms_sim/flows.py index 451e40c777..24a06cb43d 100644 --- a/pipelines/datasets/br_ms_sim/flows.py +++ b/pipelines/datasets/br_ms_sim/flows.py @@ -175,6 +175,8 @@ def table_flow( source_format=source_format, ) + # pyrefly: ignore [missing-attribute] + table_flow.job_variables = {"memory": "8Gi"} return table_flow diff --git a/pipelines/datasets/br_ms_sim/utils.py b/pipelines/datasets/br_ms_sim/utils.py index f581cb8462..f45e02663c 100644 --- a/pipelines/datasets/br_ms_sim/utils.py +++ b/pipelines/datasets/br_ms_sim/utils.py @@ -177,7 +177,9 @@ def read_dbc(filepath: Path, encoding: str = "iso-8859-1") -> pd.DataFrame: os.close(file_descriptor) try: dbc2dbf(str(filepath), tmp_path) - table = DBF(tmp_path, encoding=encoding, load=True) + # `load=True` guardaria os registros também na própria DBF, dobrando o + # pico de memória nas UFs grandes. + table = DBF(tmp_path, encoding=encoding, load=False) return pd.DataFrame(iter(table)) finally: Path(tmp_path).unlink(missing_ok=True) From b32ab64724976b207c0e822fbd8e8ac682875020 Mon Sep 17 00:00:00 2001 From: Davi Cavalcante <93160711+DaviMacielCavalcante@users.noreply.github.com> Date: Thu, 3 Sep 2026 12:06:50 -0300 Subject: [PATCH 3/9] =?UTF-8?q?chore:=20verificando=20se=20os=20recursos?= =?UTF-8?q?=20s=C3=A3o=20alocados=20corretamente?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pipelines/datasets/br_ms_sim/flows.py | 11 +++++++++- pipelines/datasets/br_ms_sim/utils.py | 31 +++++++++++++++++++++++++++ 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/pipelines/datasets/br_ms_sim/flows.py b/pipelines/datasets/br_ms_sim/flows.py index 24a06cb43d..8fcfe3cf83 100644 --- a/pipelines/datasets/br_ms_sim/flows.py +++ b/pipelines/datasets/br_ms_sim/flows.py @@ -11,6 +11,7 @@ get_source_max_year, resolve_year_source, ) +from pipelines.datasets.br_ms_sim.utils import container_memory_limit_gb from pipelines.utils.metadata.domain import ( AllFree, CoverageSpec, @@ -63,6 +64,8 @@ def run_ms_sim( prefix="Dump: ", dataset_id=dataset_id, table_id=table_id ) + print(f"Limite de memória do container: {container_memory_limit_gb()} GiB") + backfill = ano is not None source_max_year = get_source_max_year() ano = int(ano if backfill else source_max_year) @@ -175,8 +178,14 @@ def table_flow( source_format=source_format, ) + # Os dois formatos: o template do work pool descarta em silêncio a chave + # que não reconhece, e o pod cai no limite padrão. # pyrefly: ignore [missing-attribute] - table_flow.job_variables = {"memory": "8Gi"} + table_flow.job_variables = { + "memory": "8Gi", + "memory_limit": "8Gi", + "memory_request": "2Gi", + } return table_flow diff --git a/pipelines/datasets/br_ms_sim/utils.py b/pipelines/datasets/br_ms_sim/utils.py index f45e02663c..7846cdd467 100644 --- a/pipelines/datasets/br_ms_sim/utils.py +++ b/pipelines/datasets/br_ms_sim/utils.py @@ -21,6 +21,37 @@ FINAL, PRELIM = "definitivo", "preliminar" +def container_memory_limit_gb() -> float | None: + """Lê o limite de memória deste container, em GiB. + + Chave de `job_variables` que o template do work pool não reconhece é + descartada em silêncio, e o pod fica no limite padrão em vez do pedido. + + Returns: + O limite em GiB, ou None se for ilimitado ou ilegível. + """ + for path in ( + "/sys/fs/cgroup/memory.max", + "/sys/fs/cgroup/memory/memory.limit_in_bytes", + ): + try: + with open(path) as file: + raw = file.read().strip() + except OSError: + continue + if raw == "max": + return None + try: + value = int(raw) + except ValueError: + continue + # cgroup v1 usa um número enorme para "sem limite". + if value >= 1 << 62: + return None + return round(value / 1024**3, 2) + return None + + def build_paths(table_id: str, ano: int) -> tuple[Path, Path]: """Cria os diretórios de trabalho do ano e devolve os dois caminhos. From 0ef2b27dd2eb9410fa050d3ce7695f3fd1c01b43 Mon Sep 17 00:00:00 2001 From: Davi Cavalcante <93160711+DaviMacielCavalcante@users.noreply.github.com> Date: Thu, 3 Sep 2026 12:29:31 -0300 Subject: [PATCH 4/9] =?UTF-8?q?fix:=20trocando=20por=20cliente=20com=20aut?= =?UTF-8?q?entica=C3=A7=C3=A3o?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pipelines/utils/tasks.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/pipelines/utils/tasks.py b/pipelines/utils/tasks.py index 5772f1b8c3..03479104c4 100644 --- a/pipelines/utils/tasks.py +++ b/pipelines/utils/tasks.py @@ -76,7 +76,6 @@ def _sync_staging_schema( tb: bd.Table, data_path: str | Path, source_format: str, - billing_project_id: str, ) -> None: """Adiciona ao schema da staging as colunas que a fonte passou a trazer. @@ -103,14 +102,15 @@ def _sync_staging_schema( tb: tabela `basedosdados` já instanciada, apontando para a staging. data_path: arquivo ou diretório com os dados que serão carregados. source_format: `"csv"` ou `"parquet"`. - billing_project_id: projeto GCP usado para faturar a chamada. """ header_path = dump_header(data_path=data_path, source_format=source_format) incoming = tb._load_staging_schema_from_data( data_sample_path=header_path, source_format=source_format ) - client = bigquery.Client(project=billing_project_id) + # O cliente da lib é quem criou a tabela externa e escreve no prefixo. Abrir + # um `bigquery.Client` aqui cairia no ADC do pod, sem permissão de update. + client = tb.client["bigquery_staging"] table = client.get_table(tb.table_full_name["staging"]) current = {_bq_safe_column_name(field.name) for field in table.schema} @@ -190,7 +190,6 @@ def _upload_to_gcs( tb=tb, data_path=data_path, source_format=source_format, - billing_project_id=billing_project_id, ) elif dump_mode == "overwrite": From 61d4bc75a56c4722cde9b09868856602a2df9b67 Mon Sep 17 00:00:00 2001 From: Davi Cavalcante <93160711+DaviMacielCavalcante@users.noreply.github.com> Date: Thu, 3 Sep 2026 14:14:43 -0300 Subject: [PATCH 5/9] =?UTF-8?q?chore:=20removendo=20fun=C3=A7=C3=B5es=20e?= =?UTF-8?q?=20logs=20de=20depura=C3=A7=C3=A3o?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pipelines/datasets/br_ms_sim/flows.py | 3 --- pipelines/datasets/br_ms_sim/utils.py | 33 --------------------------- 2 files changed, 36 deletions(-) diff --git a/pipelines/datasets/br_ms_sim/flows.py b/pipelines/datasets/br_ms_sim/flows.py index 8fcfe3cf83..785a55ac79 100644 --- a/pipelines/datasets/br_ms_sim/flows.py +++ b/pipelines/datasets/br_ms_sim/flows.py @@ -11,7 +11,6 @@ get_source_max_year, resolve_year_source, ) -from pipelines.datasets.br_ms_sim.utils import container_memory_limit_gb from pipelines.utils.metadata.domain import ( AllFree, CoverageSpec, @@ -64,8 +63,6 @@ def run_ms_sim( prefix="Dump: ", dataset_id=dataset_id, table_id=table_id ) - print(f"Limite de memória do container: {container_memory_limit_gb()} GiB") - backfill = ano is not None source_max_year = get_source_max_year() ano = int(ano if backfill else source_max_year) diff --git a/pipelines/datasets/br_ms_sim/utils.py b/pipelines/datasets/br_ms_sim/utils.py index 7846cdd467..1d5d91105b 100644 --- a/pipelines/datasets/br_ms_sim/utils.py +++ b/pipelines/datasets/br_ms_sim/utils.py @@ -21,37 +21,6 @@ FINAL, PRELIM = "definitivo", "preliminar" -def container_memory_limit_gb() -> float | None: - """Lê o limite de memória deste container, em GiB. - - Chave de `job_variables` que o template do work pool não reconhece é - descartada em silêncio, e o pod fica no limite padrão em vez do pedido. - - Returns: - O limite em GiB, ou None se for ilimitado ou ilegível. - """ - for path in ( - "/sys/fs/cgroup/memory.max", - "/sys/fs/cgroup/memory/memory.limit_in_bytes", - ): - try: - with open(path) as file: - raw = file.read().strip() - except OSError: - continue - if raw == "max": - return None - try: - value = int(raw) - except ValueError: - continue - # cgroup v1 usa um número enorme para "sem limite". - if value >= 1 << 62: - return None - return round(value / 1024**3, 2) - return None - - def build_paths(table_id: str, ano: int) -> tuple[Path, Path]: """Cria os diretórios de trabalho do ano e devolve os dois caminhos. @@ -176,7 +145,6 @@ def download_year(ano: int, source: str, input_dir: Path) -> Path: destination = input_dir / f"DO{sigla_uf}{ano}.dbc" try: urllib.request.urlretrieve(url, filename=str(destination)) - print(f" {destination.name}: {destination.stat().st_size} bytes") except Exception as error: destination.unlink(missing_ok=True) missing.append(sigla_uf) @@ -522,7 +490,6 @@ def clean_year( partition / f"{table_id}.csv", index=False ) total += len(dataframe) - print(f" {sigla_uf}: {len(dataframe):,} linhas") if total == 0: raise RuntimeError( From 39211af907c71b89522c44176c9e60db84c6df3e Mon Sep 17 00:00:00 2001 From: Davi Cavalcante <93160711+DaviMacielCavalcante@users.noreply.github.com> Date: Thu, 3 Sep 2026 17:59:16 -0300 Subject: [PATCH 6/9] =?UTF-8?q?fix:=20adicionando=20uma=20coluna=20nova=20?= =?UTF-8?q?e=20impedindo=20do=20c=C3=B3digo=20continuar=20deletando=20uma?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- models/br_ms_sim/br_ms_sim__microdados.sql | 2 +- pipelines/datasets/br_ms_sim/README.md | 23 +++++++++++++++++----- pipelines/datasets/br_ms_sim/constants.py | 3 ++- pipelines/datasets/br_ms_sim/utils.py | 13 +++++++++--- 4 files changed, 31 insertions(+), 10 deletions(-) diff --git a/models/br_ms_sim/br_ms_sim__microdados.sql b/models/br_ms_sim/br_ms_sim__microdados.sql index ed444e882b..1210656fbb 100644 --- a/models/br_ms_sim/br_ms_sim__microdados.sql +++ b/models/br_ms_sim/br_ms_sim__microdados.sql @@ -104,5 +104,5 @@ select safe_cast(numero_dias_informacao as int64) numero_dias_informacao, safe_cast(fontes_informacao as string) fontes_informacao, safe_cast(alt_causa as string) alt_causa, - safe_cast(dado_preliminar as string) dado_preliminar + safe_cast(coalesce(dado_preliminar, '0') as string) dado_preliminar from {{ set_datalake_project("br_ms_sim_staging.microdados") }} as t diff --git a/pipelines/datasets/br_ms_sim/README.md b/pipelines/datasets/br_ms_sim/README.md index 12cc05b8b8..d3489c705d 100644 --- a/pipelines/datasets/br_ms_sim/README.md +++ b/pipelines/datasets/br_ms_sim/README.md @@ -64,13 +64,26 @@ utils.download_table("microdados", 2025, source) utils.clean_table("microdados", 2025, source) ``` +## Série histórica + +Os CSVs de 1996–2024 na staging foram gravados pela limpeza anterior, que +nulificava o código `6` de `local_ocorrencia` (aldeia indígena) antes do upload. +O dado não chegou à staging, então `full-refresh` não o recupera — só +reprocessar os `.dbc`. É o que faz `reprocess_br_ms_sim.py`, na raiz do repo. + +Enquanto o histórico não é reprocessado, ele também não traz `dado_preliminar`, +e a coluna volta nula nesses anos. O modelo aplica `coalesce(dado_preliminar, +'0')`: tudo anterior à primeira execução do pipeline veio do diretório +definitivo. Depois do reprocessamento o `coalesce` deixa de encontrar nulos e +continua correto. + ## Pontos de atenção -- `dado_preliminar` é coluna nova no modelo. Em `dump_mode="append"` o schema da - tabela externa só é ampliado por `_sync_staging_schema` - (`pipelines/utils/tasks.py`), que abre o cliente do BigQuery sem credencial e - responde `403 bigquery.tables.update`. É o primeiro caso em que esse caminho é - de fato exercitado. +- `dado_preliminar` é coluna nova no modelo, e em `dump_mode="append"` quem + amplia o schema da tabela externa é o `_sync_staging_schema` + (`pipelines/utils/tasks.py`). Esse caminho nunca tinha sido exercitado e + falhava com `403 bigquery.tables.update`, porque abria o cliente do BigQuery + sem credencial; corrigido nesta mesma branch para usar o cliente da lib. - O dicionário do conjunto vem de uma staging própria, alimentada por `models/br_ms_sim/code/update_dicionario.py`. As linhas de `dado_preliminar` precisam ser acrescentadas por lá. diff --git a/pipelines/datasets/br_ms_sim/constants.py b/pipelines/datasets/br_ms_sim/constants.py index 2a2b939196..22f258fad6 100644 --- a/pipelines/datasets/br_ms_sim/constants.py +++ b/pipelines/datasets/br_ms_sim/constants.py @@ -272,7 +272,8 @@ class constants(Enum): # Códigos que representam ausência de informação e viram NULL. NULLIFY = { - "local_ocorrencia": ["0", "6", "7", "9"], + # "6" é aldeia indígena, não ausência de informação. + "local_ocorrencia": ["0", "7", "9"], "sexo": ["0", "6", "7", "9"], "raca_cor": ["0", "6", "7", "9"], "estado_civil": ["0", "9"], diff --git a/pipelines/datasets/br_ms_sim/utils.py b/pipelines/datasets/br_ms_sim/utils.py index 1d5d91105b..b0f8053c28 100644 --- a/pipelines/datasets/br_ms_sim/utils.py +++ b/pipelines/datasets/br_ms_sim/utils.py @@ -144,7 +144,13 @@ def download_year(ano: int, source: str, input_dir: Path) -> Path: url = template.format(sigla_uf=sigla_uf, ano=ano) destination = input_dir / f"DO{sigla_uf}{ano}.dbc" try: - urllib.request.urlretrieve(url, filename=str(destination)) + # `urlretrieve` não aceita timeout: uma transferência travada + # penduraria o worker até alguém matar a execução. + with ( + urllib.request.urlopen(url, timeout=300) as response, + open(destination, "wb") as file, + ): + shutil.copyfileobj(response, file) except Exception as error: destination.unlink(missing_ok=True) missing.append(sigla_uf) @@ -302,8 +308,9 @@ def parse_hora(value: object) -> str | None: def parse_idade(value: object) -> float | None: """Converte a idade codificada do SIM em anos. - O primeiro dígito é a unidade — 1 minuto, 2 hora, 3 mês, 4 ano, 5 ano acima - de 100 — e os demais são a quantidade. + O primeiro dígito é a unidade — 0 minutos, 1 horas, 2 dias, 3 meses, 4 anos, + 5 anos acima de 100 — e os demais são a quantidade. Minutos e horas viram + zero; a unidade 0 não é tratada e devolve None, como na carga anterior. Args: value: Valor bruto do arquivo. From 2911b4d13ddb314b6fbdc0f4f1d2f50aae174243 Mon Sep 17 00:00:00 2001 From: Davi Cavalcante <93160711+DaviMacielCavalcante@users.noreply.github.com> Date: Fri, 11 Sep 2026 11:09:34 -0300 Subject: [PATCH 7/9] =?UTF-8?q?chore:=20deixando=20o=20flow=20mais=20no=20?= =?UTF-8?q?padr=C3=A3o=20do=20repo?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pipelines/datasets/br_ms_sim/README.md | 10 ++ pipelines/datasets/br_ms_sim/constants.py | 38 +++++--- pipelines/datasets/br_ms_sim/flows.py | 113 ++++++---------------- pipelines/datasets/br_ms_sim/tasks.py | 7 +- pipelines/datasets/br_ms_sim/utils.py | 55 +++++------ 5 files changed, 95 insertions(+), 128 deletions(-) diff --git a/pipelines/datasets/br_ms_sim/README.md b/pipelines/datasets/br_ms_sim/README.md index d3489c705d..e50707d7f7 100644 --- a/pipelines/datasets/br_ms_sim/README.md +++ b/pipelines/datasets/br_ms_sim/README.md @@ -42,6 +42,16 @@ Execução de teste no pool de dev, sem tocar em produção: Os padrões escrevem em **produção**, mesmo saindo do pool de teste. +## Qual ano entra + +Sem `ano`, o flow carrega o ano mais recente que existe na fonte, que é o ano em +curso: o diretório preliminar publica o ano corrente antes de ele fechar. + +O poll compara esse ano com o fim da cobertura da tabela. Depois que a cobertura +alcança o ano, as execuções seguintes encerram sem carregar nada, e as revisões +que o DATASUS publicar naquele ano preliminar não entram. Para trazê-las, executar +com `ano` preenchido ou com `force_run`. + ## Formato da staging A staging é CSV desde a carga original. O particionado sai em diff --git a/pipelines/datasets/br_ms_sim/constants.py b/pipelines/datasets/br_ms_sim/constants.py index 22f258fad6..31cd54b23e 100644 --- a/pipelines/datasets/br_ms_sim/constants.py +++ b/pipelines/datasets/br_ms_sim/constants.py @@ -8,20 +8,30 @@ class constants(Enum): """Constantes de br_ms_sim.""" - FTP_FINAL = ( - "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/" - "DO{sigla_uf}{ano}.dbc" - ) - FTP_PRELIM = ( - "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/PRELIM/DORES/" - "DO{sigla_uf}{ano}.dbc" - ) - FTP_FINAL_DIR = ( - "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/" - ) - FTP_PRELIM_DIR = ( - "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/PRELIM/DORES/" - ) + # Versões do dado no FTP, em ordem de precedência: um ano fechado pelo + # DATASUS passa a existir nas duas, e o definitivo é o que vale. + SOURCES = { + "definitivo": { + "file": ( + "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/" + "DO{sigla_uf}{ano}.dbc" + ), + "dir": ( + "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/" + ), + }, + "preliminar": { + "file": ( + "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/PRELIM/DORES/" + "DO{sigla_uf}{ano}.dbc" + ), + "dir": ( + "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/PRELIM/DORES/" + ), + }, + } + + PRELIM = "preliminar" # Área de trabalho do pod. `input/` recebe os .dbc, `output/` o particionado # que sobe para o GCS. diff --git a/pipelines/datasets/br_ms_sim/flows.py b/pipelines/datasets/br_ms_sim/flows.py index 785a55ac79..f5bd6e0cb4 100644 --- a/pipelines/datasets/br_ms_sim/flows.py +++ b/pipelines/datasets/br_ms_sim/flows.py @@ -4,7 +4,6 @@ from prefect import flow -from pipelines.datasets.br_ms_sim.constants import constants from pipelines.datasets.br_ms_sim.tasks import ( clean_table, download_table, @@ -13,7 +12,6 @@ ) from pipelines.utils.metadata.domain import ( AllFree, - CoverageSpec, DateFormat, YearOnly, ) @@ -28,36 +26,21 @@ upload_to_gcs, ) -DATE_FORMAT = DateFormat.YEAR - -def coverage(table_id: str) -> CoverageSpec: - """Devolve a cobertura da tabela. - - Raises: - ValueError: Se a tabela não constar de `constants.TABLES`. - """ - if table_id not in constants.TABLES.value: - raise ValueError(f"tabela sem cobertura definida: {table_id}") - return AllFree( - date_column=YearOnly(col="ano"), - date_format=DATE_FORMAT, - ) - - -def run_ms_sim( - *, - dataset_id: str, - table_id: str, - ano: int | None, - materialize_after_dump: bool, - update_metadata: bool, - target: str, - force_run: bool, - dump_mode: str, - source_format: str, +@flow( + name="br_ms_sim__microdados", + log_prints=True, +) +def br_ms_sim__microdados( + dataset_id: str = "br_ms_sim", + table_id: str = "microdados", + ano: int | None = None, + materialize_after_dump: bool = True, + update_metadata: bool = True, + target: str = "prod", + force_run: bool = False, ) -> None: - """Executa o ciclo baixar, limpar, subir, dbt e metadados de um ano.""" + """Carrega um ano do SIM, do FTP do DATASUS até a materialização.""" # pyrefly: ignore [unused-coroutine] rename_flow_run_dataset_table( prefix="Dump: ", dataset_id=dataset_id, table_id=table_id @@ -73,7 +56,7 @@ def run_ms_sim( table_id=table_id, source_max_date=source_max_year, env="prod", - date_format=DATE_FORMAT, + date_format=DateFormat.YEAR, compare_against="coverage", ) if not has_new_data: @@ -86,7 +69,7 @@ def run_ms_sim( table_id=table_id, source_max_date=source_max_year, env="prod", - date_format=DATE_FORMAT, + date_format=DateFormat.YEAR, update_metadata=update_metadata, materialize_after_dump=materialize_after_dump, ) @@ -102,8 +85,8 @@ def run_ms_sim( dataset_id=dataset_id, table_id=table_id, bucket_name="basedosdados-dev", - dump_mode=dump_mode, - source_format=source_format, + dump_mode="append", + source_format="csv", ) run_dbt( @@ -121,8 +104,8 @@ def run_ms_sim( dataset_id=dataset_id, table_id=table_id, bucket_name="basedosdados", - dump_mode=dump_mode, - source_format=source_format, + dump_mode="append", + source_format="csv", ) run_dbt( @@ -136,54 +119,20 @@ def run_ms_sim( register_table_materialization_task( dataset_id=dataset_id, table_id=table_id, - coverage=coverage(table_id), + coverage=AllFree( + date_column=YearOnly(col="ano"), + date_format=DateFormat.YEAR, + ), env="prod", bq_project="basedosdados", ) -def ms_sim_flow( - table_id: str, - dump_mode: str = "append", - source_format: str = "csv", -): - """Carimba o flow de uma tabela.""" - - @flow( - name=f"br_ms_sim__{table_id}", - log_prints=True, - ) - def table_flow( - dataset_id: str = "br_ms_sim", - table_id: str = table_id, - ano: int | None = None, - materialize_after_dump: bool = True, - update_metadata: bool = True, - target: str = "prod", - force_run: bool = False, - ) -> None: - """Carrega um ano do SIM, do FTP do DATASUS até a materialização.""" - run_ms_sim( - dataset_id=dataset_id, - table_id=table_id, - ano=ano, - materialize_after_dump=materialize_after_dump, - update_metadata=update_metadata, - target=target, - force_run=force_run, - dump_mode=dump_mode, - source_format=source_format, - ) - - # Os dois formatos: o template do work pool descarta em silêncio a chave - # que não reconhece, e o pod cai no limite padrão. - # pyrefly: ignore [missing-attribute] - table_flow.job_variables = { - "memory": "8Gi", - "memory_limit": "8Gi", - "memory_request": "2Gi", - } - return table_flow - - -br_ms_sim__microdados = ms_sim_flow("microdados") +# Os dois formatos: o template do work pool descarta em silêncio a chave +# que não reconhece, e o pod cai no limite padrão. +# pyrefly: ignore [missing-attribute] +br_ms_sim__microdados.job_variables = { + "memory": "8Gi", + "memory_limit": "8Gi", + "memory_request": "2Gi", +} diff --git a/pipelines/datasets/br_ms_sim/tasks.py b/pipelines/datasets/br_ms_sim/tasks.py index b1987afa3f..a99bfb46b2 100644 --- a/pipelines/datasets/br_ms_sim/tasks.py +++ b/pipelines/datasets/br_ms_sim/tasks.py @@ -29,7 +29,8 @@ def resolve_year_source(ano: int) -> str: ano: Ano a resolver. Returns: - `"definitivo"` ou `"preliminar"`. + A chave da versão em `constants.SOURCES` — `"definitivo"` ou + `"preliminar"`. """ return utils.resolve_year_source(ano) @@ -41,7 +42,7 @@ def download_table(table_id: str, ano: int, source: str) -> Path: Args: table_id: Slug da tabela. ano: Ano a baixar. - source: `"definitivo"` ou `"preliminar"`. + source: Versão do dado, chave de `constants.SOURCES`. Returns: O diretório de entrada com os arquivos baixados. @@ -56,7 +57,7 @@ def clean_table(table_id: str, ano: int, source: str) -> Path: Args: table_id: Slug da tabela. ano: Ano a limpar. - source: `"definitivo"` ou `"preliminar"`. + source: Versão do dado, chave de `constants.SOURCES`. Returns: O diretório particionado, no formato esperado por `upload_to_gcs`. diff --git a/pipelines/datasets/br_ms_sim/utils.py b/pipelines/datasets/br_ms_sim/utils.py index b0f8053c28..ea16fede47 100644 --- a/pipelines/datasets/br_ms_sim/utils.py +++ b/pipelines/datasets/br_ms_sim/utils.py @@ -18,8 +18,6 @@ from pipelines.datasets.br_ms_sim.constants import constants -FINAL, PRELIM = "definitivo", "preliminar" - def build_paths(table_id: str, ano: int) -> tuple[Path, Path]: """Cria os diretórios de trabalho do ano e devolve os dois caminhos. @@ -71,18 +69,19 @@ def list_ftp_years(directory_url: str) -> set[int]: def get_source_max_year() -> str: """Devolve o ano mais recente publicado na fonte. - Considera os dois diretórios, já que o ano corrente costuma existir apenas - no preliminar. O valor é a competência do dado, não a data da consulta. + Considera todas as versões de `constants.SOURCES`, já que o ano corrente + costuma existir apenas no preliminar. O valor é a competência do dado, não a + data da consulta. Returns: O ano mais recente, no formato `%Y`. Raises: - RuntimeError: Se nenhum arquivo for encontrado nos dois diretórios. + RuntimeError: Se nenhum arquivo for encontrado em nenhuma das versões. """ - years = list_ftp_years(constants.FTP_FINAL_DIR.value) | list_ftp_years( - constants.FTP_PRELIM_DIR.value - ) + years: set[int] = set() + for urls in constants.SOURCES.value.values(): + years |= list_ftp_years(urls["dir"]) if not years: raise RuntimeError( "nenhum arquivo DO*.dbc encontrado no FTP do DATASUS — a fonte " @@ -94,23 +93,24 @@ def get_source_max_year() -> str: def resolve_year_source(ano: int) -> str: """Diz de qual diretório o ano deve ser baixado. - O definitivo tem precedência sobre o preliminar. Um ano fechado pelo DATASUS - passa a existir nos dois diretórios, e reprocessá-lo substitui o dado - preliminar pelo definitivo. + Devolve a primeira versão de `constants.SOURCES` que tem o ano, e a ordem de + declaração lá é que dá ao definitivo precedência sobre o preliminar. Um ano + fechado pelo DATASUS passa a existir nos dois diretórios, e reprocessá-lo + substitui o dado preliminar pelo definitivo. Args: ano: Ano a resolver. Returns: - `"definitivo"` ou `"preliminar"`. + A chave da versão em `constants.SOURCES` — `"definitivo"` ou + `"preliminar"`. Raises: - FileNotFoundError: Se o ano não existir em nenhum dos dois diretórios. + FileNotFoundError: Se o ano não existir em nenhuma das versões. """ - if ano in list_ftp_years(constants.FTP_FINAL_DIR.value): - return FINAL - if ano in list_ftp_years(constants.FTP_PRELIM_DIR.value): - return PRELIM + for source, urls in constants.SOURCES.value.items(): + if ano in list_ftp_years(urls["dir"]): + return source raise FileNotFoundError( f"ano {ano} não está no FTP do DATASUS, nem definitivo nem preliminar" ) @@ -124,20 +124,17 @@ def download_year(ano: int, source: str, input_dir: Path) -> Path: Args: ano: Ano a baixar. - source: `"definitivo"` ou `"preliminar"`. + source: Versão do dado, chave de `constants.SOURCES`. input_dir: Diretório de destino. Returns: O diretório de destino. Raises: + KeyError: Se `source` não for uma chave de `constants.SOURCES`. RuntimeError: Se nenhuma das 27 UFs for baixada. """ - template = ( - constants.FTP_FINAL.value - if source == FINAL - else constants.FTP_PRELIM.value - ) + template = constants.SOURCES.value[source]["file"] missing = [] for sigla_uf in constants.UFS.value: @@ -309,8 +306,8 @@ def parse_idade(value: object) -> float | None: """Converte a idade codificada do SIM em anos. O primeiro dígito é a unidade — 0 minutos, 1 horas, 2 dias, 3 meses, 4 anos, - 5 anos acima de 100 — e os demais são a quantidade. Minutos e horas viram - zero; a unidade 0 não é tratada e devolve None, como na carga anterior. + 5 anos acima de 100 — e os demais são a quantidade. Horas viram zero; a + unidade 0 não é tratada e devolve None, como na carga anterior. Args: value: Valor bruto do arquivo. @@ -468,7 +465,7 @@ def clean_year( Args: table_id: Slug da tabela, que define as colunas de partição. ano: Ano processado. - source: `"definitivo"` ou `"preliminar"`. + source: Versão do dado, chave de `constants.SOURCES`. input_dir: Diretório com os arquivos `.dbc`. output_dir: Raiz do particionado. @@ -482,7 +479,7 @@ def clean_year( partition_columns = table["partition_columns"] file_prefix = table["file_prefix"] municipios = load_municipios() - is_prelim = source == PRELIM + is_prelim = source == constants.PRELIM.value total = 0 for filepath in sorted(input_dir.glob(f"{file_prefix}*{ano}.dbc")): @@ -513,7 +510,7 @@ def download_table(table_id: str, ano: int, source: str) -> Path: Args: table_id: Slug da tabela. ano: Ano a baixar. - source: `"definitivo"` ou `"preliminar"`. + source: Versão do dado, chave de `constants.SOURCES`. Returns: O diretório de entrada com os arquivos baixados. @@ -528,7 +525,7 @@ def clean_table(table_id: str, ano: int, source: str) -> Path: Args: table_id: Slug da tabela. ano: Ano a limpar. - source: `"definitivo"` ou `"preliminar"`. + source: Versão do dado, chave de `constants.SOURCES`. Returns: O diretório particionado, no formato esperado por `upload_to_gcs`. From 9ec90eab27b9df220d7125c2632985b7ff9577b0 Mon Sep 17 00:00:00 2001 From: Davi Cavalcante <93160711+DaviMacielCavalcante@users.noreply.github.com> Date: Fri, 11 Sep 2026 12:08:01 -0300 Subject: [PATCH 8/9] =?UTF-8?q?fix(br=5Fms=5Fsim):=20descri=C3=A7=C3=A3o?= =?UTF-8?q?=20de=20dado=5Fpreliminar=20sem=20quebra=20final?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- models/br_ms_sim/schema.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/models/br_ms_sim/schema.yml b/models/br_ms_sim/schema.yml index 7f49c8fd19..dfd995c1c7 100644 --- a/models/br_ms_sim/schema.yml +++ b/models/br_ms_sim/schema.yml @@ -214,7 +214,7 @@ models: - name: alt_causa description: Alt. Causa - name: dado_preliminar - description: > + description: >- Indica se o registro vem da versão preliminar da fonte, ainda sujeita a revisão (1), ou da versão definitiva (0) - name: br_ms_sim__dicionario From d88406c4f2438ba7e7c0d2b37c5613d10543b91a Mon Sep 17 00:00:00 2001 From: Davi Cavalcante <93160711+DaviMacielCavalcante@users.noreply.github.com> Date: Fri, 11 Sep 2026 16:00:10 -0300 Subject: [PATCH 9/9] fix(br_ms_sim): remove chave memory, que o work pool descarta --- pipelines/datasets/br_ms_sim/flows.py | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/pipelines/datasets/br_ms_sim/flows.py b/pipelines/datasets/br_ms_sim/flows.py index f5bd6e0cb4..4570928da8 100644 --- a/pipelines/datasets/br_ms_sim/flows.py +++ b/pipelines/datasets/br_ms_sim/flows.py @@ -128,11 +128,10 @@ def br_ms_sim__microdados( ) -# Os dois formatos: o template do work pool descarta em silêncio a chave -# que não reconhece, e o pod cai no limite padrão. +# `memory` não existe no template do work pool, que só conhece o par abaixo, e +# chave fora do template é descartada em silêncio — o pod ficaria no padrão. # pyrefly: ignore [missing-attribute] br_ms_sim__microdados.job_variables = { - "memory": "8Gi", "memory_limit": "8Gi", "memory_request": "2Gi", }