✅ Capítulo 14 · Nivel Avanzado

Observabilidad de Datos

La calidad de datos es la disciplina que garantiza que los datos son precisos, completos, consistentes y oportunos. Sin calidad de datos, todos los análisis y modelos ML construidos sobre ellos son potencialmente inútiles o peligrosos.

⏱️ Lectura: ~55 min
🎯 Nivel: Avanzado
📌 Herramientas: GE, Soda, dbt tests

Las 6 Dimensiones de Calidad de Datos

✅ Completeness (Completitud)

¿Están todos los datos necesarios presentes? % de valores NULL o faltantes. Ej: 95% de clientes tienen email registrado.

🎯 Accuracy (Precisión)

¿Los datos reflejan la realidad? Son correctos y libres de errores. Ej: el precio de un producto coincide con el precio real.

🔗 Consistency (Consistencia)

¿Los datos son coherentes entre sistemas? El mismo dato tiene el mismo valor en todas las fuentes. Ej: el total del pedido coincide en CRM y ERP.

⏰ Timeliness (Oportunidad)

¿Los datos están disponibles cuando se necesitan? Freshness. Ej: los datos de ventas del día deben estar disponibles antes de las 8am.

🆔 Uniqueness (Unicidad)

¿No hay duplicados no deseados? Cada entidad aparece una sola vez. Ej: un cliente no tiene dos registros con el mismo email.

✔️ Validity (Validez)

¿Los datos cumplen el formato y rango esperado? Ej: la edad está entre 0-150, el email tiene formato válido, el código de país es de 2 letras.

Great Expectations

Great Expectations (GX) es el framework Python más completo para validar, documentar y perfilar datos. Genera documentación HTML interactiva con los resultados de las validaciones.

"""
Great Expectations 0.18+ (GX Core)
"""
import great_expectations as gx
import pandas as pd

# ── SETUP ─────────────────────────────────────────────────────────────────
context = gx.get_context()  # Data Context: gestiona configuración y artefactos

# ── DEFINIR EXPECTATIONS (VALIDACIONES) ────────────────────────────────────
# Desde un DataFrame
df = pd.read_csv('orders.csv')
validator = context.sources.pandas_default.read_dataframe(df)

# Completitud
validator.expect_column_values_to_not_be_null("order_id")
validator.expect_column_values_to_not_be_null("customer_id")

# Unicidad
validator.expect_column_values_to_be_unique("order_id")

# Valores permitidos
validator.expect_column_values_to_be_in_set(
    "status", 
    {"pending", "processing", "shipped", "delivered", "cancelled", "returned"}
)

# Rangos numéricos
validator.expect_column_values_to_be_between("amount", min_value=0, max_value=100000)
validator.expect_column_min_to_be_between("amount", min_value=0)

# Tipos de datos
validator.expect_column_values_to_be_of_type("order_id", "int64")
validator.expect_column_values_to_match_regex("email", r"^[a-zA-Z0-9_.+-]+@[a-zA-Z0-9-]+\.[a-zA-Z0-9-.]+$")

# Integridad referencial (usando conjuntos de valores válidos)
valid_customer_ids = set(pd.read_csv('customers.csv')['customer_id'])
validator.expect_column_values_to_be_in_set("customer_id", valid_customer_ids)

# Estadísticas y distribuciones
validator.expect_column_mean_to_be_between("amount", min_value=50, max_value=500)
validator.expect_column_quantile_values_to_be_between(
    "amount",
    quantile_ranges={"quantiles": [0.1, 0.5, 0.9], "value_ranges": [[0, 50], [50, 300], [300, 5000]]}
)

# Tendencia de filas: asegurar que llegan suficientes registros
validator.expect_table_row_count_to_be_between(min_value=1000, max_value=100000)

# Guardar el Suite de Expectations
validator.save_expectation_suite("orders_expectations")

# ── VALIDAR EN PRODUCCIÓN ──────────────────────────────────────────────────
checkpoint = context.add_or_update_checkpoint(
    name="orders_checkpoint",
    validations=[{
        "batch_request": {"datasource_name": "orders_datasource"},
        "expectation_suite_name": "orders_expectations"
    }],
    action_list=[
        # Notificar a Slack si hay fallos
        {"name": "send_slack_alert", "action": {"class_name": "SlackNotificationAction",
         "slack_webhook": "{{ ENV.SLACK_WEBHOOK }}",
         "notify_on": "failure"}},
        # Actualizar el sitio de documentación
        {"name": "update_data_docs", "action": {"class_name": "UpdateDataDocsAction"}},
        # Guardar resultados de validación
        {"name": "store_validation_result", "action": {"class_name": "StoreValidationResultAction"}},
    ]
)

result = checkpoint.run()

