Objetivo do curso: Capacitar participantes a extrair, validar, tratar e analisar microdados do SUS com R, integrando fontes de custos/preços e técnicas aplicadas à avaliação econômica.

Ecossistema de dados do SUS com R

Ementa

  • Como estimar o uso de memória dos objetos no R?
  • fread() ainda é útil?
  • Um tibble com 100 milhões de registros cabe na RAM?
  • Como processar vários DBC sem manter todos simultaneamente na memória?
  • Como converter uma coleção mensal para Parquet?
  • Como consultar arquivos com Arrow ou DuckDB sem coletar tudo?
  • Como visualizar amostras no RStudio?

Objetivo da aula: adaptar a ETL da Parte 1 para volumes que não devem ser mantidos integralmente na memória do R.


Preparar os pacotes

Execute este bloco uma vez para instalar os pacotes usados nesta parte do módulo.

install.packages(c(
  "arrow", "data.table", "DBI", "dplyr",
  "duckdb", "lobstr", "remotes"
))

# O read.dbc pode não estar disponível no CRAN para sua versão do R
remotes::install_github("danicat/read.dbc", upgrade = "never")

Memória e volume

O que ocupa memória

Quando read.dbc() lê um arquivo, o resultado é materializado como um data.frame. O tamanho depende de:

  • quantidade de linhas;
  • quantidade de colunas;
  • tipos das colunas;
  • comprimento e repetição dos textos;
  • objetos intermediários produzidos pelas transformações;
  • cópias temporárias criadas durante seleções, junções e ordenações.

O número de linhas isoladamente não determina o consumo. Cem milhões de linhas com três inteiros são diferentes de cem milhões de linhas com dezenas de campos de texto.

Medir um objeto

library(lobstr)

# pa foi produzido na Parte 1.
obj_size(pa)

Com R base:

# Exibe uma estimativa do tamanho do objeto.
format(
  object.size(pa),
  units = "GB"
)

Observar a memória da sessão

# Mostra o uso do gerenciador de memória.
gc()

Para remover objetos que não serão mais utilizados:

# Remove partes intermediárias e solicita coleta de memória.
rm(partes_pa, partes_am, pa_exemplo)
gc()

gc() não torna um objeto grande menor. Ele apenas permite recuperar memória de objetos que deixaram de ser referenciados.


O tidyverse e 100 milhões de registros

Resposta curta

Um tibble local é um objeto em memória. Portanto, dplyr não garante que 100 milhões de linhas caibam na RAM. O resultado depende do número e dos tipos das colunas, da memória disponível e das cópias temporárias exigidas pela operação.

O tidyverse continua útil em grandes volumes quando a sintaxe dplyr é executada sobre um backend preguiçoso, como:

  • Arrow;
  • DuckDB;
  • PostgreSQL por meio de dbplyr;
  • outros bancos compatíveis com DBI e dbplyr.

Nesses casos, filter(), select(), group_by() e summarise() são traduzidos ou executados pelo backend. Somente o resultado necessário deve ser trazido para a RAM com collect().

Comparação das estratégias

Estratégia Dados principais ficam na RAM? Uso recomendado
data.frame ou tibble local Sim recortes que cabem confortavelmente na memória
data.table local Sim leitura e transformação rápida de arquivos delimitados
Arrow Dataset Não integralmente coleção de arquivos Parquet e consultas por partição
DuckDB Não integralmente SQL analítico sobre arquivos ou banco local
PostgreSQL Não integralmente dados persistentes, multiusuário e rotinas institucionais

A interface dplyr não define onde os dados são armazenados. Um pipeline pode operar sobre um tibble em RAM, um conjunto Arrow ou uma tabela PostgreSQL.


fread() ainda é utilizado?

Sim. data.table::fread() continua sendo uma opção relevante para arquivos regulares e delimitados, como CSV e TSV. Entre seus recursos estão:

  • leitura multithread;
  • detecção automática de separador e tipos;
  • seleção de colunas durante a leitura;
  • definição de tipos com colClasses;
  • inspeção da estrutura com nrows = 0;
  • preservação de zeros à esquerda com keepLeadingZeros = TRUE.

fread() não lê diretamente o formato DBC. Ele se aplica depois que os dados foram convertidos para um arquivo delimitado ou quando a fonte já é CSV/TSV.

Inspecionar sem carregar todas as linhas

library(data.table)

# Lê apenas nomes e tipos inferidos.
estrutura <- fread(
  "pa_2024.csv",
  nrows = 0
)

estrutura

Ler somente as colunas necessárias

library(data.table)

pa_dt <- fread(
  "pa_2024.csv",
  select = c(
    "PA_CODUNI",
    "PA_CMP",
    "PA_PROC_ID",
    "PA_QTDAPR",
    "PA_VALAPR"
  ),
  colClasses = list(
    character = c(
      "PA_CODUNI",
      "PA_CMP",
      "PA_PROC_ID"
    )
  ),
  keepLeadingZeros = TRUE
)

Atualizar por referência

Uma diferença importante do data.table é a atualização por referência com :=, que reduz cópias desnecessárias.

library(data.table)

# Cria a métrica dentro do mesmo objeto.
pa_dt[
  PA_QTDAPR > 0,
  valor_por_unidade := PA_VALAPR / PA_QTDAPR
]

Processar os DBC um por vez

read.dbc() materializa um DBC inteiro. Para reduzir o pico de memória, leia um arquivo mensal, selecione as colunas necessárias, grave o resultado e libere o objeto antes do próximo arquivo.

Localizar os PA já baixados

