πŸ› οΈ CapΓ­tulo 7 Β· Nivel BΓ‘sico-Intermedio

Herramientas Fundamentales

El stack tecnolΓ³gico del Data Engineer moderno. Estas herramientas son las que aparecen en casi todas las ofertas laborales y que debes dominar para ser efectivo en cualquier empresa.

⏱️ Lectura: ~65 min
🎯 Nivel: BΓ‘sico β†’ Intermedio
πŸ“Œ Foco: Uso prΓ‘ctico

El Stack Completo del Data Engineer 2026

πŸ”€

Git

git 2.43+
Control de versiones. Sin Git no hay trabajo en equipo ni CI/CD. Esencial.
version controlcollaboration
🐳

Docker

Docker 25+
ContainerizaciΓ³n. Garantiza que tu pipeline funciona igual en dev, staging y producciΓ³n.
containersreproducibility
🐘

PostgreSQL

PostgreSQL 16+
La RDBMS mΓ‘s poderosa del mundo open source. Para datos operacionales y analΓ­ticos a escala mediana.
RDBMSACIDSQL
πŸƒ

MongoDB

MongoDB 7+
Document store NoSQL lΓ­der. Para datos semi-estructurados y esquemas variables.
NoSQLdocumentsJSON
β­•

Apache Spark

Spark 3.5+
Motor de procesamiento distribuido. El estΓ‘ndar para Big Data processing a escala.
distributedbatchstreaming
πŸ“‘

Apache Kafka

Kafka 3.7+
Plataforma de event streaming distribuida. El corazΓ³n de arquitecturas event-driven.
streamingeventsqueue
πŸŒ€

Apache Airflow

Airflow 2.9+
Orquestador de workflows de datos. Programa, monitorea y gestiona pipelines con DAGs en Python.
orchestrationDAGsscheduling
βš™οΈ

dbt

dbt Core 1.8+
Data Build Tool. Transforma datos en el DWH con SQL modular, versionado y testeable.
transformationsSQLtesting
❄️

Snowflake

Snowflake 2026
Cloud Data Warehouse lΓ­der. CΓ³mputo y storage separados, auto-escale, SQL analΓ­tico supremo.
DWHcloudanalytics
πŸ”΅

BigQuery

BigQuery 2026
DWH serverless de Google. Paga por query, escala automΓ‘tico, excelente para analytics ad-hoc.
serverlessGCPpetabytes

Git para Data Engineering

Git es tan fundamental como el SQL. Un Data Engineer que no maneja Git avanzado no puede trabajar en equipo ni en CI/CD moderno.

πŸ“š Comandos Git Esenciales para DE β–Ό
# ── SETUP INICIAL ─────────────────────────────────────────────────────────
git config --global user.name "Tu Nombre"
git config --global user.email "tu@email.com"
git config --global core.editor "code --wait"      # VS Code como editor
git config --global init.defaultBranch main

# ── BRANCHING STRATEGY (GitFlow para Data Engineering) ────────────────────
# Crear feature branch para nuevo pipeline
git checkout -b feature/orders-etl-pipeline

# Trabajo en el pipeline...
git add .
git commit -m "feat(orders): add incremental ETL pipeline with dedup"

# Mantener sincronizado con main
git fetch origin main
git rebase origin/main  # Preferir rebase sobre merge para historia limpia

# ── GITIGNORE para proyectos de DE ────────────────────────────────────────
cat >> .gitignore << 'EOF'
# Datos (nunca versionar datos reales)
*.csv
*.parquet
*.json.gz
data/
raw/
processed/

# Credenciales (NUNCA versionar secrets)
.env
*.key
*.pem
secrets/
credentials.json
service_account.json

# Python
__pycache__/
*.pyc
.venv/
dist/
*.egg-info/

# Airflow
logs/
airflow.db
airflow.cfg  # Contiene passwords

# dbt
target/
dbt_packages/
profiles.yml  # Tiene conexiones de BD
EOF

# ── GIT PARA DATOS: Large File Storage ────────────────────────────────────
# Para archivos grandes (modelos ML, fixtures de test)
git lfs install
git lfs track "*.parquet"
git lfs track "*.pkl"

