Descrição do Projeto
Pipeline completo de ETL para análise de dados de uma concessionária de veículos. O projeto extrai dados operacionais de vendas, clientes, produtos e categorias de um banco de produção PostgreSQL, realiza transformações com dbt seguindo arquitetura de camadas, e disponibiliza os dados em um Data Warehouse para visualização em Looker Studio.
Toda a orquestração é feita via Apache Airflow, permitindo agendamento automático, monitoramento e retry de falhas.
Arquitetura do Pipeline
┌─────────────────┐
│ PostgreSQL │
│ (Produção) │
│ · vendas │
│ · clientes │
│ · produtos │
└────────┬────────┘
│ Extract
▼
┌─────────────────┐
│ Apache Airflow │
│ DAG Pipeline │
│ · Extração │
│ · Validação │
│ · Trigger dbt │
└────────┬────────┘
│ Load
▼
┌─────────────────┐
│ PostgreSQL │
│ (Data Warehouse)│
│ · staging │
└────────┬────────┘
│ Transform
▼
┌─────────────────┐
│ dbt │
│ · staging │
│ · dimensions │
│ · analytics │
└────────┬────────┘
│
▼
┌─────────────────┐
│ Looker Studio │
│ Dashboards │
│ · Vendas │
│ · Clientes │
│ · Produtos │
└─────────────────┘
Fluxo de Dados
- Extração: Airflow conecta ao banco de produção e extrai dados via SQL
- Carga: Dados brutos são carregados na camada staging do DW
- Transformação: dbt aplica transformações em 3 camadas (staging → dimensions → analytics)
- Visualização: Looker Studio consome as tabelas finais para dashboards
Stack Utilizada
Orquestração
- Apache Airflow 2.5+
- DAGs com agendamento diário
- Retry automático em falhas
Transformação
- dbt para transformações SQL
- Camadas: staging, dimensions, analytics
- Testes de qualidade de dados
Armazenamento
- PostgreSQL (origem e DW)
- Schemas separados por camada
Visualização
- Looker Studio
- Conexão direta com PostgreSQL
- Refresh automático
Exemplo de Código
DAG do Airflow
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.bash import BashOperator from datetime import datetime, timedelta default_args = { 'owner': 'data_team', 'retries': 3, 'retry_delay': timedelta(minutes=5), } with DAG( 'concessionaria_etl', default_args=default_args, schedule_interval='0 2 * * *', start_date=datetime(2024, 1, 1), catchup=False, ) as dag: extract = PostgresOperator( task_id='extract_vendas', sql='sql/extract_vendas.sql' ) run_dbt = BashOperator( task_id='run_dbt', bash_command='dbt run --models analytics' ) extract >> run_dbt
Modelo dbt — fact_vendas.sql
-- models/analytics/fact_vendas.sql {{ config(materialized='table', schema='analytics') }} WITH vendas_enriched AS ( SELECT v.venda_id, v.data_venda, v.valor_total, c.nome_cliente, c.cidade, p.modelo, p.categoria FROM {{ ref('stg_vendas') }} v LEFT JOIN {{ ref('dim_clientes') }} c ON v.cliente_id = c.cliente_id LEFT JOIN {{ ref('dim_produtos') }} p ON v.produto_id = p.produto_id ) SELECT * FROM vendas_enriched
Desafios e Soluções
1. Performance na Extração
Desafio: Extração completa levava mais de 2 horas, impactando o banco de produção.
Solução: Extração incremental baseada em timestamps, reduzindo tempo para ~15 minutos e carga no banco em 85%.
2. Qualidade de Dados
Desafio: Dados inconsistentes quebravam o pipeline (nulos, duplicatas).
Solução: Testes dbt para validar unicidade, não-nulidade e integridade referencial. Pipeline falha rápido com alertas claros.
3. Monitoramento e Alertas
Desafio: Falhas silenciosas não eram detectadas rapidamente.
Solução: Callbacks do Airflow para enviar alertas por email com logs detalhados.
Resultados e Métricas
Tempo de Execução Diária
Registros Processados/Dia
Taxa de Sucesso
Impacto no Negócio
- Dashboards atualizados diariamente com dados atuais
- Redução de 85% no tempo de geração de relatórios gerenciais
- Visibilidade em tempo real de vendas, estoque e performance
- Base sólida para expansão futura (ML, previsões)
Aprendizados
- Importância de testes de qualidade desde o início do projeto
- Benefícios da extração incremental para reduzir carga
- Valor do monitoramento proativo em pipelines de produção
- Estruturação de transformações em camadas lógicas (staging → dimensions → analytics)