arquivos_pa <- list.files(
  "dados_modulo_2",
  pattern = "^PA.*\\.dbc$",
  full.names = TRUE,
  ignore.case = TRUE
)

arquivos_pa
Mostrar resultado
## [1] "dados_modulo_2/PAAC2409.dbc" "dados_modulo_2/PAAC2410.dbc"
## [3] "dados_modulo_2/PAAC2411.dbc" "dados_modulo_2/PAAC2501.dbc"

Converter cada PA para Parquet

O formato Parquet armazena colunas de forma compacta e permite que mecanismos analíticos leiam apenas colunas e grupos de linhas necessários.

library(arrow)
library(dplyr)
library(read.dbc)

diretorio_parquet <- "dados_modulo_2/parquet_pa"

dir.create(
  diretorio_parquet,
  showWarnings = FALSE,
  recursive = TRUE
)

for (arquivo in arquivos_pa) {
  # Mantém somente as colunas da análise.
  pa_mes <- read.dbc(arquivo) |>
    transmute(
      PA_CODUNI = as.character(PA_CODUNI),
      PA_CMP = as.character(PA_CMP),
      PA_PROC_ID = as.character(PA_PROC_ID),
      PA_AUTORIZ = as.character(PA_AUTORIZ),
      PA_QTDAPR = as.numeric(PA_QTDAPR),
      PA_VALAPR = as.numeric(PA_VALAPR),
      arquivo_origem = basename(arquivo)
    )

  saida <- file.path(
    diretorio_parquet,
    sub("\\.dbc$", ".parquet", basename(arquivo), ignore.case = TRUE)
  )

  write_parquet(
    pa_mes,
    saida
  )

  # Libera o mês antes da próxima leitura.
  rm(pa_mes)
  gc()
}

Essa estratégia não elimina o custo de memória necessário para abrir um DBC, mas evita manter várias competências simultaneamente na RAM.


Consultar com Arrow

Abrir a coleção sem coletar

library(arrow)
library(dplyr)

pa_dataset <- open_dataset(
  "dados_modulo_2/parquet_pa",
  format = "parquet"
)

pa_dataset

open_dataset() cria uma representação da coleção de arquivos. Os registros não são convertidos imediatamente em um tibble local.

Filtrar, selecionar e resumir

library(arrow)
library(dplyr)

resumo_arrow <- pa_dataset |>
  filter(
    PA_CMP >= "202409",
    PA_CMP <= "202411"
  ) |>
  group_by(PA_CMP) |>
  summarise(
    quantidade = sum(PA_QTDAPR, na.rm = TRUE),
    valor = sum(PA_VALAPR, na.rm = TRUE)
  ) |>
  collect()

resumo_arrow

O collect() aparece somente depois da agregação. Assim, a RAM recebe três linhas resumidas, e não toda a coleção.

Obter uma amostra pequena

library(arrow)
library(dplyr)

amostra_pa <- pa_dataset |>
  select(
    PA_CMP,
    PA_CODUNI,
    PA_PROC_ID,
    PA_QTDAPR,
    PA_VALAPR
  ) |>
  head(1000) |>
  collect()

Consultar com DuckDB

DuckDB executa SQL analítico local e pode consultar arquivos Parquet sem importá-los integralmente para um data frame do R.

Criar uma conexão local

library(DBI)
library(duckdb)

con <- dbConnect(
  duckdb(),
  dbdir = "dados_modulo_2/sus.duckdb"
)

# Limita a memória usada pelo DuckDB.
dbExecute(
  con,
  "SET memory_limit = '4GB'"
)

Consultar Parquet com dplyr

library(DBI)
library(duckdb)
library(dplyr)

pa_lazy <- tbl(
  con,
  "read_parquet('dados_modulo_2/parquet_pa/*.parquet')"
)

resumo_duckdb <- pa_lazy |>
  filter(
    PA_CMP >= "202409",
    PA_CMP <= "202411"
  ) |>
  group_by(PA_CMP) |>
  summarise(
    quantidade = sum(PA_QTDAPR, na.rm = TRUE),
    valor = sum(PA_VALAPR, na.rm = TRUE)
  ) |>
  collect()

resumo_duckdb

Encerrar a conexão

library(DBI)
library(duckdb)

dbDisconnect(
  con,
  shutdown = TRUE
)

Visualizar no RStudio

O Data Viewer virtualiza a exibição das linhas, mas o objeto local continua ocupando memória. Além disso, ordenar ou filtrar um data frame muito grande pode exigir uma varredura completa.

Evite:

# Evite abrir o objeto integral.
View(pa)

Visualize uma amostra:

# Abre somente mil linhas já coletadas.
View(amostra_pa)

Para exploração inicial, prefira:

library(dplyr)

glimpse(amostra_pa)
head(amostra_pa, 20)
summary(amostra_pa)

Para dados externos à RAM, filtre e agregue no Arrow, DuckDB ou PostgreSQL; depois use collect() apenas no resultado destinado à inspeção ou ao gráfico.


Estratégia recomendada por escala

Situação Estratégia inicial
Um ou poucos DBC mensais read.dbc() + dplyr
Vários DBC que ainda cabem na RAM lista + bind_rows()
CSV grande fread(select = ..., colClasses = ...)
Muitos DBC mensais converter um por vez para Parquet
Coleção Parquet Arrow Dataset ou DuckDB
Rotina institucional persistente PostgreSQL, DuckDB persistente ou lakehouse
Resultado para gráfico ou tabela agregar antes de collect()