if not result.success:
    # Extraer expectativas fallidas
    for validation_result in result.run_results.values():
        for expectation_result in validation_result.results:
            if not expectation_result.success:
                print(f"❌ FAILED: {expectation_result.expectation_config.expectation_type}")
                print(f"   Column: {expectation_result.expectation_config.kwargs.get('column')}")
                print(f"   Observed: {expectation_result.result.get('observed_value')}")
    raise ValueError("Data quality validation failed")

Soda Core

Soda ofrece una sintaxis YAML más declarativa para data quality checks. Integra directamente con el Data Warehouse sin necesitar extraer datos.

# ── SODA CHECKS en YAML (checks.yml) ──────────────────────────────────────
# Mucho más legible para analistas no-técnicos
# Se ejecuta directamente sobre la base de datos (pushdown)

# checks.yml
checks for orders:
  # Completitud
  - missing_count(order_id) = 0:
      name: "No order_id vacíos"
      fail: when > 0
  
  - missing_percent(customer_email) < 5%:
      name: "Menos de 5% de emails faltantes"
  
  # Unicidad
  - duplicate_count(order_id) = 0
  
  # Validez
  - invalid_count(status) = 0:
      valid values: [pending, processing, shipped, delivered, cancelled]
  
  - values_in_set(currency) = 0:
      valid values: [USD, EUR, GBP, CAD, MXN, BRL]
  
  - min(amount) >= 0
  - max(amount) < 100000
  
  # Tendencias (detectar anomalías)
  - row_count between 1000 and 50000
  - row_count between yesterday - 20% and yesterday + 20%:
      name: "Filas de hoy similares a ayer"
  
  # Freshness: los datos no son muy viejos
  - freshness(created_at) < 24h:
      name: "Datos de las últimas 24 horas"
  
  # Distribución estadística
  - avg(amount) between 80 and 200:
      name: "Ticket promedio en rango esperado"
  
  # Integridad referencial (cross-check)
  - values in (customer_id) must exist in customers (id):
      name: "Todos los customer_id existen en tabla customers"

# ── EJECUTAR SODA CHECKS ────────────────────────────────────────────────────
# CLI
# soda scan -d my_postgres_connection -c checks.yml

# Python API
from soda.scan import Scan

scan = Scan()
scan.set_data_source_name("warehouse")
scan.add_configuration_yaml_file("soda_config.yml")
scan.add_sodacl_yaml_file("checks.yml")
scan.execute()

# Resultados
print(scan.get_scan_results())
if scan.has_check_failures():
    raise Exception("Data quality checks failed")

dbt Tests para Data Quality

# ── DBT TESTS: Tests integrados en el pipeline de transformación ───────────
# schema.yml en tu proyecto dbt

version: 2
models:
  - name: fct_orders
    columns:
      - name: order_id
        tests:
          - unique
          - not_null
      - name: customer_id
        tests:
          - not_null
          - relationships:  # Integridad referencial
              to: ref('dim_customers')
              field: customer_id
      - name: status
        tests:
          - accepted_values:
              values: ['pending', 'shipped', 'delivered', 'cancelled']
      - name: amount
        tests:
          - not_null
          - dbt_expectations.expect_column_values_to_be_between:
              min_value: 0
              max_value: 100000

# CUSTOM TESTS en SQL
# tests/assert_no_negative_revenue.sql
SELECT order_id
FROM {{ ref('fct_orders') }}
WHERE amount < 0  -- La query debe retornar 0 filas para pasar el test

# GENERIC TEST personalizado
# macros/test_column_not_in_future.sql
{% test column_not_in_future(model, column_name) %}
  SELECT {{ column_name }}, CURRENT_TIMESTAMP
  FROM {{ model }}
  WHERE {{ column_name }} > CURRENT_TIMESTAMP
{% endtest %}

# Uso en schema.yml
# - name: order_date
#   tests:
#     - column_not_in_future

Data Contracts

Los Data Contracts son acuerdos formales entre el productor de datos y sus consumidores. Definen el schema, semántica, SLAs y calidad esperada. Son el estándar emergente para Data Mesh en 2026.

# ── DATA CONTRACT: YAML Schema ─────────────────────────────────────────────
# orders_contract.yaml
# Basado en Data Contract Specification (datacontract.com)

dataContractSpecification: 0.9.3
id: orders-api-v1
info:
  title: Orders Data Contract
  version: 2.0.0
  description: |
    Contrato de datos del dominio de Órdenes.
    Los consumidores de este contrato se comprometen a no leer directamente 
    de las tablas fuente, sino solo a través de las interfaces definidas aquí.
  owner: data-engineering@company.com
  contact:
    name: Orders Team
    url: https://teams.microsoft.com/orders-team
    email: orders-data@company.com

servers:
  production:
    type: snowflake
    account: mycompany.us-east-1
    database: PRODUCTION
    schema: GOLD

terms:
  usage: "Solo para uso interno. No compartir con terceros sin aprobación de Legal."
  limitations: "Datos de clientes sujetos a GDPR. No usar para marketing sin consentimiento."
  noticePeriod: P3M  # 3 meses de aviso para cambios breaking

