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.
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}