# ── COMMITS SEMÁNTICOS (Conventional Commits) ─────────────────────────────
# Formato: type(scope): description
git commit -m "feat(pipeline): add Kafka consumer for orders topic"
git commit -m "fix(etl): handle null customer_id in transform stage"
git commit -m "refactor(spark): optimize join with broadcast hint"
git commit -m "test(quality): add data quality checks for fact_orders"
git commit -m "docs(readme): update setup instructions for WSL2"

# ── PRE-COMMIT HOOKS para calidad de cΓ³digo ────────────────────────────────
# Instalar pre-commit
pip install pre-commit
cat > .pre-commit-config.yaml << 'EOF'
repos:
  - repo: https://github.com/psf/black
    rev: 24.0.0
    hooks:
      - id: black  # Auto-format Python
  - repo: https://github.com/pycqa/flake8
    rev: 7.0.0
    hooks:
      - id: flake8  # Linting
  - repo: https://github.com/sqlfluff/sqlfluff
    rev: 3.0.0
    hooks:
      - id: sqlfluff-lint  # SQL linting
        args: [--dialect, postgres]
EOF
pre-commit install
πŸ”€ Branching Strategy para equipos de datos β–Ό
# TRUNK-BASED DEVELOPMENT (recomendado para DE)
# Rama principal: main (siempre deployable)
# Feature branches: corta vida (1-3 dΓ­as)
# Sin ramas de larga vida (no develop, no release)

main
β”œβ”€β”€ feature/add-customer-dim          (1-2 dΓ­as)
β”œβ”€β”€ fix/null-handling-orders-etl      (horas)
└── refactor/optimize-spark-job       (1 dΓ­a)

