Data Contracts: cómo diseñar un framework desde cero
Si alguna vez te rompieron un pipeline porque alguien cambió una columna en la fuente sin avisarte, este post es para vos.
Los Data Contracts son la solución a un problema que todos los ingenieros de datos enfrentamos: las fuentes cambian sin aviso y los pipelines se rompen en silencio.
El problema real
Situación típica en cualquier empresa:
- El equipo de backend agrega una columna al API
- Otro equipo cambia el tipo de un campo de
INTaSTRING - Un proveedor externo modifica el formato del CSV que te manda
- Tu pipeline de Silver falla a las 3 AM
- Te enterás cuando el dashboard del CEO muestra datos vacíos
Sin Data Contracts: te enterás cuando algo se rompe. Con Data Contracts: te enterás antes de que llegue a producción.
Qué es un Data Contract
Un Data Contract es un acuerdo formal entre el productor de datos (quien genera o envía datos) y el consumidor (quien los procesa). Define:
- Schema: qué columnas, qué tipos, qué es nullable
- SLAs: cuándo llegan los datos, con qué frecuencia
- Calidad: reglas de validación (no nulos, rangos, unicidad)
- Ownership: quién es responsable si algo falla
- Versionamiento: cómo se manejan los cambios
Anatomía de un Data Contract
# contracts/transactions.yml
contract:
name: transactions
version: "2.1"
owner: backend-team
description: "Transacciones de pagos del core bancario"
sla:
freshness: "1 hour"
availability: "99.9%"
schema:
- name: transaction_id
type: BIGINT
nullable: false
unique: true
description: "ID único de la transacción"
- name: customer_id
type: BIGINT
nullable: false
description: "FK al cliente"
- name: amount
type: DECIMAL(18,2)
nullable: false
checks:
- "amount > 0"
- "amount < 1000000"
- name: transaction_date
type: DATE
nullable: false
checks:
- "transaction_date >= '2020-01-01'"
- "transaction_date <= current_date()"
- name: status
type: STRING
nullable: false
allowed_values: ["completed", "pending", "failed", "reversed"]
- name: currency
type: STRING
nullable: false
pattern: "^[A-Z]{3}$"
quality_rules:
- name: no_duplicates
sql: "SELECT COUNT(*) - COUNT(DISTINCT transaction_id) FROM {table}"
threshold: 0
- name: completeness
sql: "SELECT COUNT(*) FILTER(WHERE amount IS NULL) / COUNT(*) FROM {table}"
threshold: 0.01 # máximo 1% de nulos
- name: freshness
sql: "SELECT DATEDIFF(hour, MAX(transaction_date), current_date()) FROM {table}"
threshold: 26 # máximo 26 horas de atrasoImplementación en PySpark
1. Parser del contrato
import yaml
from dataclasses import dataclass
from typing import List, Optional
@dataclass
class ColumnContract:
name: str
type: str
nullable: bool = True
unique: bool = False
checks: Optional[List[str]] = None
allowed_values: Optional[List[str]] = None
pattern: Optional[str] = None
@dataclass
class QualityRule:
name: str
sql: str
threshold: float
@dataclass
class DataContract:
name: str
version: str
owner: str
columns: List[ColumnContract]
quality_rules: List[QualityRule]
@classmethod
def from_yaml(cls, path: str) -> "DataContract":
with open(path) as f:
raw = yaml.safe_load(f)["contract"]
columns = [
ColumnContract(**col)
for col in raw["schema"]
]
rules = [
QualityRule(**rule)
for rule in raw.get("quality_rules", [])
]
return cls(
name=raw["name"],
version=raw["version"],
owner=raw["owner"],
columns=columns,
quality_rules=rules
)2. Validador de schema
from pyspark.sql import DataFrame
from pyspark.sql.types import *
TYPE_MAP = {
"BIGINT": LongType(),
"STRING": StringType(),
"DATE": DateType(),
"DECIMAL(18,2)": DecimalType(18, 2),
"BOOLEAN": BooleanType(),
"TIMESTAMP": TimestampType(),
}
def validate_schema(df: DataFrame, contract: DataContract) -> List[str]:
"""Valida que el DataFrame cumpla el schema del contrato."""
errors = []
df_fields = {f.name: f for f in df.schema.fields}
for col in contract.columns:
if col.name not in df_fields:
errors.append(f"Columna faltante: {col.name}")
continue
field = df_fields[col.name]
expected_type = TYPE_MAP.get(col.type)
if expected_type and field.dataType != expected_type:
errors.append(
f"Tipo incorrecto en {col.name}: "
f"esperado {col.type}, encontrado {field.dataType}"
)
if not col.nullable and field.nullable:
errors.append(
f"{col.name} debería ser NOT NULL"
)
# Columnas extra (warning, no error)
contract_cols = {c.name for c in contract.columns}
extra = set(df_fields.keys()) - contract_cols
if extra:
errors.append(f"Columnas no esperadas: {extra}")
return errors3. Validador de calidad
def validate_quality(
spark,
table_name: str,
contract: DataContract
) -> List[dict]:
"""Ejecuta las reglas de calidad y reporta violaciones."""
results = []
for rule in contract.quality_rules:
query = rule.sql.format(table=table_name)
value = spark.sql(query).collect()[0][0]
passed = value <= rule.threshold
results.append({
"rule": rule.name,
"value": value,
"threshold": rule.threshold,
"passed": passed
})
if not passed:
print(f"FAIL: {rule.name} = {value} "
f"(threshold: {rule.threshold})")
return results4. Integración en el pipeline
class ContractViolation(Exception):
"""Excepción para violaciones de Data Contracts."""
pass
def alert_team(owner: str, failures: list):
"""Notifica al equipo dueño del contrato sobre las fallas de calidad."""
# Implementar según tu stack: Slack webhook, email, PagerDuty, etc.
for f in failures:
print(f"[ALERT → {owner}] {f['rule']}: valor={f['value']}, "
f"threshold={f['threshold']}")
def ingest_with_contract(
spark,
source_df: DataFrame,
contract_path: str,
target_table: str
):
"""Pipeline de ingesta con validación de contrato."""
contract = DataContract.from_yaml(contract_path)
# 1. Validar schema
schema_errors = validate_schema(source_df, contract)
if schema_errors:
raise ContractViolation(
f"Schema violation en {contract.name}: "
f"{schema_errors}"
)
# 2. Escribir a tabla
source_df.write.mode("append").saveAsTable(target_table)
# 3. Validar calidad post-write
quality = validate_quality(spark, target_table, contract)
failures = [r for r in quality if not r["passed"]]
if failures:
# Alertar pero no fallar (soft contract)
alert_team(contract.owner, failures)
return qualityHard Contracts vs Soft Contracts
| Tipo | Qué pasa si falla | Cuándo usarlo |
|---|---|---|
| Hard | Pipeline se detiene, dato no entra | Fuentes críticas (core bancario, pagos) |
| Soft | Se loguea warning, dato entra igual | Fuentes externas, datos no críticos |
# Hard contract: falla y para
if schema_errors:
raise ContractViolation(...)
# Soft contract: loguea y sigue
if schema_errors:
log_violation(contract, schema_errors)
# dato entra igual a una tabla de quarantine
source_df.write.saveAsTable(f"{target_table}_quarantine")Con Databricks Expectations (DLT)
Si usás DLT / Lakeflow Declarative Pipelines, podés expresar los contracts como expectations:
import dlt
@dlt.table(name="silver_transactions")
@dlt.expect_all_or_drop({
"valid_amount": "amount > 0 AND amount < 1000000",
"valid_status": "status IN ('completed','pending','failed','reversed')",
"not_null_id": "transaction_id IS NOT NULL",
"valid_date": "transaction_date >= '2020-01-01'"
})
def clean_transactions():
return dlt.read("bronze_transactions")expect_all_or_drop es un hard contract: las filas que no cumplen se descartan. expect_all_or_fail para el pipeline entero. expect_all solo loguea.
Links
Data Contracts (Andrew Jones) — referencia de la comunidad
DLT Expectations (Databricks) — documentación oficial
Próxima semana: Data Mesh en la práctica — lo que funciona y lo que no.