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.
etl_sia() trata o PA como base e os
demais subsistemas do SIA como fontes de enriquecimento, cruzando
exatamente por:
PA_AUTORIZ ↔︎ AP_AUTORIZ
e restringindo o cruzamento aos arquivos cujo miolo
UFAAMM seja o mesmo da competência/UF do PA. Isso
generaliza a estratégia que o script original usava especificamente para
AM.
O que mudou nesta versão. A versão anterior baixava os subsistemas auxiliares, marcava-os como concluídos e não os cruzava com o PA — a documentação prometia um join que o código não executava. Agora o enriquecimento é real, o
AP_CNSPCNé desofuscado, todo texto do DBF é convertido delatin1paraUTF-8e a função pergunta se é para continuar uma execução anterior ou criar uma nova.
subsistema =:
AB, ABO, ACF, AD,
AM, AMP, AN, AQ,
AR, ATD, SAD.PA_AUTORIZ ↔︎
AP_AUTORIZ.UF + AAMM do PA.subsistema = NULL, todos os auxiliares são
considerados.enriquecer = FALSE desliga o cruzamento e volta ao
comportamento antigo.A lista de arquivos é derivada dos nomes efetivamente encontrados no FTP, em vez de presumir uma única parte. O padrão aceita:
PASP2301.DBC
PASP2301A.DBC
PASP2301B.DBC
PASP2301C.DBC
PASP2301_1.DBC
PASP2301_2.DBC
PASP23011.DBC
PASP23012.DBC
Para cada registro PA:
PA_AUTORIZ
│
├── AB → AP_AUTORIZ
├── ABO → AP_AUTORIZ
├── ACF → AP_AUTORIZ
├── AD → AP_AUTORIZ
├── AM → AP_AUTORIZ
├── AMP → AP_AUTORIZ
├── AN → AP_AUTORIZ
├── AQ → AP_AUTORIZ
├── AR → AP_AUTORIZ
├── ATD → AP_AUTORIZ
└── SAD → AP_AUTORIZ
e sempre considerando somente os arquivos do mesmo
UFAAMM.
O cruzamento não varre o auxiliar a cada chunk do
PA. Antes de processar o PA de uma competência/UF, a função grava um
índice em Parquet por AP_AUTORIZ:
<saida>/_etl_sia_estado<sufixo>/indice_auxiliares/202301_SP_AQ.parquet
Cada chunk do PA faz apenas um match() contra esse
índice já carregado. O índice é reaproveitado entre chunks, entre partes
do PA e entre execuções retomadas — é a alternativa sugerida na versão
anterior do documento, agora implementada.
Quando uma mesma APAC tem mais de um laudo no auxiliar,
mantém-se o primeiro, para que o join não multiplique
linhas do PA. Toda linha ganha também
<sub>_st_vinculado
(TRUE/FALSE), permitindo medir a taxa de
vínculo.
Adiciona os atributos conforme o subsistema:
AM
st_gestante
nu_peso
nu_altura
AQ
aq_cid10
aq_linfin
aq_estadi
aq_grahis
aq_dtiden
aq_trante
aq_cidini1
aq_dtini1
aq_cidini2
aq_dtini2
aq_cidini3
aq_dtini3
aq_conttr
aq_dtintr
aq_esqu_p1
aq_totmpl
aq_totmau
aq_esqu_p2
aq_med01 ... aq_med10
AR
ar_smrd
ar_cid10
ar_linfin
ar_estadi
ar_grahis
ar_dtiden
ar_trante
ar_cidini1
ar_dtini1
ar_cidini2
ar_dtini2
ar_cidini3
ar_dtini3
ar_conttr
ar_dtintr
ar_finali
ar_cidtr1
ar_cidtr2
ar_cidtr3
ar_numc1
ar_numc2
ar_numc3
ar_iniar1
ar_iniar2
ar_iniar3
ar_fimar1
ar_fimar2
ar_fimar3
Além de AP_CNSPCN.
Para os demais subsistemas (AB, ABO,
ACF, AD, AMP, AN,
ATD, SAD) vale a regra genérica: todos os
campos com o prefixo do subsistema presentes no cabeçalho do DBF, em
minúsculas, mais AP_CNSPCN.
As colunas esperadas são gravadas mesmo quando o auxiliar não
existe para aquela UF/competência (preenchidas com
NA), de modo que o dataset Parquet mantenha esquema estável
e possa ser lido de uma vez com arrow::open_dataset().
Como você pode selecionar, por exemplo:
A função grava, para fins de rastreabilidade:
am_ap_cnspcn
aq_ap_cnspcn
ar_ap_cnspcn
enquanto mantém os campos específicos exatamente como definidos:
st_gestante
nu_peso
nu_altura
aq_cid10
aq_estadi
...
ar_smrd
ar_cid10
...
Se ainda assim um nome colidir com uma coluna do próprio PA, o campo do auxiliar recebe o prefixo do subsistema. Isso evita perda de informação.
O AP_CNSPCN do SIA (e o PA_CNSMED do PA)
chegam ofuscados pela tabela usual dos sistemas de disseminação. A
conversão é feita durante a construção do índice, antes da gravação em
Parquet:
# 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_)
}Controles:
converter_cns = TRUE (padrão) aplica a conversão;manter_cns_original = TRUE preserva o valor bruto em
<coluna>_bruto / PA_CNSMED_BRUTO, para
auditoria.Os DBF do DATASUS estão em latin1. Todo valor extraído
passa por:
Dois cuidados que o código já resolve:
trimws() e gsub() traduzem implicitamente uma
string marcada como latin1 para UTF-8; um
iconv() depois disso produziria dupla codificação (os bytes
de à reinterpretados como latin1, gerando
à + lixo). Por isso o recorte é feito sobre o texto ainda
marcado como latin1, converte-se, e só então apara-se.td_diretriz_cuidado.csv — passam por
texto_para_utf8(), que só converte quando os bytes não são
UTF-8 válidos.###############################################################################
# etl_sia() - ETL incremental de SIA/PA para Parquet
# Saida: dataset Parquet incremental (diretorio com extensao .parquet),
# permitindo retomada sem reprocessar competencias/arquivos concluidos.
###############################################################################
suppressPackageStartupMessages({
library(curl)
library(read.dbc)
library(arrow)
})
## ---------------------------------------------------------------------------
## Utilitarios de nivel superior (reutilizaveis fora do etl_sia)
## ---------------------------------------------------------------------------
# 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_)
}
# Toda string vinda do DBF nasce em latin1; o restante do pipeline (Parquet,
# joins, regex de CID) assume UTF-8.
para_utf8 <- function(x, aparar = TRUE) {
if (!length(x)) return(x)
x <- iconv(x, from = "latin1", to = "UTF-8", sub = "")
if (isTRUE(aparar)) trimws(x) else x
}
# Para documentos inteiros (listagem do FTP, CSV do PCDT), que ja podem chegar
# em UTF-8: converte apenas se os bytes nao forem UTF-8 validos.
texto_para_utf8 <- function(x) {
if (!length(x)) return(x)
ok <- validUTF8(x)
x[!ok] <- iconv(x[!ok], from = "latin1", to = "UTF-8", sub = "")
Encoding(x) <- "UTF-8"
x
}
# Campos gravados no PA para cada subsistema auxiliar.
# Regra: nome de saida = nome do campo no DBF em minusculas.
# Excecao: AM usa a nomenclatura padronizada (st_, nu_).
campos_subsistema <- function(sub) {
sub <- toupper(sub)
if (identical(sub, "AM")) {
return(c(AM_GESTANT = "st_gestante",
AM_PESO = "nu_peso",
AM_ALTURA = "nu_altura"))
}
if (identical(sub, "AQ")) {
nomes <- c("AQ_CID10", "AQ_LINFIN", "AQ_ESTADI", "AQ_GRAHIS", "AQ_DTIDEN",
"AQ_TRANTE", "AQ_CIDINI1", "AQ_DTINI1", "AQ_CIDINI2", "AQ_DTINI2",
"AQ_CIDINI3", "AQ_DTINI3", "AQ_CONTTR", "AQ_DTINTR", "AQ_ESQU_P1",
"AQ_TOTMPL", "AQ_TOTMAU", "AQ_ESQU_P2",
sprintf("AQ_MED%02d", 1:10))
return(setNames(tolower(nomes), nomes))
}
if (identical(sub, "AR")) {
nomes <- c("AR_SMRD", "AR_CID10", "AR_LINFIN", "AR_ESTADI", "AR_GRAHIS",
"AR_DTIDEN", "AR_TRANTE", "AR_CIDINI1", "AR_DTINI1", "AR_CIDINI2",
"AR_DTINI2", "AR_CIDINI3", "AR_DTINI3", "AR_CONTTR", "AR_DTINTR",
"AR_FINALI", "AR_CIDTR1", "AR_CIDTR2", "AR_CIDTR3",
"AR_NUMC1", "AR_NUMC2", "AR_NUMC3",
"AR_INIAR1", "AR_INIAR2", "AR_INIAR3",
"AR_FIMAR1", "AR_FIMAR2", "AR_FIMAR3")
return(setNames(tolower(nomes), nomes))
}
## Demais subsistemas: regra generica por prefixo, resolvida em tempo de
## leitura contra o cabecalho real do DBF (ver mapa_campos()).
character(0)
}
# Situacao das execucoes gravadas em um diretorio de saida.
status_etl_sia <- function(diretorio_saida, imprimir = TRUE) {
if (!dir.exists(diretorio_saida)) {
if (imprimir) cat("Nenhuma execucao encontrada em ", diretorio_saida, "\n", sep = "")
return(invisible(data.frame()))
}
dirs <- list.dirs(diretorio_saida, recursive = FALSE)
estados <- dirs[grepl("^_etl_sia_estado", basename(dirs))]
if (!length(estados)) {
if (imprimir) cat("Nenhuma execucao encontrada em ", diretorio_saida, "\n", sep = "")
return(invisible(data.frame()))
}
linhas <- lapply(estados, function(e) {
sufixo <- sub("^_etl_sia_estado", "", basename(e))
done <- list.files(e, pattern = "\\.done$", recursive = TRUE)
ds <- list.dirs(diretorio_saida, recursive = FALSE)
ds <- ds[grepl(paste0(sufixo, "\\.parquet$"), basename(ds))]
if (!nzchar(sufixo)) ds <- ds[!grepl("_[0-9]{4}-[0-9]{2}-[0-9]{2}-[0-9]{4}\\.parquet$", basename(ds))]
partes <- if (length(ds)) length(list.files(ds, pattern = "\\.parquet$", recursive = TRUE)) else 0L
bytes <- if (length(ds)) sum(file.size(list.files(ds, pattern = "\\.parquet$",
recursive = TRUE, full.names = TRUE)), na.rm = TRUE) else 0
mt <- suppressWarnings(max(file.mtime(c(e, list.files(e, full.names = TRUE))), na.rm = TRUE))
data.frame(
execucao = if (nzchar(sufixo)) sub("^_", "", sufixo) else "(inicial)",
sufixo = sufixo,
arquivos_concluidos = length(done),
datasets = length(ds),
partes_parquet = partes,
mb = round(bytes / 1024^2, 1),
ultima_atividade = format(as.POSIXct(mt, origin = "1970-01-01"), "%Y-%m-%d %H:%M"),
stringsAsFactors = FALSE
)
})
out <- do.call(rbind, linhas)
out <- out[order(out$ultima_atividade, decreasing = TRUE), , drop = FALSE]
if (imprimir) {
cat("Execucoes em ", diretorio_saida, ":\n", sep = "")
print(out[, setdiff(names(out), "sufixo")], row.names = FALSE)
}
invisible(out)
}
## ---------------------------------------------------------------------------
## Funcao principal
## ---------------------------------------------------------------------------
etl_sia <- function(
cmp_min,
cmp_max,
uf = NULL,
cid = NULL,
sigtap = NULL,
subsistema = NULL,
diretorio_entrada,
diretorio_saida = NULL,
manter_dbc = TRUE,
enriquecer = TRUE,
converter_cns = TRUE,
manter_cns_original = FALSE,
modo = c("perguntar", "continuar", "novo"),
execucao = NULL,
limite_memoria = 0.75,
limite_disco = 0.10,
url_ftp = "ftp://ftp.datasus.gov.br/dissemin/publicos/SIASUS/200801_/Dados/",
url_pcdt = "https://raw.githubusercontent.com/LabXiabr/bd_geral/refs/heads/main/td_diretriz_cuidado.csv",
chunk = NULL,
verbose = TRUE
) {
## -------------------------------------------------------------------------
## 1. Validacao e normalizacao dos parametros
## -------------------------------------------------------------------------
modo <- match.arg(modo)
cmp_min <- as.character(cmp_min); cmp_max <- as.character(cmp_max)
if (!grepl("^[0-9]{6}$", cmp_min) || !grepl("^[0-9]{6}$", cmp_max))
stop("cmp_min e cmp_max devem estar no formato YYYYMM.")
if (cmp_min > cmp_max) stop("cmp_min nao pode ser maior que cmp_max.")
cmp_seq <- format(seq(as.Date(paste0(substr(cmp_min,1,4), "-", substr(cmp_min,5,6), "-01")),
as.Date(paste0(substr(cmp_max,1,4), "-", substr(cmp_max,5,6), "-01")),
by = "month"), "%Y%m")
ufs_validas <- c("AC","AL","AP","AM","BA","CE","DF","ES","GO","MA","MT","MS",
"MG","PA","PB","PR","PE","PI","RJ","RN","RS","RO","RR","SC","SP","SE","TO")
if (is.null(uf) || !length(uf) || all(!nzchar(trimws(uf)))) uf <- ufs_validas
uf <- toupper(trimws(as.character(uf)))
if (any(!uf %in% ufs_validas)) stop("UF invalida: ", paste(uf[!uf %in% ufs_validas], collapse = ", "))
subs_default <- c("AB","ABO","ACF","AD","AM","AMP","AN","AQ","AR","ATD","SAD")
if (is.null(subsistema) || !length(subsistema) || all(!nzchar(trimws(subsistema)))) subsistema <- subs_default
subsistema <- toupper(trimws(as.character(subsistema)))
if (any(!subsistema %in% subs_default)) {
stop("Subsistema(s) nao suportado(s): ", paste(subsistema[!subsistema %in% subs_default], collapse = ", "),
". Use: ", paste(subs_default, collapse = ", "))
}
if (is.null(cid) || !length(cid) || all(!nzchar(trimws(cid)))) {
cid <- character(0)
} else {
cid <- toupper(gsub("[^A-Z0-9]", "", trimws(as.character(cid))))
if (any(!nchar(cid) %in% c(3L,4L))) stop("Cada CID deve ter 3 ou 4 caracteres.")
}
if (is.null(sigtap) || !length(sigtap) || all(!nzchar(trimws(as.character(sigtap))))) {
sigtap <- character(0)
} else {
sigtap <- gsub("\\D", "", as.character(sigtap))
sigtap <- sub("^0+", "", sigtap)
sigtap[sigtap == ""] <- "0"
if (any(nchar(sigtap) > 9L)) stop("SIGTAP deve ter no maximo 9 digitos.")
}
if (!dir.exists(diretorio_entrada)) stop("Diretorio de entrada inexistente: ", diretorio_entrada)
diretorio_entrada <- normalizePath(diretorio_entrada, winslash = "/", mustWork = TRUE)
if (is.null(diretorio_saida) || !nzchar(diretorio_saida))
diretorio_saida <- file.path(diretorio_entrada, "sia_parquet")
dir.create(diretorio_saida, recursive = TRUE, showWarnings = FALSE)
diretorio_saida <- normalizePath(diretorio_saida, winslash = "/", mustWork = TRUE)
## Arquivos temporarios ficam dentro da entrada, evitando /tmp.
dir_tmp <- file.path(diretorio_entrada, "tmp_sia")
dir.create(dir_tmp, recursive = TRUE, showWarnings = FALSE)
Sys.setenv(TMPDIR = dir_tmp, TMP = dir_tmp, TEMP = dir_tmp)
## -------------------------------------------------------------------------
## 2. Execucao: continuar de onde parou ou criar nova (sem sobrescrever)
## -------------------------------------------------------------------------
existentes <- status_etl_sia(diretorio_saida, imprimir = FALSE)
ha_execucao <- is.data.frame(existentes) && nrow(existentes) > 0L
carimbo <- format(Sys.time(), "%Y-%m-%d-%H%M")
escolher_sufixo <- function() {
if (!ha_execucao) return("")
if (!is.null(execucao) && nzchar(execucao)) {
alvo <- if (identical(execucao, "(inicial)")) "" else paste0("_", sub("^_", "", execucao))
if (!alvo %in% existentes$sufixo)
stop("Execucao inexistente: ", execucao,
". Disponiveis: ", paste(existentes$execucao, collapse = ", "))
return(alvo)
}
if (identical(modo, "novo")) return(paste0("_", carimbo))
if (identical(modo, "continuar")) return(existentes$sufixo[1L])
## modo == "perguntar"
cat("\nJa existe saida em ", diretorio_saida, ".\n", sep = "")
status_etl_sia(diretorio_saida, imprimir = TRUE)
if (!interactive()) {
stop("Execucao nao interativa: informe modo = \"continuar\" (retoma a mais recente), ",
"modo = \"novo\" (cria saida com sufixo ", carimbo,
") ou execucao = \"<identificador>\".")
}
opcoes <- c(sprintf("Continuar a execucao mais recente [%s]", existentes$execucao[1L]),
sprintf("Criar nova execucao [%s]", carimbo))
if (nrow(existentes) > 1L)
opcoes <- c(opcoes, sprintf("Continuar outra execucao [%s]",
paste(existentes$execucao[-1L], collapse = ", ")))
escolha <- utils::menu(opcoes, title = "O que deseja fazer?")
if (escolha == 0L) stop("Execucao cancelada pelo usuario.")
if (escolha == 1L) return(existentes$sufixo[1L])
if (escolha == 2L) return(paste0("_", carimbo))
alvo <- readline("Informe o identificador da execucao: ")
alvo <- if (identical(trimws(alvo), "(inicial)")) "" else paste0("_", sub("^_", "", trimws(alvo)))
if (!alvo %in% existentes$sufixo) stop("Execucao inexistente: ", alvo)
alvo
}
sufixo <- escolher_sufixo()
retomando <- ha_execucao && sufixo %in% existentes$sufixo
if (verbose) {
cat(if (retomando) "Retomando execucao: " else "Nova execucao: ",
if (nzchar(sufixo)) sub("^_", "", sufixo) else "(inicial)", "\n", sep = "")
}
## -------------------------------------------------------------------------
## 3. Memoria, chunk e espaco em disco
## -------------------------------------------------------------------------
memoria_disponivel <- function() {
if (.Platform$OS.type == "windows") {
# fallback conservador no Windows
NA_real_
} else {
x <- tryCatch(readLines("/proc/meminfo"), error = function(e) character())
z <- grep("^MemAvailable:", x, value = TRUE)
if (!length(z)) return(NA_real_)
as.numeric(sub(".*:[[:space:]]*([0-9]+).*", "\\1", z)) * 1024
}
}
espaco_livre <- function(path) {
path <- normalizePath(path, winslash = "/", mustWork = FALSE)
if (.Platform$OS.type == "windows") return(NA_real_)
z <- tryCatch(system2("df", c("-Pk", shQuote(path)), stdout = TRUE, stderr = FALSE),
error = function(e) character())
if (length(z) < 2) return(NA_real_)
p <- strsplit(trimws(z[length(z)]), "[[:space:]]+")[[1]]
if (length(p) < 4) return(NA_real_)
as.numeric(p[4]) * 1024
}
checar_disco <- function() {
livre <- espaco_livre(diretorio_saida)
if (is.na(livre)) return(invisible(TRUE))
total <- tryCatch({
z <- system2("df", c("-Pk", shQuote(diretorio_saida)), stdout = TRUE, stderr = FALSE)
p <- strsplit(trimws(z[length(z)]), "[[:space:]]+")[[1]]
as.numeric(p[2]) * 1024
}, error = function(e) NA_real_)
if (!is.na(total) && livre / total <= limite_disco)
stop("ETL interrompido: espaco livre em disco <= ", limite_disco * 100, "%.")
invisible(TRUE)
}
if (is.null(chunk)) {
mem <- memoria_disponivel()
chunk <- if (is.na(mem)) 50000L else max(5000L, min(500000L, as.integer(mem * limite_memoria / 2500)))
}
chunk <- as.integer(chunk)
## -------------------------------------------------------------------------
## 4. Tabela de PCDT
## -------------------------------------------------------------------------
baixar_texto <- function(url) {
texto_para_utf8(rawToChar(curl::curl_fetch_memory(url)$content))
}
pcdt <- tryCatch({
txt <- baixar_texto(url_pcdt)
utils::read.csv(text = txt, stringsAsFactors = FALSE, check.names = FALSE,
na.strings = c("", "NA"), quote = "\"")
}, error = function(e) {
stop("Nao foi possivel ler td_diretriz_cuidado.csv: ", conditionMessage(e))
})
if (!all(c("co_cid", "sg_pcdt") %in% names(pcdt)))
stop("td_diretriz_cuidado.csv precisa conter co_cid e sg_pcdt.")
## Retorna os sg_pcdt correspondentes ao CID da linha.
cache_pcdt <- new.env(parent = emptyenv())
mapear_pcdt <- function(cids) {
if (!length(cids)) return(NA_character_)
cids <- toupper(trimws(cids))
## Prefixo evita nome de variavel vazio quando o CID vem em branco.
chave <- paste0("k:", paste(cids, collapse = "|"))
if (!is.null(cache_pcdt[[chave]])) return(cache_pcdt[[chave]])
out <- character(0)
for (x in cids[nzchar(cids)]) {
hit <- pcdt[grepl(paste0("(^|,)", x, "(,|$)"), pcdt$co_cid, ignore.case = TRUE), "sg_pcdt"]
if (length(hit)) out <- c(out, hit)
## Para CID de 3 digitos, aceitar tambem os codigos de 4 que iniciem pelo prefixo.
if (nchar(x) == 3L) {
hit2 <- pcdt$sg_pcdt[vapply(strsplit(pcdt$co_cid, ",", fixed = TRUE),
function(v) any(startsWith(toupper(trimws(v)), x)), logical(1))]
out <- c(out, hit2)
}
}
res <- unique(out[nzchar(out)])
cache_pcdt[[chave]] <- res
res
}
sg_pcdt <- if (length(cid)) unique(unlist(lapply(cid, mapear_pcdt), use.names = FALSE)) else "TODOS"
sg_pcdt <- if (!length(sg_pcdt)) "SEM_PCDT" else sg_pcdt
## -------------------------------------------------------------------------
## 5. Leitor DBF em blocos (com iconv latin1 -> UTF-8)
## -------------------------------------------------------------------------
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)]))
nomes <- toupper(para_utf8(nomes))
lens <- as.integer(fd[17, ])
ini <- 2L + c(0L, cumsum(lens))[seq_len(nf)]
list(n_rec=n_rec, h_len=h_len, r_len=r_len, nomes=nomes, ini=ini, fim=ini+lens-1L)
}
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) {
checar_disco()
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) break
length(blk) <- n * hd$r_len
blk[blk == as.raw(0)] <- as.raw(32)
s <- rawToChar(blk); rm(blk); Encoding(s) <- "latin1"
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]
## O recorte e feito sobre o texto marcado como latin1 (posicoes do DBF
## sao em bytes) e a conversao vem ANTES de qualquer funcao de string:
## trimws()/gsub() traduziriam implicitamente para UTF-8 e um iconv
## posterior produziria dupla codificacao (JOAO -> JOA~O).
para_utf8(substring(s, b+hd$ini[j], b+hd$fim[j]))
}
FUN(getcol, validos, hd$nomes)
lidos <- lidos + n
rm(s); gc(FALSE)
}
invisible(hd$n_rec)
}
## -------------------------------------------------------------------------
## 6. FTP / arquivos locais
## -------------------------------------------------------------------------
texto_ftp <- texto_para_utf8(rawToChar(curl::curl_fetch_memory(url_ftp)$content))
arquivos_ftp <- unique(regmatches(texto_ftp, gregexpr("[[:alnum:]_-]+\\.[dD][bB][cC]", texto_ftp))[[1]])
arquivos_ftp <- arquivos_ftp[nzchar(arquivos_ftp)]
## Aceita PASP2301.DBC, PASP2301A.DBC, PASP2301_1.DBC, PASP23011.DBC etc.
listar_partes <- function(sub, u, cmp) {
yy <- substr(cmp, 3, 4); mm <- substr(cmp, 5, 6)
padrao <- sprintf("^%s%s%s%s([A-Z]|_?[0-9]+)?\\.DBC$", sub, u, yy, mm)
arquivos_ftp[grepl(padrao, toupper(arquivos_ftp))]
}
preparar_dbf <- function(arq) {
local_dbc <- file.path(diretorio_entrada, arq)
if (!file.exists(local_dbc) || file.size(local_dbc) == 0) {
ok <- tryCatch({
curl::curl_download(
url = paste0(url_ftp, arq),
destfile = local_dbc,
quiet = TRUE
)
TRUE
}, error = function(e) {
warning(
"Falha no download de ", arq, ": ",
conditionMessage(e)
)
FALSE
})
if (!isTRUE(ok) ||
!file.exists(local_dbc) ||
file.size(local_dbc) == 0) {
if (file.exists(local_dbc) &&
file.size(local_dbc) == 0) {
unlink(local_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 <- tryCatch({
read.dbc::dbc2dbf(
local_dbc,
dest_dbf
)
TRUE
}, error = function(e) {
warning(
"Falha ao converter ", arq, ": ",
conditionMessage(e)
)
FALSE
})
if (!isTRUE(ok) ||
!file.exists(dest_dbf) ||
file.size(dest_dbf) == 0) {
if (file.exists(dest_dbf)) {
unlink(dest_dbf)
}
return(NULL)
}
}
list(
dbc = local_dbc,
dbf = dest_dbf
)
}
limpar_temporarios <- function(prep) {
if (is.null(prep)) return(invisible(NULL))
if (!manter_dbc && file.exists(prep$dbc)) unlink(prep$dbc)
if (file.exists(prep$dbf)) unlink(prep$dbf)
invisible(NULL)
}
## -------------------------------------------------------------------------
## 7. Filtros
## -------------------------------------------------------------------------
match_prefix <- function(x, filtros) {
if (!length(filtros)) return(rep(TRUE, length(x)))
x <- trimws(as.character(x))
Reduce(`|`, lapply(filtros, function(f) startsWith(x, f)))
}
match_cid <- function(c1, c2) {
if (!length(cid)) return(rep(TRUE, length(c1)))
a <- toupper(substr(trimws(c1), 1L, 4L)); b <- toupper(substr(trimws(c2), 1L, 4L))
Reduce(`|`, lapply(cid, function(f) {
if (nchar(f) == 3L) startsWith(a, f) | startsWith(b, f)
else a == f | b == f
}))
}
match_sigtap <- function(x) {
if (!length(sigtap)) return(rep(TRUE, length(x)))
x <- gsub("\\D", "", trimws(as.character(x))); x <- sub("^0+", "", x)
x[x == ""] <- "0"
Reduce(`|`, lapply(sigtap, function(f) if (nchar(f) < 9L) startsWith(x, f) else x == f))
}
## -------------------------------------------------------------------------
## 8. Estado incremental
## -------------------------------------------------------------------------
estado_dir <- file.path(diretorio_saida, paste0("_etl_sia_estado", sufixo))
idx_dir <- file.path(estado_dir, "indice_auxiliares")
dir.create(estado_dir, recursive=TRUE, showWarnings=FALSE)
dir.create(idx_dir, recursive=TRUE, showWarnings=FALSE)
nome_base <- function(sg) sprintf("sia_pa_%s_%s_%s%s", sg, substr(cmp_min,3,6), substr(cmp_max,3,6), sufixo)
## Cada sg_pcdt tem um dataset Parquet independente. O nome com .parquet e
## uma pasta-dataset: isso permite append incremental sem reescrever todo o arquivo.
dataset_dir <- function(sg) file.path(diretorio_saida, paste0(nome_base(sg), ".parquet"))
status_file <- function(cmp, uf, sub, arq) file.path(estado_dir, paste0(cmp,"_",uf,"_",sub,"_",arq,".done"))
marcar_done <- function(path) {
dir.create(
dirname(path),
recursive = TRUE,
showWarnings = FALSE
)
file.create(path)
}
## -------------------------------------------------------------------------
## 9. Indice dos subsistemas auxiliares (AP_AUTORIZ -> campos especificos)
## -------------------------------------------------------------------------
## Resolve, contra o cabecalho real do DBF, quais campos serao levados ao PA.
mapa_campos <- function(sub, nomes) {
mapa <- campos_subsistema(sub)
if (!length(mapa)) {
alvo <- grep(paste0("^", sub, "_"), nomes, value = TRUE)
alvo <- setdiff(alvo, "AP_AUTORIZ")
mapa <- setNames(tolower(alvo), alvo)
}
mapa <- mapa[names(mapa) %in% nomes]
## AP_CNSPCN e comum a todos os auxiliares: recebe prefixo do subsistema.
if ("AP_CNSPCN" %in% nomes)
mapa <- c(mapa, setNames(paste0(tolower(sub), "_ap_cnspcn"), "AP_CNSPCN"))
mapa
}
## Um arquivo Parquet por (competencia, UF, subsistema). Construido uma unica
## vez e reaproveitado por todos os chunks do PA e por execucoes retomadas.
construir_indice <- function(cmp, u, sub) {
destino <- file.path(idx_dir, sprintf("%s_%s_%s.parquet", cmp, u, sub))
if (file.exists(destino)) return(destino)
arqs <- listar_partes(sub, u, cmp)
if (!length(arqs)) return(NULL)
partes <- list()
for (arq in arqs) {
checar_disco()
prep <- preparar_dbf(arq)
if (is.null(prep)) { warning("Falha ao preparar auxiliar ", arq); next }
dbf_ler_chunks(prep$dbf, function(getcol, validos, nomes) {
if (!"AP_AUTORIZ" %in% nomes) return(invisible(NULL))
aut <- getcol("AP_AUTORIZ")
ok <- which(!is.na(aut) & nzchar(aut))
if (!length(ok)) return(invisible(NULL))
mapa <- mapa_campos(sub, nomes)
d <- list(AP_AUTORIZ = aut[ok])
for (nm in names(mapa)) {
v <- getcol(nm)[ok]
if (isTRUE(converter_cns) && identical(nm, "AP_CNSPCN")) {
if (isTRUE(manter_cns_original)) d[[paste0(mapa[[nm]], "_bruto")]] <- v
v <- converte_cns(v)
}
d[[ mapa[[nm]] ]] <- v
}
partes[[length(partes) + 1L]] <<- as.data.frame(d, stringsAsFactors = FALSE)
invisible(NULL)
})
limpar_temporarios(prep)
gc(FALSE)
}
if (!length(partes)) return(NULL)
ind <- do.call(rbind, partes); rm(partes)
## Uma APAC pode ter mais de um laudo; mantem-se o primeiro para que o join
## nao multiplique linhas do PA.
ind <- ind[!duplicated(ind$AP_AUTORIZ), , drop = FALSE]
tmp <- paste0(destino, ".tmp")
arrow::write_parquet(ind, tmp, compression = "zstd")
file.rename(tmp, destino)
rm(ind); gc(FALSE)
destino
}
## Colunas auxiliares esperadas na saida. Sao emitidas mesmo quando o
## subsistema nao existe para aquela UF/competencia (preenchidas com NA),
## para que o dataset Parquet mantenha um esquema estavel.
cols_aux <- new.env(parent = emptyenv())
for (s0 in subsistema) {
mapa0 <- campos_subsistema(s0)
base0 <- c(unname(mapa0), paste0(tolower(s0), "_ap_cnspcn"))
if (isTRUE(manter_cns_original))
base0 <- c(base0, paste0(tolower(s0), "_ap_cnspcn_bruto"))
assign(s0, unique(base0), envir = cols_aux)
}
carregar_indices <- function(cmp, u) {
if (!isTRUE(enriquecer)) return(list())
out <- list()
for (sub in subsistema) {
p <- tryCatch(construir_indice(cmp, u, sub), error = function(e) {
warning("Indice ", sub, " ", u, cmp, ": ", conditionMessage(e)); NULL
})
if (is.null(p) || !file.exists(p)) next
ind <- as.data.frame(arrow::read_parquet(p), stringsAsFactors = FALSE)
if (!nrow(ind)) next
out[[sub]] <- ind
assign(sub, unique(c(get(sub, envir = cols_aux), setdiff(names(ind), "AP_AUTORIZ"))),
envir = cols_aux)
if (verbose) cat(sprintf(" indice %-3s: %d APAC\n", sub, nrow(ind)))
}
out
}
## Acrescenta ao data.frame do PA as colunas dos auxiliares.
enriquecer_pa <- function(d, indices) {
if (!isTRUE(enriquecer)) return(d)
chave <- trimws(as.character(d$PA_AUTORIZ))
for (sub in subsistema) {
ind <- indices[[sub]]
m <- if (is.null(ind)) rep(NA_integer_, nrow(d)) else match(chave, ind$AP_AUTORIZ)
for (cl in get(sub, envir = cols_aux)) {
## Colisao de nome com o proprio PA recebe o prefixo do subsistema.
nome_final <- if (cl %in% names(d)) paste0(tolower(sub), "_", cl) else cl
d[[nome_final]] <- if (!is.null(ind) && cl %in% names(ind)) ind[[cl]][m] else NA_character_
}
d[[paste0(tolower(sub), "_st_vinculado")]] <- !is.na(m)
}
d
}
## -------------------------------------------------------------------------
## 10. ETL
## -------------------------------------------------------------------------
total <- 0L
for (cmp in cmp_seq) {
for (u in uf) {
arqs_pa <- listar_partes("PA", u, cmp)
if (!length(arqs_pa)) next
if (verbose) cat(sprintf("%s | %s | PA=%d\n", cmp, u, length(arqs_pa)))
pendentes <- arqs_pa[!vapply(arqs_pa, function(a) file.exists(status_file(cmp,u,"PA",a)), logical(1))]
if (!length(pendentes)) next
## Indice dos auxiliares desta competencia/UF, construido antes do PA.
indices <- carregar_indices(cmp, u)
for (sub in names(indices)) marcar_done(status_file(cmp, u, sub, "INDICE"))
for (arq in pendentes) {
checar_disco()
prep <- preparar_dbf(arq)
if (is.null(prep)) { warning("Falha ao preparar ", arq); next }
## Para cada bloco, somente registros que passam nos filtros viram data.frame.
## Partes tem nomes deterministicos: se o processo cair, o arquivo e limpo
## no reinicio e nao ha duplicacao.
chunk_id <- 0L
for (sg0 in sg_pcdt) {
dd <- dataset_dir(sg0)
if (dir.exists(dd)) {
olds <- list.files(dd, pattern=paste0("^", cmp, "_", u, "_", tools::file_path_sans_ext(arq), "_"), full.names=TRUE)
if (length(olds)) unlink(olds)
}
}
dbf_ler_chunks(prep$dbf, function(getcol, validos, nomes) {
chunk_id <<- chunk_id + 1L
cid1 <- getcol("PA_CIDPRI")
cid2 <- getcol("PA_CIDSEC")
proc <- getcol("PA_PROC_ID")
hit <- which(match_cid(cid1,cid2) & match_sigtap(proc))
if (!length(hit)) return(invisible(NULL))
linhas <- validos[hit]
d <- lapply(nomes, function(nm) getcol(nm, linhas)); names(d) <- nomes
d <- as.data.frame(d, stringsAsFactors=FALSE)
## CNS do profissional tambem chega ofuscado no PA.
if (isTRUE(converter_cns) && "PA_CNSMED" %in% names(d)) {
if (isTRUE(manter_cns_original)) d$PA_CNSMED_BRUTO <- d$PA_CNSMED
d$PA_CNSMED <- converte_cns(d$PA_CNSMED)
}
d$UF_ETL <- u
d$CMP_ETL <- cmp
## Enriquecimento PA_AUTORIZ <-> AP_AUTORIZ, mesmo UF + AAMM.
d <- enriquecer_pa(d, indices)
## Determina PCDT usando CID primario/secundario. Quando nenhum CID foi
## informado, preserva TODOS; quando foi informado mas nao ha correspondencia,
## usa SEM_PCDT.
p1 <- lapply(d$PA_CIDPRI, mapear_pcdt)
p2 <- lapply(d$PA_CIDSEC, mapear_pcdt)
pcdt_row <- mapply(function(a,b) unique(c(a,b)), p1, p2, SIMPLIFY=FALSE)
if (!length(cid)) pcdt_row <- replicate(nrow(d), "TODOS", simplify=FALSE)
if (length(cid)) pcdt_row <- lapply(pcdt_row, function(x) if(length(x)) x else "SEM_PCDT")
## Um registro pode pertencer a mais de um PCDT quando o CID de 3 digitos
## sobrepoe mais de uma diretriz. Nesse caso ele e escrito em cada dataset.
for (sg in unique(unlist(pcdt_row, use.names=FALSE))) {
ix <- which(vapply(pcdt_row, function(x) sg %in% x, logical(1)))
if (!length(ix)) next
ds <- dataset_dir(sg)
dir.create(ds, recursive=TRUE, showWarnings=FALSE)
part <- file.path(ds, sprintf("%s_%s_%s_%s_%06d.parquet", cmp, u, tools::file_path_sans_ext(arq), sg, chunk_id))
arrow::write_parquet(d[ix,,drop=FALSE], part, compression="zstd")
total <<- total + length(ix)
checar_disco()
}
invisible(NULL)
})
limpar_temporarios(prep)
marcar_done(status_file(cmp,u,"PA",arq))
gc(FALSE)
}
rm(indices); gc(FALSE)
}
}
cat("ETL concluida. Registros gravados: ", total, "\n", sep="")
cat("Execucao: ", if (nzchar(sufixo)) sub("^_", "", sufixo) else "(inicial)", "\n", sep="")
cat("Saida: ", diretorio_saida, "\n", sep="")
invisible(list(
diretorio_saida=diretorio_saida,
execucao=if (nzchar(sufixo)) sub("^_", "", sufixo) else "(inicial)",
retomada=retomando,
datasets=unique(vapply(sg_pcdt, dataset_dir, character(1))),
competencias=cmp_seq,
uf=uf,
subsistema=subsistema,
enriquecer=enriquecer,
total=total,
chunk=chunk,
pcdt=sg_pcdt
))
}Se já existir saída no diretorio_saida, a função
não sobrescreve. Ela exibe o status e pergunta o que
fazer:
Ja existe saida em /meudiscogrande/dbc/sia_parquet.
Execucoes em /meudiscogrande/dbc/sia_parquet:
execucao arquivos_concluidos datasets partes_parquet mb ultima_atividade
(inicial) 812 3 4193 9120 2026-08-18 03:11
O que deseja fazer?
1: Continuar a execucao mais recente [(inicial)]
2: Criar nova execucao [2026-08-20-2339]
.done
e o índice dos auxiliares; nada é reprocessado.YYYY-MM-DD-HHMM ao nome dos datasets e ao diretório de
estado:sia_pa_DINS_2001_2606.parquet <- execução inicial
sia_pa_DINS_2001_2606_2026-08-20-2339.parquet <- nova execução
_etl_sia_estado
_etl_sia_estado_2026-08-20-2339
Em execução não interativa (Rscript), a
função interrompe e exige a decisão explícita por parâmetro:
modo = "continuar" — retoma a execução mais
recente;modo = "novo" — cria uma nova com carimbo;execucao = "2026-08-20-2339" — retoma uma execução
específica.Para consultar o status sem executar nada:
O script trabalha em chunks e descarta os objetos depois da gravação,
para evitar acumular o resultado em memória. O tamanho do bloco é
derivado de MemAvailable
(limite_memoria = 0.75) e a execução é interrompida quando
o espaço livre em disco cai abaixo de limite_disco (10% por
padrão).
Cada arquivo concluído recebe uma marca .done no
diretório de estado; ao reiniciar, as partes Parquet do arquivo em curso
são apagadas e regravadas, o que garante retomada sem duplicação.
etl_sia(
cmp_min, # "AAAAMM"
cmp_max, # "AAAAMM"
uf = NULL, # NULL = todas as UFs
cid = NULL, # NULL = sem filtro de CID
sigtap = NULL, # NULL = sem filtro de procedimento
subsistema = NULL, # NULL = todos os auxiliares
diretorio_entrada, # cache dos .DBC
diretorio_saida = NULL, # padrão: <entrada>/sia_parquet
manter_dbc = TRUE,
enriquecer = TRUE, # cruzamento PA_AUTORIZ <-> AP_AUTORIZ
converter_cns = TRUE, # desofusca AP_CNSPCN e PA_CNSMED
manter_cns_original = FALSE, # preserva o valor bruto
modo = c("perguntar", "continuar", "novo"),
execucao = NULL, # retomar uma execução específica
limite_memoria = 0.75,
limite_disco = 0.10,
url_ftp = "ftp://ftp.datasus.gov.br/dissemin/publicos/SIASUS/200801_/Dados/",
url_pcdt = "https://raw.githubusercontent.com/LabXiabr/bd_geral/refs/heads/main/td_diretriz_cuidado.csv",
chunk = NULL,
verbose = TRUE
)Para executar AM + AQ + AR, no período 202001 a 202606, e considerando os CIDs que você passou:
source("etl_sia.R")
etl_sia(
cmp_min = "202001",
cmp_max = "202606",
uf = NULL,
cid = c(
"E232", # Diabetes Insípido
"D500", "D508", # Anemia por Deficiência de Ferro
"D841" # Angioedema
),
sigtap = NULL,
subsistema = c("AM", "AQ", "AR"),
diretorio_entrada = "/meudiscogrande/dbc/",
diretorio_saida = NULL,
manter_dbc = TRUE
)Como você não especificou UF, uf = NULL fará o
processamento para todas as UFs.
Em Rscript não há prompt: informe modo
explicitamente.
Os datasets já saem separados por sg_pcdt
(DINS, FER, ANGI), derivado do
td_diretriz_cuidado.csv; uma execução única com os quatro
CIDs já produz os três datasets sem misturar resultados. Rodar
separadamente continua possível:
Diabetes Insípido — E232 →
DINS:
etl_sia(
cmp_min = "202001",
cmp_max = "202606",
uf = NULL,
cid = "E232",
sigtap = NULL,
subsistema = c("AM", "AQ", "AR"),
diretorio_entrada = "/meudiscogrande/dbc/",
diretorio_saida = NULL,
manter_dbc = TRUE,
modo = "continuar"
)Anemia por Deficiência de Ferro — D500,
D508 → FER:
etl_sia(
cmp_min = "202001",
cmp_max = "202606",
uf = NULL,
cid = c("D500", "D508"),
sigtap = NULL,
subsistema = c("AM", "AQ", "AR"),
diretorio_entrada = "/meudiscogrande/dbc/",
diretorio_saida = NULL,
manter_dbc = TRUE,
modo = "continuar"
)Angioedema — D841 →
ANGI:
library(arrow); library(dplyr)
ds <- open_dataset("/meudiscogrande/dbc/sia_parquet/sia_pa_DINS_2001_2606.parquet")
ds %>%
filter(am_st_vinculado) %>%
select(PA_AUTORIZ, PA_CIDPRI, st_gestante, nu_peso, nu_altura,
aq_cid10, aq_estadi, am_ap_cnspcn) %>%
collect()Taxa de vínculo por subsistema, útil para validar o cruzamento: