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.

Visão geral

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 de latin1 para UTF-8 e a função pergunta se é para continuar uma execução anterior ou criar uma nova.

Subsistemas do SIA

  • PA é sempre a base do dataset de saída.
  • Subsistemas auxiliares aceitos em subsistema =: AB, ABO, ACF, AD, AM, AMP, AN, AQ, AR, ATD, SAD.
  • Todos os auxiliares são cruzados por PA_AUTORIZ ↔︎ AP_AUTORIZ.
  • O arquivo auxiliar precisa corresponder ao mesmo UF + AAMM do PA.
  • Se subsistema = NULL, todos os auxiliares são considerados.
  • enriquecer = FALSE desliga o cruzamento e volta ao comportamento antigo.

Particionamento dos DBC

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

Enriquecimento

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.

Índice incremental dos auxiliares

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.

Campos específicos

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().

Nomes duplicados

Como você pode selecionar, por exemplo:

subsistema = c("AM", "AQ", "AR")

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.

Desofuscação do CNS

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.

Codificação de caracteres

Os DBF do DATASUS estão em latin1. Todo valor extraído passa por:

iconv(x, from = "latin1", to = "UTF-8", sub = "")

Dois cuidados que o código já resolve:

  1. A conversão vem antes de qualquer função de string. 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.
  2. Documentos que já chegam em UTF-8 — a listagem do FTP e o td_diretriz_cuidado.csv — passam por texto_para_utf8(), que só converte quando os bytes não são UTF-8 válidos.

Execuções: continuar ou criar nova

###############################################################################
# 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]
  • Continuar reaproveita as marcas .done e o índice dos auxiliares; nada é reprocessado.
  • Criar nova acrescenta o carimbo 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:

source("etl_sia.R")
status_etl_sia("/meudiscogrande/dbc/sia_parquet")

Memória e retomada

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.

Assinatura

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
)

Exemplos de Execução

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.

Se quiser executar pelo terminal

Em Rscript não há prompt: informe modo explicitamente.

Rscript -e 'source("etl_sia.R"); etl_sia(
  cmp_min="202001",
  cmp_max="202606",
  uf=NULL,
  cid=c("E232","D500","D508","D841"),
  sigtap=NULL,
  subsistema=c("AM","AQ","AR"),
  diretorio_entrada="/meudiscogrande/dbc/",
  diretorio_saida=NULL,
  manter_dbc=TRUE,
  modo="continuar"
)'

Para cada PCDT separadamente

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 — E232DINS:

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, D508FER:

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 — D841ANGI:

etl_sia(
  cmp_min = "202001",
  cmp_max = "202606",
  uf = NULL,
  cid = "D841",
  sigtap = NULL,
  subsistema = c("AM", "AQ", "AR"),
  diretorio_entrada = "/meudiscogrande/dbc/",
  diretorio_saida = NULL,
  manter_dbc = TRUE,
  modo = "continuar"
)

Leitura do resultado

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:

ds %>%
  summarise(
    n = n(),
    am = mean(am_st_vinculado),
    aq = mean(aq_st_vinculado),
    ar = mean(ar_st_vinculado)
  ) %>%
  collect()

Limitações conhecidas

  • Uma APAC com múltiplos laudos no auxiliar contribui apenas com o primeiro registro. Se a análise exigir todos, o índice precisa ser lido separadamente e agregado.
  • O índice de uma competência/UF é carregado inteiro em memória durante o PA daquela competência/UF. Isso é pequeno em relação ao PA, mas cresce com o número de subsistemas selecionados.
  • Todas as colunas são gravadas como texto; conversões de tipo (datas, numéricos) ficam a cargo da etapa de análise.