Regras práticas

  • selecione somente as colunas necessárias;
  • preserve identificadores como texto;
  • processe uma partição por vez;
  • remova objetos intermediários;
  • evite ordenar uma base inteira sem necessidade;
  • agregue antes de trazer os resultados para a RAM;
  • não use View() como método de validação de bases gigantes;
  • mantenha DBC como fonte bruta e Parquet ou banco como camada analítica;
  • registre arquivo, competência, layout e transformação aplicada.

Síntese

fread() continua adequado para arquivos delimitados, mas não substitui read.dbc() na leitura do DBC. Um tibble local com 100 milhões de linhas pode exceder a RAM, principalmente quando contém muitas colunas de texto ou quando as transformações criam cópias.

Para grandes volumes, mantenha a sintaxe familiar do dplyr, mas altere o backend:

DBC mensal
  → seleção de colunas
  → Parquet
  → Arrow ou DuckDB
  → filtro e agregação preguiçosos
  → collect() somente do resultado

Tabela de Diretrizes de Cuidado (PCDT)

Para realizar avaliações econômicas voltadas a linhas de cuidado específicas, precisamos saber quais códigos CID-10 buscar nas bases de produção. O Ministério da Saúde, por meio da Conitec, publica os Protocolos Clínicos e Diretrizes Terapêuticas (PCDT).

Utilizamos uma tabela de metadados padronizada para mapear a doença aos seus respectivos CIDs.

Consultando a tabela de diretrizes

O chunk abaixo lê o arquivo diretamente do repositório GitHub e exibe suas 20 primeiras linhas.

library(readr)
library(dplyr)
library(knitr)

# URL da tabela de metadados declarada como texto simples
url_tabela <- "https://raw.githubusercontent.com/LabXiabr/bd_geral/refs/heads/main/td_diretriz_cuidado.csv"

# Leitura do CSV diretamente da web
tabela_diretrizes <- read_csv(url_tabela, show_col_types = FALSE)

# Exibicao das 20 primeiras linhas
tabela_diretrizes |>
  head(20) |>
  kable(caption = "Primeiras 20 linhas da Tabela de Diretrizes de Cuidado")
Mostrar resultado
Primeiras 20 linhas da Tabela de Diretrizes de Cuidado
co_diretriz_cuidado…1 co_diretriz_cuidado…2 co_cid sg_pcdt no_abreviado no_diretriz_extenso no_diretriz50 sg_tipo st_rara
1 Acidente Vascular Cerebral Isquêmico Agudo I630,I631,I632,I633,I634,I635,I636,I638,I639,I650,I651,I652,I653,I658,I659,I660,I661,I662,I663,I664,I668,I669 AVC avc_isquem Trombólise no Acidente Vascular Cerebral (AVC) Isquêmico Agudo Acidente Vascular Cerebral Isquêmico Agudo pcdt 0
2 Acne grave L700,L701,L708 ACN acne Isotretinoína no tratamento da acne grave Acne grave Protocolos 0
3 Acromegalia E220 ACRO acromegali Acromegalia Acromegalia pcdt 1
4 Adenocarcinoma de Próstata C61,D75 PROS c_prostata Adenocarcinoma de Próstata Adenocarcinoma de Próstata ddt 0
5 Anemia Hemolítica Autoimune D590,D591,D460,D461,D467,D600,D608,D610,D611,D612,D613,D618,D70,Z948 AHA anemia_hem Anemia Hemolítica Autoimune Anemia Hemolítica Autoimune pcdt 1
6 Anemia por Deficiência de Ferro D500,D508 FER anemia_fer Anemia por Deficiência de Ferro Anemia por Deficiência de Ferro pcdt 0
7 Angioedema D841 ANGI angiodema Angioedema associado à deficiência de C1 esterase Angioedema pcdt 1
8 Aorta Torácica Descendente I712 ATD aorta_tora Diretrizes Brasileiras para Utilização de Endoprótese em Aorta Torácica Descendente Aorta Torácica Descendente Diretrizes 0
9 Artrite Idiopática Juvenil M080,M081,M082,M083,M084,M088,M089 AIJ artrite_ju Artrite Idiopática Juvenil (AIJ) Artrite Idiopática Juvenil pcdt 1
10 Artrite Psoríaca M070,M072,M073 APSO artrite_ps Artrite Psoríaca Artrite Psoríaca pcdt 0
11 Artrite Reativa M021,M023,M032,M036 AREA artrite_rea Artrite Reativa Artrite Reativa pcdt 0
12 Artrite Reumatoide M050,M051,M052,M053,M058,M060,M068 AREU artrite_reu Artrite Reumatoide NA NA 0
13 Asma J450,J451,J458,J45 ASMA asma Asma Asma pcdt;csap 0
14 Atrofia Muscular Espinhal G120,G121 AME atrofia Atrofia Muscular Espinhal 5q Tipo 1 Atrofia Muscular Espinhal pcdt 1
15 COVID-19 B342,U04,U049,U071,U072,U089,U099,U109,U129 COV covid_19 Diretrizes Brasileiras para Tratamento Hospitalar do Paciente com COVID-19 COVID-19 Diretrizes 0
16 Carcinoma Diferenciado de Tireoide C73 CTIR c_tireoide Carcinoma Diferenciado de Tireoide NA NA 0
17 Carcinoma de Células Renais C64 CREN c_renal Carcinoma de Células Renais Carcinoma de Células Renais pcdt 0
18 Carcinoma de Mama C500,C501,C502,C503,C504,C505,C506,C508,C509 CMAM c_mama Carcinoma de Mama Carcinoma de Mama ddt 0
19 Carcinoma de estômago e esôfago C150,C151,C152,C153,C154,C155,C158,C159,C160,C161,C162,C163,C164,C165,C166,C168,C169,C170,C171,C172,C173,C178,C179,C18,C180,C181,C182,C183,C184,C185,C186,C187,C188,C189,C19,C20,C268,C474,C481,C493 CEES c_estomago Carcinoma de estômago e esôfago NA NA 0
20 Carcinoma hepatocelular C22 CHEP c_hepatocel Carcinoma hepatocelular NA NA 0

