Este projeto é um desafio técnico desenvolvido para demonstrar a capacidade de planejar, implementar e manter uma pipeline de dados que extrai informações de múltiplas fontes, armazena os dados localmente e os carrega em um banco de dados de destino.
A solução consiste em duas etapas principais:
-
Extração e Armazenamento Local:
- Extrair dados de:
- Banco de dados PostgreSQL (Northwind).
- Arquivo CSV representando a tabela
order_details.
- Armazenar os dados extraídos como arquivos Parquet organizados por tabelas e por data de execução.
- Extrair dados de:
-
Carregamento no Banco de Destino:
- Carregar os arquivos Parquet gerados no banco de dados PostgreSQL de destino, no schema
processed. - Garantir que as tabelas
orderseorder_detailsestejam disponíveis para consultas.
- Carregar os arquivos Parquet gerados no banco de dados PostgreSQL de destino, no schema
- Airflow: Gerenciamento e agendamento de tarefas.
- Meltano: Pipeline ETL para extração e carregamento de dados.
- PostgreSQL: Banco de dados de origem e destino.
- Docker Compose: Orquestração dos serviços.
├── data
│ │── db_destination_persisted (persistir os dados quando populados no passo 2)
│ ├── processed
│ │ ├── postgres
│ │ │ └── orders (criada dinamicamente na primeira execução)
│ │ │ └── 2024-01-01 (criada dinamicamente)
│ │ │ └── orders.parquet
│ │ └── csv
│ │ └── 2024-01-01 (criada dinamicamente)
│ │ └── order_details.parquet
│ └── source
│ └── northwind.sql
│ └── order_details.csv
├── meltano_project
│ ├── meltano.yml
├── airflow
│ ├── dags
│ │ ├── postgres_csv_to_local_pipeline.py
│ │ ├── local_to_postgres_pipeline.py
│ └── logs
│ └── plugins
├── docker-compose.yml
└── requirements.txt
└── README.md
- Docker e Docker Compose instalados.
- Python 3.8+ instalado (se necessário rodar localmente).
Certifique-se de instalar as bibliotecas necessárias:
pip install -r requirements.txtSuba os serviços:
docker-compose up -dIsso iniciará:
- PostgreSQL.
- Airflow.
- Meltano.
A DAG extract_to_parquet executa a extração dos dados de:
- Tabelas do banco
northwind. - Arquivo CSV de
order_details.
Os dados extraídos serão armazenados em:
/data/processed/{source}/{table}/{execution_date}/{table}.parquet
Para rodar manualmente:
airflow dags trigger extract_to_parquetA DAG load_parquet_to_postgres carrega os arquivos Parquet no banco de destino (destination), no schema processed.
Para rodar manualmente:
airflow dags trigger load_parquet_to_postgresApós executar o pipeline, você pode consultar os dados no banco de destino:
SELECT o.order_id, o.customer_id, d.product_id, d.unit_price, d.quantity, d.discount
FROM processed.orders o
JOIN processed.order_details d ON o.order_id = d.order_id;- Data Retroativa:
- No Airflow, configure a data de execução para reprocessar os dados de uma data que desejar.
- Uso das ferramentas especificadas: Airflow, Meltano, PostgreSQL.
- Idempotência: Arquivos e carregamentos são organizados por data, evitando duplicação.
- Dependências explícitas: Step 2 depende do sucesso do Step 1.
- Extração completa: Todas as tabelas do PostgreSQL e o arquivo CSV são extraídos.
- Monitoramento e Logs: Airflow fornece logs e visualização detalhada do pipeline.
- Reprocessamento por data: Suporta execução diária e retroativa.
- Instruções claras: Este README explica como executar e monitorar o pipeline.