# GITFLOW (para releases formales, menos comΓΊn en DE moderno)
main          # ProducciΓ³n
β”œβ”€β”€ develop   # IntegraciΓ³n
β”œβ”€β”€ feature/* # Nuevas funcionalidades
β”œβ”€β”€ release/* # PreparaciΓ³n de release
└── hotfix/*  # Fixes de producciΓ³n urgentes

# Pull Request Template para proyectos de DE
cat > .github/pull_request_template.md << 'EOF'
## DescripciΓ³n


## Tipo de cambio
- [ ] Nuevo pipeline
- [ ] ModificaciΓ³n de pipeline existente  
- [ ] Fix de bug en producciΓ³n
- [ ] Refactoring (sin cambio de funcionalidad)
- [ ] Cambio de schema/modelo

## Tests realizados
- [ ] Unit tests pasan (pytest)
- [ ] Data quality checks pasan
- [ ] Probado en staging con datos reales
- [ ] Performance benchmarked (si aplica)

## Impacto en datos
- Schema cambia: SΓ­ / No
- Datos histΓ³ricos requeridos re-procesar: SΓ­ / No
- SLA afectado: SΓ­ / No
EOF

Docker para Data Engineering

Docker garantiza que tus pipelines corren igual en tu laptop, en CI/CD y en producciΓ³n. Es la herramienta mΓ‘s importante para la reproducibilidad.

# ── DOCKERFILE para un pipeline Python ────────────────────────────────────
# Dockerfile
FROM python:3.11-slim AS base

# Sistema
RUN apt-get update && apt-get install -y \
    gcc g++ curl \
    && rm -rf /var/lib/apt/lists/*

WORKDIR /app

# Dependencies primero (mejor cache de Docker)
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

# CΓ³digo de aplicaciΓ³n
COPY src/ ./src/
COPY config/ ./config/

# No correr como root en producciΓ³n
RUN useradd -m -u 1000 dataengineer
USER dataengineer

ENV PYTHONPATH=/app
ENV PYTHONUNBUFFERED=1

ENTRYPOINT ["python", "-m", "src.pipeline"]
CMD ["--config", "config/production.yaml"]
# ── DOCKER-COMPOSE: Stack completo local de desarrollo ────────────────────
# docker-compose.yml
version: '3.8'

services:
  # Base de datos operacional
  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_DB: dataengineering
      POSTGRES_USER: de_user
      POSTGRES_PASSWORD: ${POSTGRES_PASSWORD}
    ports: ["5432:5432"]
    volumes:
      - postgres_data:/var/lib/postgresql/data
      - ./init.sql:/docker-entrypoint-initdb.d/init.sql
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U de_user"]
      interval: 10s
      timeout: 5s
      retries: 5

  # Kafka + Zookeeper
  zookeeper:
    image: confluentinc/cp-zookeeper:7.6.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181

  kafka:
    image: confluentinc/cp-kafka:7.6.0
    depends_on: [zookeeper]
    ports: ["9092:9092"]
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092

  # Apache Airflow
  airflow-webserver:
    image: apache/airflow:2.9.0
    command: webserver
    ports: ["8080:8080"]
    environment:
      - AIRFLOW__CORE__EXECUTOR=LocalExecutor
      - AIRFLOW__DATABASE__SQL_ALCHEMY_CONN=postgresql+psycopg2://de_user:${POSTGRES_PASSWORD}@postgres/airflow
    volumes:
      - ./dags:/opt/airflow/dags
      - ./plugins:/opt/airflow/plugins
    depends_on:
      postgres:
        condition: service_healthy

  airflow-scheduler:
    image: apache/airflow:2.9.0
    command: scheduler
    environment:
      - AIRFLOW__CORE__EXECUTOR=LocalExecutor
      - AIRFLOW__DATABASE__SQL_ALCHEMY_CONN=postgresql+psycopg2://de_user:${POSTGRES_PASSWORD}@postgres/airflow
    volumes:
      - ./dags:/opt/airflow/dags
    depends_on: [postgres]

  # Redis para cachΓ©
  redis:
    image: redis:7-alpine
    ports: ["6379:6379"]

volumes:
  postgres_data:

# Comandos ΓΊtiles:
# docker-compose up -d            β†’ Levanta todos los servicios
# docker-compose logs -f airflow  β†’ Ver logs de Airflow
# docker-compose exec postgres psql -U de_user -d dataengineering  β†’ Shell de PG
# docker-compose down -v          β†’ Baja y elimina volΓΊmenes

PostgreSQL Avanzado

-- ── CONFIGURACIONES CLAVE PARA DE ─────────────────────────────────────────
-- postgresql.conf (ajustar para analytics workloads)
-- shared_buffers = 256MB (25% de RAM total)
-- effective_cache_size = 768MB (75% de RAM total)
-- work_mem = 64MB (para sorts y hash joins)
-- maintenance_work_mem = 512MB (para VACUUM, CREATE INDEX)
-- max_parallel_workers_per_gather = 4 (paralelismo en queries)

-- ── JSONB: PostgreSQL como document store ──────────────────────────────────
CREATE TABLE events (
    event_id    BIGSERIAL PRIMARY KEY,
    event_type  TEXT NOT NULL,
    occurred_at TIMESTAMP NOT NULL DEFAULT NOW(),
    data        JSONB NOT NULL,  -- Datos arbitrarios
    tags        TEXT[]           -- Array de tags
);

-- Índice GIN para búsqueda eficiente en JSONB
CREATE INDEX idx_events_data ON events USING GIN (data);

-- Insertar evento con datos arbitrarios
INSERT INTO events (event_type, data, tags) VALUES
('order_placed', 
 '{"order_id": 1001, "customer_id": 501, "items": [{"sku": "LAP-001", "qty": 1}], "total": 999.99}',
 ARRAY['ecommerce', 'web']
),
('user_login', 
 '{"user_id": 501, "ip": "192.168.1.1", "browser": "Chrome", "device": "desktop"}',
 ARRAY['auth']
);

-- Queries sobre JSONB (usa el Γ­ndice GIN)
SELECT * FROM events WHERE data @> '{"order_id": 1001}';  -- Contiene
SELECT data->>'customer_id' AS customer_id FROM events WHERE event_type = 'order_placed';
SELECT jsonb_path_query(data, '$.items[*].sku') FROM events WHERE event_type = 'order_placed';

-- ── FULL-TEXT SEARCH ─────────────────────────────────────────────────────
-- Agregar columna de bΓΊsqueda
ALTER TABLE products ADD COLUMN search_vector tsvector
    GENERATED ALWAYS AS (
        setweight(to_tsvector('spanish', coalesce(name, '')), 'A') ||
        setweight(to_tsvector('spanish', coalesce(description, '')), 'B')
    ) STORED;

CREATE INDEX idx_products_search ON products USING GIN (search_vector);

-- BΓΊsqueda eficiente
SELECT name, ts_rank(search_vector, query) AS rank
FROM products, to_tsquery('spanish', 'laptop & gaming & portable') query
WHERE search_vector @@ query
ORDER BY rank DESC LIMIT 10;

-- ── MATERIALIZED VIEWS para analytics ─────────────────────────────────────
CREATE MATERIALIZED VIEW mv_sales_summary AS
SELECT
    DATE_TRUNC('day', order_date) AS sale_date,
    product_category,
    customer_segment,
    COUNT(*)                 AS orders,
    SUM(amount)              AS revenue,
    AVG(amount)              AS avg_ticket
FROM fact_orders
GROUP BY 1, 2, 3;

CREATE INDEX ON mv_sales_summary (sale_date, product_category);

-- Refrescar periΓ³dicamente (cron job o Airflow)
REFRESH MATERIALIZED VIEW CONCURRENTLY mv_sales_summary;
-- CONCURRENTLY: actualiza sin lock (tabla sigue siendo consultable)

dbt: Data Build Tool

dbt transforma SQL de scripts ad-hoc en cΓ³digo de ingenierΓ­a con versionado, testing, documentaciΓ³n y dependencias automΓ‘ticas. Es posiblemente la herramienta que mΓ‘s impacto ha tenido en el Data Engineering moderno (2018-2026).

Estructura de proyecto dbt

# ── ESTRUCTURA DE PROYECTO DBT ─────────────────────────────────────────────
my_dbt_project/
β”œβ”€β”€ dbt_project.yml          # ConfiguraciΓ³n del proyecto
β”œβ”€β”€ profiles.yml             # Conexiones a bases de datos (NO en git)
β”œβ”€β”€ models/
β”‚   β”œβ”€β”€ staging/             # Bronze β†’ Silver: una fuente, 1-1
β”‚   β”‚   β”œβ”€β”€ stg_orders.sql
β”‚   β”‚   β”œβ”€β”€ stg_customers.sql
β”‚   β”‚   └── schema.yml       # Tests y documentaciΓ³n
β”‚   β”œβ”€β”€ intermediate/        # LΓ³gica de negocio compleja
β”‚   β”‚   └── int_customer_orders.sql
β”‚   └── marts/               # Silver β†’ Gold: tablas finales de negocio
β”‚       β”œβ”€β”€ finance/
β”‚       β”‚   └── fct_revenue.sql
β”‚       └── marketing/
β”‚           └── dim_customers.sql
β”œβ”€β”€ tests/                   # Tests de datos custom
β”‚   └── assert_orders_positive.sql
β”œβ”€β”€ macros/                  # ReutilizaciΓ³n de SQL
β”‚   └── generate_surrogate_key.sql
└── snapshots/               # SCD Tipo 2 automΓ‘tico
    └── snapshot_customers.sql
-- ── MODELO STAGING: Limpieza bΓ‘sica ────────────────────────────────────────
-- models/staging/stg_orders.sql
{{ config(materialized='view') }}  -- Las vistas son mΓ‘s eficientes en staging

WITH source AS (
    -- Referencia a la fuente de datos
    SELECT * FROM {{ source('raw', 'orders') }}
),

renamed AS (
    SELECT
        -- Renombrar y tipar correctamente
        order_id::INTEGER      AS order_id,
        customer_id::INTEGER   AS customer_id,
        order_date::TIMESTAMP  AS ordered_at,
        total_amount::DECIMAL  AS order_amount_usd,
        LOWER(TRIM(status))    AS status,
        CURRENT_TIMESTAMP      AS _loaded_at
    FROM source
),

cleaned AS (
    SELECT *
    FROM renamed
    WHERE order_id IS NOT NULL
      AND customer_id IS NOT NULL
      AND order_amount_usd >= 0
)

SELECT * FROM cleaned

-- ── MODELO MART: Tabla de hechos final ──────────────────────────────────────
-- models/marts/finance/fct_orders.sql
{{ config(
    materialized='incremental',
    unique_key='order_id',
    incremental_strategy='merge'
) }}

WITH orders AS (
    SELECT * FROM {{ ref('stg_orders') }}
    {% if is_incremental() %}
        -- Solo procesar registros nuevos en modo incremental
        WHERE ordered_at >= (SELECT MAX(ordered_at) FROM {{ this }})
    {% endif %}
),

customers AS (
    SELECT * FROM {{ ref('dim_customers') }}
),

final AS (
    SELECT
        {{ dbt_utils.generate_surrogate_key(['order_id']) }} AS order_pk,
        o.order_id,
        o.ordered_at,
        o.order_amount_usd,
        o.status,
        c.customer_segment,
        c.country,
        -- MΓ©tricas derivadas
        CASE WHEN o.order_amount_usd >= 500 THEN 'large'
             WHEN o.order_amount_usd >= 100 THEN 'medium'
             ELSE 'small' END AS order_size
    FROM orders o
    LEFT JOIN customers c ON o.customer_id = c.customer_id
)

SELECT * FROM final

-- ── SCHEMA.YML: Tests y DocumentaciΓ³n ──────────────────────────────────────
# models/marts/finance/schema.yml
version: 2
models:
  - name: fct_orders
    description: "Tabla de hechos con todas las Γ³rdenes entregadas"
    columns:
      - name: order_pk
        description: "Surrogate key del pedido"
        tests:
          - unique
          - not_null
      - name: order_id
        tests:
          - unique
          - not_null
      - name: order_amount_usd
        tests:
          - not_null
          - dbt_expectations.expect_column_values_to_be_between:
              min_value: 0
              max_value: 100000
      - name: status
        tests:
          - accepted_values:
              values: ['pending', 'processing', 'shipped', 'delivered', 'cancelled']

# ── COMANDOS DBT ESENCIALES ────────────────────────────────────────────────
# dbt debug               β†’ Verificar conexiΓ³n
# dbt run                 β†’ Ejecutar todos los modelos
# dbt run -s stg_orders   β†’ Solo un modelo
# dbt run -s +fct_orders  β†’ fct_orders y sus dependencias upstream
# dbt test                β†’ Ejecutar todos los tests
# dbt test -s fct_orders  β†’ Tests de un modelo especΓ­fico
# dbt docs generate && dbt docs serve  β†’ Generar y ver documentaciΓ³n
# dbt snapshot            β†’ Ejecutar snapshots (SCD2)

Snowflake

-- ── ARQUITECTURA SNOWFLAKE: CΓ³mputo y Storage separados ─────────────────
-- Virtual Warehouses: unidades de cΓ³mputo independientes y escalables
CREATE WAREHOUSE analytics_wh
    WITH WAREHOUSE_SIZE = 'MEDIUM'     -- XS, S, M, L, XL, 2XL, ...
    AUTO_SUSPEND = 300                  -- Se apaga si no hay queries (ahorra costos)
    AUTO_RESUME = TRUE                  -- Se activa automΓ‘ticamente
    SCALING_POLICY = 'ECONOMY';        -- Solo escala si hay mucha carga

-- MΓΊltiples warehouses = sin contenciΓ³n entre equipos
CREATE WAREHOUSE bi_wh WITH WAREHOUSE_SIZE = 'SMALL' AUTO_SUSPEND = 60;
CREATE WAREHOUSE etl_wh WITH WAREHOUSE_SIZE = 'LARGE' AUTO_SUSPEND = 180;

-- ── TIME TRAVEL: Consultar datos histΓ³ricos ────────────────────────────────
-- Ver tabla como era hace 2 dΓ­as
SELECT * FROM fact_orders AT (OFFSET => -172800);  -- -172800 segundos = 2 dΓ­as

-- Antes de un error especΓ­fico
SELECT * FROM fact_orders BEFORE (STATEMENT => '0196b9c8-0033-8c6b-0002-697c000ae17e');

-- Clonar tabla a estado anterior (instantΓ‘neo, zero-copy)
CREATE TABLE fact_orders_backup CLONE fact_orders AT (OFFSET => -3600);

-- ── ZERO-COPY CLONING: Clonar sin duplicar storage ────────────────────────
-- Clone de base de datos completa para ambiente de testing
CREATE DATABASE prod_clone CLONE production;
-- Inmediato y sin duplicar datos. Solo el delta ocupa storage adicional.

-- ── STREAMS & TASKS: CDC automatizado ─────────────────────────────────────
-- Stream: captura cambios (INSERT/UPDATE/DELETE) en una tabla
CREATE STREAM orders_stream ON TABLE raw_orders
    SHOW_INITIAL_ROWS = TRUE;

-- Task: procesa el stream periΓ³dicamente
CREATE TASK process_orders_task
    WAREHOUSE = etl_wh
    SCHEDULE = '5 MINUTE'
    WHEN SYSTEM$STREAM_HAS_DATA('orders_stream')
AS
    MERGE INTO silver_orders AS target
    USING (
        SELECT order_id, customer_id, amount, status,
               METADATA$ACTION, METADATA$ISUPDATE
        FROM orders_stream
    ) AS source
    ON target.order_id = source.order_id
    WHEN MATCHED AND source.METADATA$ACTION = 'DELETE' THEN DELETE
    WHEN MATCHED AND source.METADATA$ACTION = 'INSERT' THEN UPDATE SET ...
    WHEN NOT MATCHED AND source.METADATA$ACTION = 'INSERT' THEN INSERT ...;

-- Activar la task
ALTER TASK process_orders_task RESUME;

BigQuery

-- ── BIGQUERY: SQL analΓ­tico serverless a escala de petabytes ──────────────
-- Particionamiento automΓ‘tico por fecha (reduce costos drasticamente)
CREATE TABLE `project.dataset.orders`
PARTITION BY DATE(created_at)
CLUSTER BY customer_id, status  -- Clustering: como Γ­ndice de columnas
OPTIONS (
    partition_expiration_days = 365,  -- Auto-expire particiones viejas
    require_partition_filter = TRUE    -- Fuerza usar filtro de particiΓ³n
) AS
SELECT * FROM `project.dataset.raw_orders`;

-- Query eficiente (usa partition pruning + clustering)
SELECT customer_id, SUM(amount) AS total
FROM `project.dataset.orders`
WHERE DATE(created_at) BETWEEN '2026-01-01' AND '2026-03-31'  -- Solo esas particiones
  AND status = 'delivered'
GROUP BY customer_id;

-- ── BQML: Machine Learning directo en SQL ─────────────────────────────────
CREATE OR REPLACE MODEL `project.models.customer_churn`
OPTIONS (
    model_type = 'logistic_reg',
    input_label_cols = ['churned'],
    auto_class_weights = TRUE
) AS
SELECT
    recency_days, frequency, monetary_value,
    avg_session_duration, product_categories_count,
    churned
FROM `project.dataset.customer_features`
WHERE train_date < '2025-10-01';

-- Predicciones
SELECT customer_id, predicted_churned, predicted_churned_probs
FROM ML.PREDICT(
    MODEL `project.models.customer_churn`,
    (SELECT * FROM `project.dataset.active_customers`)
);

-- ── BIGQUERY STORAGE API: Lectura masiva desde Python ─────────────────────
from google.cloud import bigquery, bigquery_storage

bq_client = bigquery.Client()
bqs_client = bigquery_storage.BigQueryReadClient()

# Leer tabla grande eficientemente con Arrow
table = f"projects/{PROJECT}/datasets/{DATASET}/tables/{TABLE}"
requested_session = bqs_client.create_read_session(
    parent=f"projects/{PROJECT}",
    read_session=bigquery_storage.ReadSession(
        table=table,
        data_format=bigquery_storage.DataFormat.ARROW,
    ),
    max_stream_count=10  # Paralelismo
)

# Leer en paralelo
import pyarrow as pa
frames = []
for stream in requested_session.streams:
    reader = bqs_client.read_rows(stream.name)
    frames.append(reader.to_dataframe())
df = pd.concat(frames)