Estudo de Caso 1: Artrite Reumatoide em São Paulo (2023-2024)

Para tratar volumes massivos como os do estado de São Paulo sem estourar a memória RAM, aplicamos uma rotina iterativa de ETL (mês a mês):

  1. Definição de Parâmetros: Configuramos o diretório /media/ferre/ferre500gb/Downloads/dbc/, o parâmetro booleano manter_dbc (para apagar os arquivos brutos e economizar espaço em disco) e a lista de CIDs para Artrite Reumatoide (M050, M051, M052, M053, M058, M060, M068).
  2. Download e Leitura: Baixamos individualmente os arquivos PA (Produção) e AM (Medicamentos) do mês.
  3. Enriquecimento: Cruzamos o número da autorização (APAC) para integrar a coluna AP_CNSPCN (CNS do paciente) do arquivo AM no arquivo PA.
  4. Filtragem e Limpeza: Filtramos apenas os registros que apresentam os CIDs de Artrite Reumatoide em PA_CIDPRI ou PA_CIDSEC, descartando os registros irrelevantes. Removemos os arquivos .dbc do disco (caso manter_dbc = FALSE) e liberamos a RAM com gc().
  5. Exportação: Salvamos a base consolidada e reduzida em formato .csv.gz.

Observações:

  • Versao de BAIXA MEMORIA:
  • não usa /tmp ou tempdir() (roda em dir_dbc)
  • DBF lido em blocos de registros (não roda inteiro na RAM)
  • filtro de CID aplicado dentro do bloco (aplicar SIGTAP se desejável)
  • resultado gravado incrementalmente em disco (não acumula data.frame)
  • O pico de RAM esperado é inferior a 500 MB

Script ETL para Artrite Reumatoide (SP)

###############################################################################
# SIASUS / SIA-PA  ->  artrite reumatoide, SP, 2023-2024
###############################################################################

suppressPackageStartupMessages({
  library(curl)
  library(read.dbc)   # usado apenas para dbc2dbf() (descompactacao para disco)
})

## ---------------------------------------------------------------- parametros
dir_dbc      <- "/media/ferre/ferre500gb/Downloads/dbc/"
dir_tmp      <- file.path(dir_dbc, "tmp")
dir_saida    <- file.path(dir_dbc, "saida")
dir_comp     <- file.path(dir_saida, "competencias")   # 1 CSV por competencia
dir_pq       <- file.path(dir_saida, "parquet")        # dataset particionado
manter_dbc   <- TRUE          # manter os .dbc baixados
manter_dbf   <- FALSE         # NUNCA deixe TRUE: .dbf ocupa GBs
refazer      <- FALSE         # TRUE ignora os checkpoints e reprocessa tudo
url_ftp      <- "ftp://ftp.datasus.gov.br/dissemin/publicos/SIASUS/200801_/Dados/"
uf_alvo      <- "SP"
anos_alvo    <- c(23, 24)
meses_alvo   <- 1:12
cids_artrite <- c("M050", "M051", "M052", "M053", "M058", "M060", "M068")
CHUNK        <- 50000L        # registros por bloco; reduza se ainda pesar
decodifica_cns   <- TRUE      # traduz o AP_CNSPCN ofuscado para digitos
guarda_cns_bruto <- FALSE     # TRUE mantem tambem a coluna AP_CNSPCN_BRUTO

usar_parquet <- requireNamespace("arrow", quietly = TRUE)

# Desofusca o AP_CNSPCN do SIA (mesma tabela dos seds usuais).
converte_cns <- function(x) {
  de   <- c("{", "}", "~", "\u007f", "\u00c7", "\u00e4", "\u00fc", "\u00e9", "|", "\u00e2")
  para <- c("0", "9", "8", "7", "6", "5", "4", "3", "2", "1")
  for (i in seq_along(de)) x <- gsub(de[i], para[i], x, fixed = TRUE)
  x <- gsub("[^0-9]", "4", x)
  ifelse(nzchar(x), x, NA_character_)
}

for (d in c(dir_dbc, dir_tmp, dir_saida, dir_comp, dir_pq))
  dir.create(d, showWarnings = FALSE, recursive = TRUE)
Sys.setenv(TMPDIR = dir_tmp, TMP = dir_tmp, TEMP = dir_tmp)

## ------------------------------------------------- leitor de DBF em blocos
# Le o cabecalho dBase III: n de registros, tamanho do cabecalho/registro e
# a tabela de campos (nome, tipo, largura).
dbf_header <- function(path) {
  con <- file(path, "rb"); on.exit(close(con))
  h     <- readBin(con, "raw", 32L)
  n_rec <- readBin(h[5:8],   "integer", size = 4, endian = "little")
  h_len <- readBin(h[9:10],  "integer", size = 2, signed = FALSE, endian = "little")
  r_len <- readBin(h[11:12], "integer", size = 2, signed = FALSE, endian = "little")
  nf    <- (h_len - 32L - 1L) %/% 32L
  fd    <- readBin(con, "raw", nf * 32L)
  dim(fd) <- c(32L, nf)
  nomes <- apply(fd[1:11, , drop = FALSE], 2L, function(x) rawToChar(x[x != as.raw(0)]))
  tipos <- rawToChar(fd[12, ], multiple = TRUE)
  lens  <- as.integer(fd[17, ])
  ini   <- 2L + c(0L, cumsum(lens))[seq_len(nf)]   # +1 do flag de exclusao
  list(n_rec = n_rec, h_len = h_len, r_len = r_len,
       nomes = nomes, tipos = tipos, ini = ini, fim = ini + lens - 1L)
}

