Instalar com Codex ou Claude Copie este prompt, cole no Codex, Claude ou outro assistente e deixe que ele revise a página da skill e instale para você.
Um comando direto ignora o prompt de revisão. Verifique a origem antes de executá-lo.
Instruções da origem · Visualização somente leitura
name
data-engineering-medallion-pipeline
description
End-to-end data engineering pipeline using MinIO, Airbyte, PostgreSQL, DBT, and Airflow with medallion architecture (Bronze/Silver/Gold layers)
triggers
["set up a medallion architecture data pipeline","configure airbyte with minio and postgres","create dbt bronze silver gold models","orchestrate data pipeline with airflow","implement data quality tests in dbt","build an elt pipeline with docker compose","create a data lakehouse with medallion layers","deploy data engineering stack locally"]
This skill enables AI agents to work with a complete data engineering pipeline implementing the Medallion Architecture (Bronze → Silver → Gold) using modern open-source tools: MinIO (S3-compatible storage), Airbyte (data ingestion), PostgreSQL (data warehouse), DBT (transformations), Apache Airflow (orchestration), and Grafana (monitoring).
What This Project Does
The data-engineering-medallion project provides a complete end-to-end data pipeline that:
Ingests raw data from MinIO object storage into PostgreSQL using Airbyte
Transforms data through three layers (Bronze/Silver/Gold) using DBT
Orchestrates the entire pipeline with Apache Airflow DAGs
Validates data quality with automated DBT tests
Monitors infrastructure health with Prometheus and Grafana
Visualizes business metrics in Power BI dashboards
The architecture follows ELT (Extract-Load-Transform) pattern with clear separation of concerns:
Bronze: Raw immutable data from sources (JSONB format)
git clone https://github.com/LucasGoulartCouto/data-engineering-medallion.git
cd data-engineering-medallion
# Setup environment and start all services
make setup
make start
# Verify all containers are healthy
make status
# Infrastructure
make setup # Create .env, directories, install dependencies
make start # Start all Docker services
make stop # Stop all services
make restart # Restart all services
make status # Check container health
make logs SERVICE=airflow # View logs for specific service# DBT Operations
make dbt-run # Run all DBT models (bronze → silver → gold)
make dbt-test # Run data quality tests
make dbt-docs # Generate and serve documentation
make dbt-snapshot # Capture SCD Type 2 snapshots
make dbt-clean # Clean compiled artifacts# Data Pipeline
make upload-data # Upload sample data to MinIO
make trigger-dag DAG_ID=bronze_ingestion_dag # Manually trigger Airflow DAG# Development
make lint # Lint Python and SQL code
make format # Format Python code with black
make validate # Validate Airflow DAGs and DBT models# Cleanup
make clean # Remove volumes and stop services
make clean-all # Full cleanup including Docker images
Silver models clean, cast types, deduplicate, and add calculated fields:
-- dbt/models/silver/silver_orders.sql
{{ config(
materialized='table',
schema='silver',
unique_key='order_id'
) }}
WITH deduplicated AS (
SELECT*,
ROW_NUMBER() OVER (
PARTITIONBY order_id
ORDERBY _airbyte_emitted_at DESC
) AS rn
FROM {{ ref('bronze_orders') }}
)
SELECT
order_id::INTEGER,
customer_id::INTEGER,
order_date::DATE,
total_amount::DECIMAL(10,2),
UPPER(TRIM(status)) AS status,
CASEWHEN total_amount::DECIMAL>1000THEN'high_value'WHEN total_amount::DECIMAL>500THEN'medium_value'ELSE'low_value'ENDAS order_value_segment,
_airbyte_emitted_at AS ingested_at,
CURRENT_TIMESTAMPAS transformed_at
FROM deduplicated
WHERE rn =1AND order_id ISNOT NULLAND order_date::DATE<=CURRENT_DATE
Gold Layer (Business Metrics)
Gold models aggregate data for business consumption:
-- dbt/models/gold/gold_product_performance.sql
{{ config(
materialized='table',
schema='gold'
) }}
SELECT
p.product_id,
p.product_name,
p.category,
p.supplier,
COUNT(DISTINCT oi.order_id) AS total_orders,
SUM(oi.quantity) AS total_quantity_sold,
SUM(oi.line_total) AS total_revenue,
SUM(oi.quantity * p.cost) AS total_cost,
SUM(oi.line_total) -SUM(oi.quantity * p.cost) AS total_profit,
ROUND(
(SUM(oi.line_total) -SUM(oi.quantity * p.cost)) /NULLIF(SUM(oi.line_total), 0) *100,
2
) AS profit_margin_percentage,
AVG(oi.unit_price) AS avg_unit_price,
MAX(o.order_date) AS last_sale_date
FROM {{ ref('silver_products') }} p
INNERJOIN {{ ref('silver_order_items') }} oi
ON p.product_id = oi.product_id
INNERJOIN {{ ref('silver_orders') }} o
ON oi.order_id = o.order_id
WHERE o.status ='completed'GROUPBY p.product_id, p.product_name, p.category, p.supplier
DBT Testing & Data Quality
Schema Tests (schema.yml)
# dbt/models/silver/schema.ymlversion:2models:-name:silver_ordersdescription:"Cleaned and validated orders"columns:-name:order_iddescription:"Primary key"tests:-unique-not_null-name:customer_iddescription:"Foreign key to customers"tests:-not_null-relationships:to:ref('silver_customers')field:customer_id-name:statustests:-accepted_values:values: ['PENDING', 'COMPLETED', 'CANCELLED', 'REFUNDED']
-name:total_amounttests:-dbt_utils.accepted_range:min_value:0max_value:1000000
Custom Tests
-- dbt/tests/assert_no_future_dates_in_sales.sqlSELECT order_id, order_date
FROM {{ ref('silver_orders') }}
WHERE order_date >CURRENT_DATE
-- dbt/tests/assert_no_negative_profit.sqlSELECT product_id, total_profit
FROM {{ ref('gold_product_performance') }}
WHERE total_profit <0
Running Tests
# Run all tests
make dbt-test
# Run tests for specific model
docker-compose exec dbt dbt test --select silver_orders
# Run specific test type
docker-compose exec dbt dbt test --select test_type:generic
docker-compose exec dbt dbt test --select test_type:singular
Airflow DAG Development
Bronze Ingestion DAG
# airflow/dags/bronze_ingestion_dag.pyfrom datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.airbyte.operators.airbyte import AirbyteTriggerSyncOperator
from airflow.operators.python import PythonOperator
default_args = {
'owner': 'data-engineering',
'depends_on_past': False,
'email_on_failure': True,
'email_on_retry': False,
'retries': 3,
'retry_delay': timedelta(minutes=5),
}
with DAG(
'bronze_ingestion_dag',
default_args=default_args,
description='Ingest data from MinIO to PostgreSQL Bronze layer',
schedule_interval='@daily',
start_date=datetime(2024, 1, 1),
catchup=False,
tags=['bronze', 'ingestion', 'airbyte'],
) as dag:
trigger_airbyte_sync = AirbyteTriggerSyncOperator(
task_id='trigger_airbyte_orders_sync',
airbyte_conn_id='airbyte_default',
connection_id='{{ var.value.airbyte_connection_id }}',
asynchronous=False,
timeout=3600,
)
defvalidate_ingestion(**context):
from airflow.providers.postgres.hooks.postgres import PostgresHook
pg_hook = PostgresHook(postgres_conn_id='postgres_default')
# Check row count
result = pg_hook.get_first(
"SELECT COUNT(*) FROM airbyte_raw._airbyte_raw_orders"
)
if result[0] == 0:
raise ValueError("No data ingested to bronze layer")
print(f"Validated {result[0]} rows in bronze layer")
validate_task = PythonOperator(
task_id='validate_bronze_ingestion',
python_callable=validate_ingestion,
)
trigger_airbyte_sync >> validate_task
Silver Transformation DAG
# airflow/dags/silver_transformation_dag.pyfrom datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.dbt.cloud.operators.dbt import DbtRunOperator
from airflow.operators.bash import BashOperator
default_args = {
'owner': 'data-engineering',
'retries': 2,
'retry_delay': timedelta(minutes=3),
}
with DAG(
'silver_transformation_dag',
default_args=default_args,
description='Transform bronze to silver layer with DBT',
schedule_interval='@daily',
start_date=datetime(2024, 1, 1),
catchup=False,
tags=['silver', 'transformation', 'dbt'],
) as dag:
run_silver_models = BashOperator(
task_id='run_dbt_silver_models',
bash_command='cd /opt/dbt && dbt run --select silver.*',
)
test_silver_models = BashOperator(
task_id='test_dbt_silver_models',
bash_command='cd /opt/dbt && dbt test --select silver.*',
)
run_silver_models >> test_silver_models
Gold Aggregation DAG
# airflow/dags/gold_aggregation_dag.pyfrom datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import BranchPythonOperator
default_args = {
'owner': 'data-engineering',
'retries': 2,
'retry_delay': timedelta(minutes=3),
}
with DAG(
'gold_aggregation_dag',
default_args=default_args,
description='Build gold layer business metrics',
schedule_interval='@daily',
start_date=datetime(2024, 1, 1),
catchup=False,
tags=['gold', 'aggregation', 'metrics'],
) as dag:
run_gold_models = BashOperator(
task_id='run_dbt_gold_models',
bash_command='cd /opt/dbt && dbt run --select gold.*',
)
test_gold_models = BashOperator(
task_id='test_dbt_gold_models',
bash_command='cd /opt/dbt && dbt test --select gold.*',
)
snapshot_gold = BashOperator(
task_id='snapshot_gold_metrics',
bash_command='cd /opt/dbt && dbt snapshot',
)
run_gold_models >> test_gold_models >> snapshot_gold
Data Upload to MinIO
# scripts/upload_to_minio.pyimport os
from minio import Minio
from minio.error import S3Error
import pandas as pd
from pathlib import Path
defupload_data_to_minio():
"""Upload sample CSV data to MinIO bucket"""# Initialize MinIO client
client = Minio(
os.getenv('MINIO_ENDPOINT', 'localhost:9000'),
access_key=os.getenv('MINIO_ACCESS_KEY', 'minioadmin'),
secret_key=os.getenv('MINIO_SECRET_KEY'),
secure=False
)
bucket_name = os.getenv('MINIO_BUCKET', 'raw-data')
# Create bucket if not existstry:
ifnot client.bucket_exists(bucket_name):
client.make_bucket(bucket_name)
print(f"Created bucket: {bucket_name}")
except S3Error as e:
print(f"Error creating bucket: {e}")
return# Upload files from data directory
data_dir = Path('data/raw')
for file_path in data_dir.glob('*.csv'):
try:
client.fput_object(
bucket_name,
file_path.name,
str(file_path),
content_type='text/csv'
)
print(f"Uploaded: {file_path.name}")
except S3Error as e:
print(f"Error uploading {file_path.name}: {e}")
if __name__ == '__main__':
upload_data_to_minio()
Run the upload:
make upload-data
# or
python scripts/upload_to_minio.py
-- dbt/macros/test_no_orphans.sql
{% test no_orphan_records(model, column_name, parent_model, parent_column) %}
SELECT {{ column_name }}
FROM {{ model }}
WHERE {{ column_name }} ISNOT NULLAND {{ column_name }} NOTIN (
SELECT {{ parent_column }}
FROM {{ parent_model }}
)
{% endtest %}
Querying the Data Warehouse
Bronze Layer Query
-- Raw JSONB data from AirbyteSELECT
_airbyte_ab_id,
_airbyte_emitted_at,
_airbyte_data->>'order_id'AS order_id,
_airbyte_data
FROM airbyte_raw._airbyte_raw_orders
LIMIT 10;
Silver Layer Query
-- Cleaned typed dataSELECT
order_id,
customer_id,
order_date,
total_amount,
status,
order_value_segment
FROM silver.silver_orders
WHERE order_date >=CURRENT_DATE-INTERVAL'30 days'ORDERBY order_date DESC;
Gold Layer Query
-- Business metricsSELECT
product_name,
category,
total_revenue,
total_profit,
profit_margin_percentage,
total_quantity_sold
FROM gold.gold_product_performance
WHERE profit_margin_percentage >20ORDERBY total_revenue DESC
LIMIT 10;
-- Customer segmentationSELECT
customer_segment,
COUNT(*) AS customer_count,
AVG(total_spent) AS avg_lifetime_value,
AVG(recency_days) AS avg_recency
FROM gold.gold_customer_metrics
GROUPBY customer_segment
ORDERBY avg_lifetime_value DESC;
Troubleshooting
Container Won't Start
# Check logs
make logs SERVICE=postgres
make logs SERVICE=airflow-webserver
# Verify resource allocation
docker system df
docker system prune # Clean up if needed# Reset specific service
docker-compose restart postgres
# Run with debug logging
docker-compose exec dbt dbt run --select silver_orders --debug
# Check compiled SQLcat dbt/target/compiled/datawarehouse/models/silver/silver_orders.sql
# Test single model
docker-compose exec dbt dbt run --select silver_orders --full-refresh
# Validate syntax without execution
docker-compose exec dbt dbt parse
# Wait for PostgreSQL to be ready
docker-compose exec postgres pg_isready -U ${POSTGRES_USER}# Check connection from within container
docker-compose exec dbt psql -h postgres -U ${POSTGRES_USER} -d ${POSTGRES_DB}# Verify connection string in .env
grep POSTGRES .env
Data Quality Test Failures
# Run tests with detailed output
docker-compose exec dbt dbt test --select silver_orders --store-failures
# Query failed test results
SELECT * FROM silver.test_failures
WHERE test_name = 'unique_order_id';
# Run specific test
docker-compose exec dbt dbt test --select test_name:unique_order_id
MinIO Upload Failing
# Check MinIO service status
docker-compose ps minio
# Test MinIO API directly
docker-compose exec minio mc aliassetlocal http://localhost:9000 \
${MINIO_ROOT_USER}${MINIO_ROOT_PASSWORD}
docker-compose exec minio mc lslocal/${MINIO_BUCKET}# Upload test file
docker-compose exec minio mc cp /tmp/test.csv local/${MINIO_BUCKET}/
Monitoring & Observability
Check Pipeline Health
# Airflow task status
docker-compose exec airflow-webserver \
airflow tasks state bronze_ingestion_dag trigger_airbyte_orders_sync 2024-01-01
# DBT run results
docker-compose exec dbt dbt run-operation log_run_results
# PostgreSQL query performance
docker-compose exec postgres psql -U ${POSTGRES_USER} -d ${POSTGRES_DB} -c \
"SELECT query, calls, total_time FROM pg_stat_statements ORDER BY total_time DESC LIMIT 10;"
groups:-name:dbt_pipelinerules:-alert:DBTModelFailedexpr:dbt_model_run_status{status="error"}>0for:5mlabels:severity:criticalannotations:summary:"DBT model {{ $labels.model }} failed"
This skill covers the complete data-engineering-medallion pipeline from setup through production operations, enabling AI agents to assist with implementation, troubleshooting, and extension of medallion architecture data pipelines.