models:
  orders:
    description: "Todas las órdenes de clientes procesadas."
    fields:
      order_id:
        type: long
        required: true
        unique: true
        description: "Identificador único de la orden."
      customer_id:
        type: integer
        required: true
        description: "ID del cliente. Referencia a customers.customer_id."
        pii: false
      amount:
        type: double
        required: true
        minimum: 0
        description: "Monto total de la orden en USD."
      status:
        type: string
        required: true
        enum: [pending, processing, shipped, delivered, cancelled]
      created_at:
        type: timestamp
        required: true
        description: "Timestamp UTC de creación de la orden."

quality:
  - type: sql
    description: "No debe haber order_ids nulos"
    query: "SELECT COUNT(*) FROM orders WHERE order_id IS NULL"
    mustBe: 0
  - type: sql
    description: "Cantidad de órdenes diarias esperada"
    query: "SELECT COUNT(*) FROM orders WHERE DATE(created_at) = CURRENT_DATE - 1"
    mustBeBetween: [100, 100000]

serviceLevel:
  availability: "99.9%"
  retention: "3 años"
  latency: "Los datos del día anterior disponibles antes de las 06:00 UTC"
  freshness: "Actualizado cada hora"
  support: "Respuesta en < 4 horas en días laborales"

# ── VALIDAR EL CONTRATO ────────────────────────────────────────────────────
# CLI: datacontract lint orders_contract.yaml
# CLI: datacontract test orders_contract.yaml
# CLI: datacontract publish orders_contract.yaml  # al Data Catalog

Monitoreo Continuo con Anomaly Detection

"""
Detección de anomalías en datos con Monte Carlo / código propio
"""
import pandas as pd
import numpy as np
from datetime import datetime, timedelta
from scipy import stats

class DataAnomalyDetector:
    """
    Detecta anomalías estadísticas en métricas de datos.
    Útil para detectar problemas de calidad antes de que impacten dashboards.
    """
    
    def __init__(self, historical_df: pd.DataFrame, z_score_threshold: float = 3.0):
        self.history = historical_df
        self.threshold = z_score_threshold
    
    def detect_row_count_anomaly(self, table: str, current_count: int) -> dict:
        """Detecta si el número de filas de hoy es anómalo vs historia."""
        historical_counts = self.history[
            (self.history['table'] == table) & 
            (self.history['date'] >= datetime.now() - timedelta(days=30))
        ]['row_count']
        
        if len(historical_counts) < 7:
            return {"status": "insufficient_history"}
        
        mean = historical_counts.mean()
        std = historical_counts.std()
        z_score = (current_count - mean) / std if std > 0 else 0
        
        is_anomaly = abs(z_score) > self.threshold
        
        return {
            "table": table,
            "current_count": current_count,
            "historical_mean": round(mean),
            "historical_std": round(std),
            "z_score": round(z_score, 2),
            "is_anomaly": is_anomaly,
            "severity": "critical" if abs(z_score) > 5 else "warning" if is_anomaly else "ok",
            "message": f"Row count {current_count:,} is {'anomalous' if is_anomaly else 'normal'} "
                      f"(expected ~{round(mean):,} ± {round(std):,})"
        }
    
    def detect_null_rate_anomaly(self, column_stats: pd.Series, historical_null_rates: pd.Series) -> dict:
        """Detecta cambios súbitos en la tasa de nulos."""
        current_null_rate = column_stats['null_rate']
        historical_avg = historical_null_rates.mean()
        
        # Cambio de más del 5% en tasa de nulos es sospechoso
        if current_null_rate - historical_avg > 0.05:
            return {
                "status": "anomaly",
                "type": "null_rate_spike",
                "current": f"{current_null_rate:.1%}",
                "expected": f"{historical_avg:.1%}",
                "delta": f"+{(current_null_rate - historical_avg):.1%}",
            }
        return {"status": "ok"}
    
    def run_full_scan(self, current_stats: pd.DataFrame) -> list[dict]:
        """Ejecuta detección completa y retorna lista de alertas."""
        alerts = []
        
        for _, row in current_stats.iterrows():
            result = self.detect_row_count_anomaly(row['table'], row['row_count'])
            if result.get('is_anomaly'):
                alerts.append(result)
        
        return alerts

# Integración con Airflow
from airflow.decorators import task

@task
def run_anomaly_detection(**context) -> dict:
    """Task de Airflow para ejecutar detección de anomalías."""
    current_stats = get_current_table_stats()
    historical_stats = load_historical_stats(days=30)
    
    detector = DataAnomalyDetector(historical_stats)
    alerts = detector.run_full_scan(current_stats)
    
    if alerts:
        send_slack_alert(f"⚠️ {len(alerts)} data anomalies detected:\n" + 
                         "\n".join([f"- {a['table']}: {a['message']}" for a in alerts]))
    
    return {"total_anomalies": len(alerts), "alerts": alerts}