# Percorre o DBF em blocos. Para cada bloco chama FUN(getcol, validos, nomes):
#   getcol(nome, linhas) -> vetor character (trim aplicado)
#   validos              -> indices das linhas nao marcadas como excluidas
dbf_ler_chunks <- function(path, FUN) {
  hd  <- dbf_header(path)
  con <- file(path, "rb"); on.exit(close(con))
  invisible(readBin(con, "raw", hd$h_len))
  lidos <- 0L
  while (lidos < hd$n_rec) {
    n   <- min(CHUNK, hd$n_rec - lidos)
    blk <- readBin(con, "raw", n * hd$r_len)
    if (length(blk) < n * hd$r_len) n <- length(blk) %/% hd$r_len
    if (n == 0L) break
    length(blk) <- n * hd$r_len
    blk[blk == as.raw(0)] <- as.raw(32)              # nulos -> espaco
    s <- rawToChar(blk); rm(blk)
    Encoding(s) <- "latin1"                          # 1 byte = 1 char (DATASUS)
    base <- (seq_len(n) - 1L) * hd$r_len

    validos <- which(substring(s, base + 1L, base + 1L) != "*")

    getcol <- function(nome, linhas = validos) {
      j <- match(nome, hd$nomes)
      if (is.na(j)) return(rep(NA_character_, length(linhas)))
      b <- base[linhas]
      x <- substring(s, b + hd$ini[j], b + hd$fim[j])
      # o AP_CNSPCN traz bytes altos; sem o iconv explicito o trimws (PCRE)
      # tenta traduzir para UTF-8 e falha com "string de entrada invalida"
      x <- iconv(x, from = "latin1", to = "UTF-8", sub = "")
      trimws(x)
    }

    FUN(getcol, validos, hd$nomes)
    lidos <- lidos + n
    rm(s); gc(FALSE)
  }
  invisible(hd$n_rec)
}

## ------------------------------------------------------------ FTP / download
texto_ftp    <- rawToChar(curl_fetch_memory(url_ftp)$content)
arquivos_ftp <- unique(
  regmatches(texto_ftp, gregexpr("[[:alnum:]_-]+\\.[dD][bB][cC]", texto_ftp))[[1]]
)
rm(texto_ftp)

listar_partes <- function(tipo, ano, mes) {
  padrao <- sprintf("^%s%s%02d%02d([A-Z]|_[0-9]+)?\\.DBC$", tipo, uf_alvo, ano, mes)
  arquivos_ftp[grepl(padrao, toupper(arquivos_ftp))]
}

baixar_se_ausente <- function(url, destino) {
  if (file.exists(destino) && file.size(destino) > 0) return(TRUE)
  parcial <- paste0(destino, ".part")
  ok <- try(curl_download(url, parcial, quiet = TRUE), silent = TRUE)
  if (inherits(ok, "try-error")) {
    if (file.exists(parcial)) file.remove(parcial)
    return(FALSE)
  }
  file.rename(parcial, destino)   # so vira definitivo se o download completou
  TRUE
}

# .dbc -> .dbf DENTRO de dir_tmp (nunca em /tmp), devolve o caminho do .dbf
preparar_dbf <- function(arq) {
  dest_dbc <- file.path(dir_dbc, arq)
  if (!baixar_se_ausente(paste0(url_ftp, arq), dest_dbc)) return(NULL)
  dest_dbf <- file.path(dir_tmp, sub("\\.[dD][bB][cC]$", ".dbf", arq))
  if (!file.exists(dest_dbf) || file.size(dest_dbf) == 0) {
    ok <- try(dbc2dbf(dest_dbc, dest_dbf), silent = TRUE)
    if (inherits(ok, "try-error")) {
      if (file.exists(dest_dbf)) file.remove(dest_dbf)
      return(NULL)
    }
  }
  dest_dbf
}

## ------------------------------------------------ checkpoint por competencia
# Estados possiveis de uma competencia:
#   <tag>.csv        -> processada, com registros
#   <tag>.vazio      -> processada, nenhum registro de artrite (evita refazer)
#   nenhum dos dois  -> pendente
caminho_csv   <- function(ano, mes)
  file.path(dir_comp, sprintf("artrite_%s_20%02d-%02d.csv", uf_alvo, ano, mes))
caminho_vazio <- function(ano, mes) paste0(caminho_csv(ano, mes), ".vazio")

ja_processada <- function(ano, mes)
  !refazer && (file.exists(caminho_csv(ano, mes)) || file.exists(caminho_vazio(ano, mes)))

# Gravacao atomica: escreve em .part e so entao renomeia.
gravar_competencia <- function(df, ano, mes) {
  destino <- caminho_csv(ano, mes)
  parcial <- paste0(destino, ".part")
  write.table(df, parcial, sep = ",", row.names = FALSE, col.names = TRUE,
              qmethod = "double", na = "", fileEncoding = "UTF-8")
  file.rename(parcial, destino)

  if (usar_parquet) {
    dir_part <- file.path(dir_pq, sprintf("ano=20%02d", ano), sprintf("mes=%02d", mes))
    dir.create(dir_part, showWarnings = FALSE, recursive = TRUE)
    alvo_pq <- file.path(dir_part, "part-0.parquet")
    tmp_pq  <- paste0(alvo_pq, ".part")
    ok <- try(arrow::write_parquet(df, tmp_pq, compression = "zstd"), silent = TRUE)
    if (inherits(ok, "try-error")) {
      if (file.exists(tmp_pq)) file.remove(tmp_pq)
      cat("  ! falha ao gravar parquet desta competencia\n")
    } else {
      file.rename(tmp_pq, alvo_pq)
    }
  }
  invisible(nrow(df))
}

