Pipeline completo para ingestão, transformação e armazenamento de dados públicos brasileiros provenientes da API do IBGE.
Automatizar a coleta de dados do IBGE, aplicar transformações com PySpark e armazenar em PostgreSQL para análises e consultas.
- Linguagem: Python 3.10+ (container: Python 3.12)
- Orquestração: Apache Airflow 2.11.2 (última série 2.x, EOL abr/2026)
- Processamento: PySpark 3.5+
- Data Quality: Pandera (portão pós-coleta + portão pós-transformação)
- Banco de Dados: PostgreSQL
- APIs: IBGE (Localidades v1 + Agregados SIDRA v3)
05-pipeline-dados-publicos/
├── README.md
├── requirements.txt
├── docker-compose.yml
├── .env.example
├── .gitignore
├── 00-create-dados-publicos.sh
├── src/
│ ├── __init__.py
│ ├── collectors/
│ │ └── ibge_collector.py
│ ├── transformers/
│ │ └── spark_transformer.py
│ └── loaders/
│ └── postgres_loader.py
│ └── validation/
│ └── schemas.py
├── sql/
│ └── create_schema.sql
├── airflow/
│ ├── Dockerfile
│ └── dags/
│ ├── __init__.py
│ ├── common.py
│ └── pipeline_ibge.py
├── notebooks/
│ └── 01_exploracao_dados_publicos.ipynb
├── docs/
│ └── revisao-engenharia-dados.md
└── tests/
├── __init__.py
├── test_ibge_collector.py
├── test_postgres_loader.py
└── test_validation.py
pip install -r requirements.txt --constraint https://raw.githubusercontent.com/apache/airflow/constraints-2.11.2/constraints-3.12.txtSem o
--constraint, o pip resolve dependências incompatíveis do Airflow 2 e a instalação quebra. Ajuste o3.12para a sua versão local (python3 --version).
cp .env.example .env
# Editar .env com suas credenciais| Variável | Descrição | Padrão |
|---|---|---|
PG_HOST / PG_PORT / PG_DATABASE / PG_USER / PG_PASSWORD |
Conexão PostgreSQL | localhost:5432/dados_publicos |
IBGE_API_BASE |
API de localidades (v1) | https://servicodados.ibge.gov.br/api/v1 |
IBGE_AGREGADOS_BASE |
API de agregados SIDRA (v3) | https://servicodados.ibge.gov.br/api/v3/agregados |
IBGE_TIMEOUT |
Timeout das requisições (s) | 30 |
SPARK_MASTER / SPARK_APP_NAME |
Spark local | local[*] |
AIRFLOW_ADMIN_USER / AIRFLOW_ADMIN_PASSWORD |
Login da UI (DEV ONLY) | admin/admin |
AIRFLOW_UID |
UID dos containers (escrever em ./data) |
1000 (seu id -u) |
AIRFLOW_UIDprecisa ser o seu UID do host (id -u): sem isso o Spark no container não consegue sobrescreverdata/processed/e a task de transformação falha comUnable to clear output directory.
createdb -U postgres dados_publicos
psql -U postgres -d dados_publicos -f sql/create_schema.sqlpython -m src.collectors.ibge_collectorExemplo de agregado SIDRA (população estimada do Brasil, 2024):
from src.collectors.ibge_collector import IBGECollector
collector = IBGECollector()
dados = collector.buscar_agregado(6579, 2024, 9324, "N1[all]")python -m src.transformers.spark_transformerpytest tests/ -q14 testes (coletores com HTTP mockado, carga com engine mockada, schemas de DQ — sem rede nem banco).
- Portão 1 (pós-coleta): lote vazio ou sem campos estruturais (id/sigla,
hierarquia até a UF) levanta
DataQualityErrorantes de persistir o raw. - Portão 2 (pós-transformação): schemas Pandera barram
idduplicado ou nulo,sigla_uffora das 27 UFs e tipos errados — antes do DELETE, de modo que um lote reprovado nunca apaga o dado bom do dia.
jupyter notebook notebooks/01_exploracao_dados_publicos.ipynbO notebook lê data/processed/ primeiro (parquet do transformer) e só usa o
raw como fallback. Sem dados coletados, a célula de carga falha alto em vez
de pular as análises em silêncio.
Opção A — Podman (recomendado):
cp .env.example .env # se ainda não existe
podman compose up -d --buildSobe Postgres (porta 5433 no host) + scheduler + webserver.
UI em http://localhost:8080 — login admin/admin (padrão de DEV;
troque com AIRFLOW_ADMIN_USER/AIRFLOW_ADMIN_PASSWORD antes do up).
As credenciais PG_* dentro do compose apontam para o banco do container;
fora dele vale o .env local.
Aplicar o schema no banco do container:
podman compose exec postgres psql -U postgres -d dados_publicos -f /schema/create_schema.sqlOpção B — local:
export AIRFLOW__CORE__DAGS_FOLDER=$PWD/airflow/dags
airflow db upgrade
airflow users create --username admin --password admin --firstname Admin --lastname User --role Admin --email admin@example.com
airflow scheduler &
airflow webserver &Uma DAG: pipeline_dados_ibge (diária).
Ela executa coleta → transformação PySpark → carga nas tabelas landing
(ibge.*_raw) via src/loaders/postgres_loader.py.
A carga é idempotente por dia (DELETE + INSERT na mesma transação) e qualquer
falha reprova a task para acionar o retry — ver sql/create_schema.sql
(UNIQUE por dia nas *_raw).
- Regiões, estados, mesorregiões, microrregiões e municípios
- Distritos e subdistritos, regiões imediatas/intermediárias e metropolitanas
- Séries por tabela/período/variável e recorte territorial (ex.: tabela 6579, população estimada)
- Nunca commite o arquivo
.env(já coberto pelo.gitignore) - Credenciais devem ser gerenciadas via variáveis de ambiente
- Dados públicos não possuem restrição de acesso