Data Contracts: cómo diseñar un framework desde cero

Data Architecture
Data Engineering
Podcast
Qué son los Data Contracts, por qué los necesitás, y cómo implementé un framework en producción con Databricks y PySpark.
Autor
Publicado

14 de abril de 2026

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:

  1. El equipo de backend agrega una columna al API
  2. Otro equipo cambia el tipo de un campo de INT a STRING
  3. Un proveedor externo modifica el formato del CSV que te manda
  4. Tu pipeline de Silver falla a las 3 AM
  5. 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

Listado 1: Definición YAML de un Data Contract con schema, SLAs y reglas de calidad
# 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 atraso

Implementación en PySpark

1. Parser del contrato

Listado 2: Parser de contratos YAML a dataclasses de Python
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

Listado 3: Validador de schema: compara DataFrame contra el contrato
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 errors

3. Validador de calidad

Listado 4: Validador de calidad: ejecuta reglas SQL y reporta violaciones
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 results

4. Integración en el pipeline

Listado 5: Pipeline de ingesta con validación de contrato integrada
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 quality

Hard 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
Listado 6: Hard contract vs Soft contract: fallar o enviar a cuarentena
# 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:

Listado 7: Data Contracts como DLT Expectations en Databricks
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.