## -------------------------------------------------------------- ETL iterativa
for (ano in anos_alvo) for (mes in meses_alvo) {

  if (ja_processada(ano, mes)) {
    cat(sprintf("20%02d-%02d | ja processada, pulando\n", ano, mes))
    next
  }

  arqs_pa <- listar_partes("PA", ano, mes)
  arqs_am <- listar_partes("AM", ano, mes)
  cat(sprintf("20%02d-%02d | PA: %d parte(s) | AM: %d parte(s)\n",
              ano, mes, length(arqs_pa), length(arqs_am)))
  # sem arquivo no FTP: nao marca nada, para reprocessar quando for publicado
  if (length(arqs_pa) == 0L) next

  # ---- 1) PA: filtra CID dentro de cada bloco; so o filtrado vira data.frame
  buf   <- list()
  falha <- FALSE
  for (a in arqs_pa) {
    dbf <- preparar_dbf(a)
    if (is.null(dbf)) { cat("  ! falha em", a, "\n"); falha <- TRUE; next }
    cat("  lendo", a, "...")
    dbf_ler_chunks(dbf, function(getcol, validos, nomes) {
      c1  <- substr(getcol("PA_CIDPRI"), 1, 4)
      c2  <- substr(getcol("PA_CIDSEC"), 1, 4)
      hit <- which(c1 %in% cids_artrite | c2 %in% cids_artrite)
      if (!length(hit)) return(invisible(NULL))
      linhas <- validos[hit]
      d <- lapply(nomes, function(nm) getcol(nm, linhas))
      names(d) <- nomes
      d$ORIGEM <- a                       # arquivo .dbc de procedencia
      buf[[length(buf) + 1L]] <<- as.data.frame(d, stringsAsFactors = FALSE)
      invisible(NULL)
    })
    cat(" ok\n")
    if (!manter_dbf && file.exists(dbf)) file.remove(dbf)
    if (!manter_dbc) file.remove(file.path(dir_dbc, a))
    gc(FALSE)
  }

  # alguma parte falhou: nao marca a competencia, para tentar de novo depois
  if (falha) { cat("  ! competencia incompleta, sera refeita\n"); rm(buf); gc(FALSE); next }

  if (!length(buf)) {
    file.create(caminho_vazio(ano, mes))
    cat("  -> 0 registros de artrite\n")
    rm(buf); gc(FALSE); next
  }
  df_pa <- do.call(rbind, buf); rm(buf); gc(FALSE)
  cat(sprintf("  -> %d registros de artrite\n", nrow(df_pa)))

  # ---- 2) AM: le em blocos e guarda APENAS as chaves que casam com o PA
  chaves <- paste0(df_pa$PA_CMP, "|", df_pa$PA_AUTORIZ)
  alvo   <- unique(chaves)
  mapa   <- character(0)

  for (a in arqs_am) {
    dbf <- preparar_dbf(a)
    if (is.null(dbf)) next
    dbf_ler_chunks(dbf, function(getcol, validos, nomes) {
      k   <- paste0(getcol("AP_CMP"), "|", getcol("AP_AUTORIZ"))
      hit <- which(k %in% alvo)
      if (!length(hit)) return(invisible(NULL))
      v <- getcol("AP_CNSPCN", validos[hit])
      names(v) <- k[hit]
      mapa <<- c(mapa, v[!duplicated(names(v))])
      invisible(NULL)
    })
    if (!manter_dbf && file.exists(dbf)) file.remove(dbf)
    if (!manter_dbc) file.remove(file.path(dir_dbc, a))
    gc(FALSE)
  }
  mapa <- mapa[!duplicated(names(mapa))]

  cns <- if (length(mapa)) unname(mapa[chaves]) else rep(NA_character_, nrow(df_pa))
  if (guarda_cns_bruto) df_pa$AP_CNSPCN_BRUTO <- cns
  df_pa$AP_CNSPCN <- if (decodifica_cns) converte_cns(cns) else cns
  cat(sprintf("  -> CNS recuperado para %d de %d registros\n",
              sum(!is.na(df_pa$AP_CNSPCN)), nrow(df_pa)))
  rm(chaves, alvo, mapa, cns)

  # ---- 3) grava a competencia (checkpoint) e descarta da memoria
  gravar_competencia(df_pa, ano, mes)
  rm(df_pa); gc(FALSE)
}
Mostrar resultado
## 2023-01 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2301a.dbc ... ok
##   lendo PASP2301b.dbc ... ok
##   lendo PASP2301c.dbc ... ok
##   -> 78522 registros de artrite
##   -> CNS recuperado para 77604 de 78522 registros
## 2023-02 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2302a.dbc ... ok
##   lendo PASP2302b.dbc ... ok
##   lendo PASP2302c.dbc ... ok
##   -> 74137 registros de artrite
##   -> CNS recuperado para 73261 de 74137 registros
## 2023-03 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2303a.dbc ... ok
##   lendo PASP2303b.dbc ... ok
##   lendo PASP2303c.dbc ... ok
##   -> 73533 registros de artrite
##   -> CNS recuperado para 72404 de 73533 registros
## 2023-04 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2304a.dbc ... ok
##   lendo PASP2304b.dbc ... ok
##   lendo PASP2304c.dbc ... ok
##   -> 72376 registros de artrite
##   -> CNS recuperado para 71377 de 72376 registros
## 2023-05 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2305a.dbc ... ok
##   lendo PASP2305b.dbc ... ok
##   lendo PASP2305c.dbc ... ok
##   -> 75174 registros de artrite
##   -> CNS recuperado para 74074 de 75174 registros
## 2023-06 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2306a.dbc ... ok
##   lendo PASP2306b.dbc ... ok
##   lendo PASP2306c.dbc ... ok
##   -> 75128 registros de artrite
##   -> CNS recuperado para 73942 de 75128 registros
## 2023-07 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2307a.dbc ... ok
##   lendo PASP2307b.dbc ... ok
##   lendo PASP2307c.dbc ... ok
##   -> 75719 registros de artrite
##   -> CNS recuperado para 75448 de 75719 registros
## 2023-08 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2308a.dbc ... ok
##   lendo PASP2308b.dbc ... ok
##   lendo PASP2308c.dbc ... ok
##   -> 77065 registros de artrite
##   -> CNS recuperado para 76775 de 77065 registros
## 2023-09 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2309a.dbc ... ok
##   lendo PASP2309b.dbc ... ok
##   lendo PASP2309c.dbc ... ok
##   -> 75587 registros de artrite
##   -> CNS recuperado para 75327 de 75587 registros
## 2023-10 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2310a.dbc ... ok
##   lendo PASP2310b.dbc ... ok
##   lendo PASP2310c.dbc ... ok
##   -> 75754 registros de artrite
##   -> CNS recuperado para 75492 de 75754 registros
## 2023-11 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2311a.dbc ... ok
##   lendo PASP2311b.dbc ... ok
##   lendo PASP2311c.dbc ... ok
##   -> 74428 registros de artrite
##   -> CNS recuperado para 74209 de 74428 registros
## 2023-12 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2312a.dbc ... ok
##   lendo PASP2312b.dbc ... ok
##   lendo PASP2312c.dbc ... ok
##   -> 74165 registros de artrite
##   -> CNS recuperado para 73906 de 74165 registros
## 2024-01 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2401a.dbc ... ok
##   lendo PASP2401b.dbc ... ok
##   lendo PASP2401c.dbc ... ok
##   -> 70650 registros de artrite
##   -> CNS recuperado para 70339 de 70650 registros
## 2024-02 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2402a.dbc ... ok
##   lendo PASP2402b.dbc ... ok
##   lendo PASP2402c.dbc ... ok
##   -> 72667 registros de artrite
##   -> CNS recuperado para 72363 de 72667 registros
## 2024-03 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2403a.dbc ... ok
##   lendo PASP2403b.dbc ... ok
##   lendo PASP2403c.dbc ... ok
##   -> 77155 registros de artrite
##   -> CNS recuperado para 76835 de 77155 registros
## 2024-04 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2404a.dbc ... ok
##   lendo PASP2404b.dbc ... ok
##   lendo PASP2404c.dbc ... ok
##   -> 77515 registros de artrite
##   -> CNS recuperado para 77141 de 77515 registros
## 2024-05 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2405a.dbc ... ok
##   lendo PASP2405b.dbc ... ok
##   lendo PASP2405c.dbc ... ok
##   -> 77305 registros de artrite
##   -> CNS recuperado para 77026 de 77305 registros
## 2024-06 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2406a.dbc ... ok
##   lendo PASP2406b.dbc ... ok
##   lendo PASP2406c.dbc ... ok
##   -> 78190 registros de artrite
##   -> CNS recuperado para 77871 de 78190 registros
## 2024-07 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2407a.dbc ... ok
##   lendo PASP2407b.dbc ... ok
##   lendo PASP2407c.dbc ... ok
##   -> 79742 registros de artrite
##   -> CNS recuperado para 79325 de 79742 registros
## 2024-08 | PA: 4 parte(s) | AM: 1 parte(s)
##   lendo PASP2408a.dbc ... ok
##   lendo PASP2408b.dbc ... ok
##   lendo PASP2408c.dbc ... ok
##   lendo PASP2408d.dbc ... ok
##   -> 80772 registros de artrite
##   -> CNS recuperado para 80266 de 80772 registros
## 2024-09 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2409a.dbc ... ok
##   lendo PASP2409b.dbc ... ok
##   lendo PASP2409c.dbc ... ok
##   -> 76866 registros de artrite
##   -> CNS recuperado para 76446 de 76866 registros
## 2024-10 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2410a.dbc ... ok
##   lendo PASP2410b.dbc ... ok
##   lendo PASP2410c.dbc ... ok
##   -> 81322 registros de artrite
##   -> CNS recuperado para 80864 de 81322 registros
## 2024-11 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2411a.dbc ... ok
##   lendo PASP2411b.dbc ... ok
##   lendo PASP2411c.dbc ... ok
##   -> 82093 registros de artrite
##   -> CNS recuperado para 80771 de 82093 registros
## 2024-12 | PA: 3 parte(s) | AM: 1 parte(s)
##   lendo PASP2412a.dbc ... ok
##   lendo PASP2412b.dbc ... ok
##   lendo PASP2412c.dbc ... ok
##   -> 81925 registros de artrite
##   -> CNS recuperado para 80157 de 81925 registros
## ------------------------------------------- consolidacao final dos CSVs
# Os CSVs por competencia sao preservados. Este passo apenas produz um
# arquivo unico, em streaming, sem carregar nada inteiro na memoria.
consolidar <- function() {
  arqs <- sort(list.files(dir_comp, pattern = "\\.csv$", full.names = TRUE))
  if (!length(arqs)) { cat("Nenhuma competencia processada.\n"); return(invisible(NULL)) }

  cabecalhos <- lapply(arqs, function(f)
    scan(f, what = "", sep = ",", nlines = 1L, quiet = TRUE))
  iguais <- all(vapply(cabecalhos, identical, logical(1), cabecalhos[[1]]))

  final   <- file.path(dir_saida, sprintf("artrite_reumatoide_%s.csv", tolower(uf_alvo)))
  parcial <- paste0(final, ".part")
  out <- file(parcial, open = "wt", encoding = "UTF-8")

  if (iguais) {
    for (i in seq_along(arqs)) {
      inp <- file(arqs[i], open = "rt", encoding = "UTF-8")
      primeiro <- TRUE
      repeat {
        linhas <- readLines(inp, n = 50000L)
        if (!length(linhas)) break
        if (primeiro && i > 1L) linhas <- linhas[-1L]  # descarta cabecalho repetido
        primeiro <- FALSE
        if (length(linhas)) writeLines(linhas, out)
      }
      close(inp)
    }
  } else {
    # layouts diferentes entre anos: alinha pela uniao das colunas
    uniao <- unique(unlist(cabecalhos))
    cat("! layouts divergentes; alinhando", length(uniao), "colunas\n")
    for (i in seq_along(arqs)) {
      d <- read.csv(arqs[i], colClasses = "character", check.names = FALSE,
                    na.strings = "", encoding = "UTF-8")
      for (nm in setdiff(uniao, names(d))) d[[nm]] <- NA_character_
      d <- d[, uniao, drop = FALSE]
      write.table(d, out, sep = ",", row.names = FALSE, col.names = (i == 1L),
                  qmethod = "double", na = "")
      rm(d); gc(FALSE)
    }
  }
  close(out)
  file.rename(parcial, final)

  # contagem em streaming: nao carrega o arquivo final na memoria
  n <- 0L
  cc <- file(final, open = "rt", encoding = "UTF-8")
  repeat { k <- length(readLines(cc, n = 50000L)); if (!k) break; n <- n + k }
  close(cc)
  cat("Registros consolidados:", n - 1L, "\n")   # -1 do cabecalho

  if (nzchar(Sys.which("gzip"))) {
    system2("gzip", c("-f", shQuote(final)))
    cat("CSV unico:", paste0(final, ".gz"), "\n")
  } else {
    cat("CSV unico:", final, "\n")
  }
}

