--- title: "Pipeline Completo com DuckDB" output: rmarkdown::html_vignette vignette: > %\VignetteIndexEntry{Pipeline Completo com DuckDB} %\VignetteEngine{knitr::rmarkdown} %\VignetteEncoding{UTF-8} --- ```{r setup, include = FALSE} knitr::opts_chunk$set( collapse = TRUE, comment = "#>", eval = FALSE ) ``` Este guia demonstra o pipeline completo do `datacaged`: desde a verificação do repositório HuggingFace até consultas analíticas no DuckDB com `dplyr` e SQL. ## 1. Verificar disponibilidade antes de baixar Antes de iniciar qualquer download, verifique se o repositório HuggingFace está acessível e quais competências estão disponíveis. ```{r status} library(datacaged) # Verifica se o repositório HuggingFace está online e mede a latência caged_status() #> OK Repositório HuggingFace online #> ℹ Latência: 312 ms #> ℹ Dataset: https://huggingface.co/datasets/alexsandroprado/caged ``` ```{r ftp-files} # Lista os 12 months mais recentes disponíveis no Novo CAGED caged_hf_files() # Lista todas as competências do CAGED antigo caged_hf_files(type = "antigo", n = Inf) # Lista competências do CAGED Ajustes caged_hf_files(type = "ajustes", n = 24) # Capturar o resultado para uso programático disponivel <- caged_hf_files(n = 1, verbose = FALSE) cat("Competência mais recente:", disponivel$competencia, "\n") ``` ## 2. Download dos microdados O `caged_download()` gerencia cache local automaticamente: arquivos já baixados não são re-baixados. O retorno é um `manifest` — data.frame com o status de cada arquivo. ```{r download-novo} # Baixa Novo CAGED de 2023 inteiro (MOV + FOR + EXC por competência) manifest <- caged_download( years = 2023, months = seq_len(12L), destdir = "~/dados/caged_cache" # omitir para usar cache padrão ) # Inspecionar o manifest dplyr::count(manifest, type, status) #> # A tibble: 4 × 3 #> type status n #> #> 1 EXC baixado 12 #> 2 FOR baixado 12 #> 3 MOV baixado 12 #> 4 MOV cache 0 # Arquivos disponíveis para processar manifest |> dplyr::filter(status %in% c("baixado", "cache")) |> dplyr::select(type, competencia, nome_arquivo, arquivo) ``` ```{r download-antigo} # Baixa CAGED antigo — arquivo nacional único por competência manifest_ant <- caged_download(years = 2018:2019, months = seq_len(12L)) # Baixar também os ajustes do mesmo período caged_adjustments_load( years = 2018:2019, months = seq_len(12L), db_path = "caged_historico.duckdb" ) ``` ## 3. Parse dos arquivos O `caged_parse()` extrai o `.7z`, lê o `.txt` e devolve um tibble normalizado. ```{r parse-individual} # Parsear um único arquivo df_jan <- caged_parse("~/dados/caged_cache/NOVO_CAGED/2023/202301/CAGEDMOV202301.7z") # Inspecionar dplyr::glimpse(df_jan) #> Rows: 3,665,155 #> Columns: 29 #> $ competenciamov 202301, 202301, ... #> $ uf 11, 11, ... #> $ municipio 1100015, ... #> $ saldomovimentacao 1, -1, ... #> $ salario 1412.00, 2800.50, ... # Parsear vários arquivos do mesmo type (devem ser todos MOV, ou todos FOR, etc.) arquivos_mov <- list.files( "~/dados/caged_cache/NOVO_CAGED/2023", pattern = "CAGEDMOV", recursive = TRUE, full.names = TRUE ) df_mov_2023 <- caged_parse_batch(arquivos_mov) ``` ## 4. Gravar no DuckDB ```{r gravar} # Gravar o tibble parseado no banco caged_to_duckdb(df_mov_2023, db_path = "caged.duckdb") # Re-runs são seguros: competências já no banco são puladas caged_to_duckdb(df_mov_2023, db_path = "caged.duckdb") #> ℹ 12 competências já no banco — pulando. # Para regravar (ex: após corrigir dados): caged_to_duckdb(df_mov_2023, db_path = "caged.duckdb", overwrite_competencies = TRUE) # Gravar numa table específica sem autodetectar caged_to_duckdb(df_jan, db_path = "caged.duckdb", table = "caged_mov") ``` ## 5. Pipeline completo em um único comando Para o caso mais comum — baixar tudo e gravar no banco — use `caged_load()`: ```{r pipeline-completo} library(datacaged) # ── Novo CAGED 2022-2023 ────────────────────────────────────────────────────── caged_load( years = 2022:2023, months = seq_len(12L), db_path = "caged.duckdb" ) # ── CAGED antigo 2015-2019 ──────────────────────────────────────────────────── caged_load( years = 2015:2019, months = seq_len(12L), db_path = "caged.duckdb" ) # ── CAGED Ajustes 2015-2019 ─────────────────────────────────────────────────── caged_adjustments_load( years = 2015:2019, months = seq_len(12L), db_path = "caged.duckdb" ) # ── Verificar o banco resultante ────────────────────────────────────────────── caged_info("caged.duckdb") #> ── caged.duckdb ────────────────────────────────────────────────────────── #> Tamanho do arquivo: 4.2 GB #> ── Tabelas ─────────────────────────────────────────────────────────────── #> * "caged_mov" Registros: 43,981,860 Competências: 202201 – 202312 #> * "caged_for" Registros: 1,093,176 Competências: 202201 – 202312 #> * "caged_exc" Registros: 94,800 Competências: 202201 – 202312 #> * "caged_antigo" Registros: 28,450,000 Competências: 201501 – 201912 #> * "caged_ajustes" Registros: 342,000 Competências: 201501 – 201912 ``` ## 6. Consultas com dplyr ```{r conexao} library(dplyr) con <- caged_connect("caged.duckdb") ``` ```{r consultas-dplyr} # ── Tabelas disponíveis ─────────────────────────────────────────────────────── DBI::dbListTables(con) # ── Referência lazy (não carrega na memória) ────────────────────────────────── mov <- tbl(con, "caged_mov") # ── Saldo mensal ────────────────────────────────────────────────────────────── mov |> group_by(competenciamov) |> summarise( admissoes = sum(saldomovimentacao == 1, na.rm = TRUE), desligamentos = sum(saldomovimentacao == -1, na.rm = TRUE), saldo = sum(saldomovimentacao, na.rm = TRUE) ) |> arrange(competenciamov) |> collect() # ── Saldo por UF e setor ────────────────────────────────────────────────────── mov |> filter(competenciamov >= 202301) |> group_by(uf, secao) |> summarise(saldo = sum(saldomovimentacao, na.rm = TRUE)) |> arrange(desc(saldo)) |> collect() # ── Salário médio por escolaridade e sexo ───────────────────────────────────── mov |> filter(!is.na(salario), salario > 0) |> group_by(graudeinstrucao, sexo) |> summarise( salario_medio = mean(salario, na.rm = TRUE), n = n() ) |> collect() # ── Top 10 municípios por admissões ────────────────────────────────────────── mov |> filter(saldomovimentacao == 1) |> group_by(municipio) |> summarise(admissoes = n()) |> slice_max(admissoes, n = 10) |> collect() ``` ## 7. Consultas com SQL direto ```{r consultas-sql} # Saldo líquido por seção da CNAE e mês DBI::dbGetQuery(con, " SELECT competenciamov AS competencia, secao, SUM(CASE WHEN saldomovimentacao = 1 THEN 1 ELSE 0 END) AS admissoes, SUM(CASE WHEN saldomovimentacao = -1 THEN 1 ELSE 0 END) AS desligamentos, SUM(saldomovimentacao) AS saldo, ROUND(AVG(salario), 2) AS salario_medio FROM caged_mov WHERE competenciamov >= 202301 GROUP BY competenciamov, secao ORDER BY competenciamov, saldo DESC ") # Distribuição etária por faixa e type de movimentação DBI::dbGetQuery(con, " SELECT CASE WHEN idade < 25 THEN 'Até 24 anos' WHEN idade BETWEEN 25 AND 34 THEN '25-34 anos' WHEN idade BETWEEN 35 AND 44 THEN '35-44 anos' WHEN idade BETWEEN 45 AND 54 THEN '45-54 anos' ELSE '55 anos ou mais' END AS faixa_etaria, SUM(CASE WHEN saldomovimentacao = 1 THEN 1 ELSE 0 END) AS admissoes, SUM(CASE WHEN saldomovimentacao = -1 THEN 1 ELSE 0 END) AS desligamentos FROM caged_mov GROUP BY faixa_etaria ORDER BY faixa_etaria ") ``` ## 8. Série histórica: unindo Novo CAGED e CAGED antigo Para construir séries longas que cruzem o ponto de ruptura de 2020, normalize as colunas antes de empilhar. ```{r serie-historica} # Novo CAGED — colunas a normalizar novo <- tbl(con, "caged_mov") |> select( competencia = competenciamov, uf, saldo = saldomovimentacao, salario, sexo, idade, escolaridade = graudeinstrucao ) |> mutate(serie = "novo") # CAGED antigo — já tem coluna competencia e saldomovimentacao antigo <- tbl(con, "caged_antigo") |> select( competencia, uf, saldo = saldomovimentacao, salario, sexo, idade, escolaridade ) |> mutate(serie = "antigo") # Empilha e agrega serie_hist <- bind_rows( novo |> group_by(competencia, serie) |> summarise(saldo = sum(saldo, na.rm = TRUE)), antigo |> group_by(competencia, serie) |> summarise(saldo = sum(saldo, na.rm = TRUE)) ) |> collect() |> arrange(competencia) ``` ## 9. Exportar resultados ```{r exportar} library(dplyr) # Exportar para CSV resultado <- tbl(con, "caged_mov") |> group_by(competenciamov, uf) |> summarise(saldo = sum(saldomovimentacao, na.rm = TRUE)) |> collect() readr::write_csv(resultado, "saldo_uf_mensal.csv") # Exportar todas as tabelas do banco para Parquet (nativo DuckDB, muito mais rápido) caged_to_parquet("caged.duckdb", output_dir = "~/exports") # Exportar só caged_mov particionado por UF caged_to_parquet( "caged.duckdb", output_dir = "~/exports", tables = "caged_mov", partition_by = "uf" ) # Fechar conexão DBI::dbDisconnect(con, shutdown = TRUE) ``` ## 10. Atualização incremental Use `caged_update()` para manter o banco atualizado sem re-baixar tudo. A função consulta o banco para descobrir a competência máxima e baixa apenas o que é mais recente no HuggingFace. ```{r update} # Atualiza o Novo CAGED com os meses ainda não presentes no banco caged_update(db_path = "caged.duckdb") # Atualizar todas as séries caged_update( db_path = "caged.duckdb", series = c("novo", "antigo", "ajustes") ) ``` ## 11. Boas práticas - **Sempre feche a conexão** com `DBI::dbDisconnect(con, shutdown = TRUE)` ao terminar. - **Use `tbl()` antes de `collect()`**: o DuckDB executa o máximo possível antes de trazer dados para o R. - **`select()` cedo**: traga apenas as colunas necessárias para economizar memória. - **Re-runs seguros**: `caged_load()` detecta competências já gravadas e as pula automaticamente. - **Cache local**: os `.7z` ficam em `tools::R_user_dir("datacaged", "cache")` e são reutilizados entre sessões. - **Grandes volumes**: o pipeline processa uma competência de cada vez para não estourar a RAM.