diff --git a/models/br_ms_sim/br_ms_sim__microdados.sql b/models/br_ms_sim/br_ms_sim__microdados.sql index 82c0698df6..1210656fbb 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(coalesce(dado_preliminar, '0') 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..dfd995c1c7 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..e50707d7f7 --- /dev/null +++ b/pipelines/datasets/br_ms_sim/README.md @@ -0,0 +1,102 @@ +# 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. + +## 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 +`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) +``` + +## 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, 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á. +- 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..31cd54b23e --- /dev/null +++ b/pipelines/datasets/br_ms_sim/constants.py @@ -0,0 +1,493 @@ +""" +Constantes de br_ms_sim. +""" + +from enum import Enum + + +class constants(Enum): + """Constantes de br_ms_sim.""" + + # 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. + 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 = { + # "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"], + "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..4570928da8 --- /dev/null +++ b/pipelines/datasets/br_ms_sim/flows.py @@ -0,0 +1,137 @@ +""" +Flows de br_ms_sim — Prefect 3. +""" + +from prefect import flow + +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, + 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, +) + + +@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: + """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 + ) + + 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=DateFormat.YEAR, + 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=DateFormat.YEAR, + 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="append", + source_format="csv", + ) + + 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="append", + source_format="csv", + ) + + 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=AllFree( + date_column=YearOnly(col="ano"), + date_format=DateFormat.YEAR, + ), + env="prod", + bq_project="basedosdados", + ) + + +# `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_limit": "8Gi", + "memory_request": "2Gi", +} diff --git a/pipelines/datasets/br_ms_sim/tasks.py b/pipelines/datasets/br_ms_sim/tasks.py new file mode 100644 index 0000000000..a99bfb46b2 --- /dev/null +++ b/pipelines/datasets/br_ms_sim/tasks.py @@ -0,0 +1,65 @@ +""" +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: + A chave da versão em `constants.SOURCES` — `"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: Versão do dado, chave de `constants.SOURCES`. + + 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: Versão do dado, chave de `constants.SOURCES`. + + 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..ea16fede47 --- /dev/null +++ b/pipelines/datasets/br_ms_sim/utils.py @@ -0,0 +1,540 @@ +""" +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 + + +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 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 em nenhuma das versões. + """ + 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 " + "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. + + 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: + A chave da versão em `constants.SOURCES` — `"definitivo"` ou + `"preliminar"`. + + Raises: + FileNotFoundError: Se o ano não existir em nenhuma das versões. + """ + 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" + ) + + +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: 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.SOURCES.value[source]["file"] + + 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: + # `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) + 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) + # `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) + + +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 — 0 minutos, 1 horas, 2 dias, 3 meses, 4 anos, + 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. + + 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: Versão do dado, chave de `constants.SOURCES`. + 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 == constants.PRELIM.value + 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) + + 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: Versão do dado, chave de `constants.SOURCES`. + + 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: Versão do dado, chave de `constants.SOURCES`. + + 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, + ) 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":