consolidar()
Mostrar resultado
## Registros consolidados: 1837790 
## CSV unico: /media/ferre/ferre500gb/Downloads/dbc//saida/artrite_reumatoide_sp.csv.gz
cat("CSVs por competencia:", dir_comp, "\n")
Mostrar resultado
## CSVs por competencia: /media/ferre/ferre500gb/Downloads/dbc//saida/competencias
if (usar_parquet) cat("Dataset parquet (particionado ano/mes):", dir_pq, "\n")
Mostrar resultado
## Dataset parquet (particionado ano/mes): /media/ferre/ferre500gb/Downloads/dbc//saida/parquet

Leitura analítica do arquivo compactado sem sobrecarga de RAM

Após gerar o arquivo .csv.gz, podemos inspecioná-lo utilizando o pacote arrow, que acessa os dados diretamente do disco sem carregá-los integralmente na RAM.

library(arrow)
library(dplyr)

arquivo_saida_sp <- file.path("/media/ferre/ferre500gb/Downloads/dbc/", "artrite_reumatoide_sp_23_24.csv.gz")

if (file.exists(arquivo_saida_sp)) {
  
  # Abertura preguiçosa (lazy execution)
  dataset_artrite <- open_dataset(arquivo_saida_sp, format = "csv")
  
  # Selecao e coleta de apenas 20 linhas para exibicao
  amostra_sp <- dataset_artrite |>
    select(PA_CMP, PA_AUTORIZ, AP_CNSPCN, PA_CIDPRI, PA_PROC_ID, PA_QTDAPR, PA_VALAPR) |>
    head(20) |>
    collect()
  
  print(amostra_sp)
}

Estudo de Caso 2: Processamento Nacional Completo (Doença de Paget)

Expandimos a solução para cobrir todas as 27 UFs do Brasil a partir de 2020 (2020 a 2024), buscando especificamente a Doença de Paget (M880, M888).

Script ETL Nacional para Doença de Paget


Síntese do Pipeline de Extração e Validação

O pipeline consolidado nesta aula resolve o desafio de extrair dados do SUS em escala nacional utilizando computadores normais:

  1. Mapeamento Prévio de Metadados: Identificação dos CIDs e procedimentos no PCDT antes do download para evitar o carregamento de dados desnecessários.
  2. Gerenciamento Transparente de Armazenamento: Uso de parâmetros booleanos (manter_dbc) e diretórios locais estruturados (/media/ferre/...) para controlar a retenção de arquivos brutos.
  3. Fracionamento Estruturado (Staging Iterativo): Processamento em laços UF -> Ano -> Mês, limitando o pico de consumo de RAM ao tamanho de um único arquivo mensal.
  4. Enriquecimento On-the-fly: Junção entre subsistemas (SIA/PA com SIA/AM) realizada individualmente por competência antes da filtragem.
  5. Redução Agressiva e Limpeza: Filtro imedioto por diagnósticos de interesse seguido da remoção de objetos da RAM (rm(), gc()) e exclusão de arquivos no disco (file.remove()).
  6. Persistência Analítica Eficiente: Salvamento em formatos compactados e leitura preguiçosa (arrow::open_dataset), permitindo consultas exploratórias ágeis sem reidratar gigabytes na RAM.