Saltar al contenido

Calidad

Menú

Puedes leer sin crear un espacio. Créalo solo cuando quieras guardar.

Guardar progreso

Expectations, cuarentena y event logs

Haz visibles las decisiones de calidad y separa observación, descarte, fallo y remediación.

Lectura pública
Detalles
Reto observable

Implementa reglas con acciones de observar, aislar y fallar, y demuestra tasas, cuarentena trazable y remediación mediante el event log.

Al terminar podrás
  • Elegir EXPECT, DROP o FAIL
  • Diseñar una cuarentena trazable
  • Consultar métricas en event logs
Prerrequisitos
m18
Última revisión
25 ago 2026
Nivel
Professional
Ruta relacionada
pipelines
Dominios blueprint
Data Transformation, Cleansing, and Quality · Monitoring
Estado
Revisión editorial interna
Fuentes principales
Manage data quality with pipeline expectations · Databricks · Pipeline event log schema · Databricks
Reportar un error
01
Modelo mental

Dimensiones y contratos de calidad

Las expectations convierten reglas de calidad en métricas y acciones declarativas: observar, descartar o fallar la actualización.

`expect` conserva tanto filas válidas como inválidas y registra métricas; es útil durante adopción o para reglas informativas. `expect_or_drop` elimina las filas que incumplen y permite continuar. `expect_or_fail` detiene el flow y revierte atómicamente la actualización afectada cuando aceptar datos incorrectos sería peor que retrasar publicación.

La acción se elige por impacto y capacidad de remediación, no por severidad nominal. Un `order_id` nulo puede ir a cuarentena si el resto del pipeline debe continuar; una violación de unicidad en un saldo regulado puede justificar fallo. Toda regla necesita owner, definición y umbral operativo.

PySparkTres acciones de calidad diferenciadas
from pyspark import pipelines as dp

@dp.table(name="orders_silver")
@dp.expect("known_currency", "currency IN ('EUR', 'USD', 'GBP')")
@dp.expect_or_drop("valid_order_id", "order_id IS NOT NULL")
@dp.expect_or_fail("non_negative_amount", "amount >= 0")
def orders_silver():
    return spark.readStream.table("orders_bronze")

Aplicar tres acciones al mismo dataset solo es correcto si negocio ha decidido explícitamente el tratamiento de cada violación.

¿Qué acción conserva filas inválidas pero produce métricas?

Profundiza

Una expectation es una regla booleana aplicada a cada fila que combina observación con una política de respuesta. La misma expresión puede conservar filas y registrar métricas, descartar las inválidas o fallar la actualización. La elección no expresa severidad estética, sino el daño de publicar el dato y la capacidad de remediarlo. `warn` —comportamiento de retención— sirve para medir y explorar; `drop` evita contaminar el target cuando perder esas filas está aceptado y existe trazabilidad; `fail` protege invariantes cuya violación invalida todo el resultado. En una actualización fallida, la transacción del flow se revierte, pero el alcance sobre flows paralelos y dependientes varía según el modo del pipeline. Además, las métricas de `fail` tienen limitaciones porque el update no se confirma como una ejecución normal. La regla necesita nombre estable, owner, umbral y ruta de investigación.

Expectation

Restricción nombrada basada en una expresión booleana que evalúa calidad durante el procesamiento y registra o aplica una acción configurada.

Integra controles de calidad con métricas y transacciones del pipeline en lugar de depender de comprobaciones posteriores aisladas.
Retain, drop, fail

Tres políticas que respectivamente conservan y miden, excluyen filas inválidas, o abortan la actualización al detectar una violación.

Permiten alinear cada regla con impacto, remediación y disponibilidad, evitando usar una única respuesta para toda anomalía.
Lógica de tres valores

Semántica SQL donde una expresión con null puede resultar unknown en vez de verdadero o falso explícito.

Obliga a formular constraints de nulabilidad cuidadosamente para no clasificar datos de forma distinta a la intención.
Resumen

Puntos clave

  • `expect` mide sin eliminar; `expect_or_drop` continúa sin la fila; `expect_or_fail` aborta el flow afectado.
  • Los nombres de expectation deben ser estables y describir el contrato.
  • Fallar una actualización protege el target, pero puede consumir el presupuesto de frescura.

Evita

  • Usar `expect_or_fail` para toda anomalía y convertir una fila reparable en una caída completa del flow.
  • Usar `expect` para una clave obligatoria y permitir que inválidos lleguen silenciosamente a consumidores.
02
Implementación

EXPECT y observación

Una cuarentena útil conserva la fila, procedencia, regla incumplida y versión para que pueda corregirse y reprocesarse.

Descartar con una expectation no crea automáticamente una tabla de cuarentena. El patrón explícito clasifica una vista común en salida válida e inválida. La rama inválida añade `failure_reasons`, timestamp, source file u offset y versión de pipeline; la rama válida aplica las mismas condiciones complementarias.

La cuarentena necesita política de acceso porque puede contener PII, retención y un flujo de remediación. Reinyectar datos corregidos directamente en silver salta trazabilidad; conviene publicar una nueva entrada con identidad estable o una tabla de correcciones que pase por las mismas reglas.

PySparkClasificación con motivos de rechazo
from pyspark.sql import functions as F

classified = (
    spark.readStream.table("orders_bronze")
      .withColumn(
          "failure_reasons",
          F.array_compact(F.array(
              F.when(F.col("order_id").isNull(), F.lit("MISSING_ORDER_ID")),
              F.when(F.col("amount") < 0, F.lit("NEGATIVE_AMOUNT")),
              F.when(F.col("event_ts").isNull(), F.lit("INVALID_EVENT_TS")),
          ))
      )
)

valid = classified.where("size(failure_reasons) = 0")
quarantine = classified.where("size(failure_reasons) > 0")

Usa una función compartida para que la condición válida sea exactamente el complemento de la cuarentena.

¿Qué propiedad evita perder o duplicar filas entre silver y cuarentena?

Profundiza

Cuarentena no significa una carpeta donde los datos malos desaparecen; es un producto operativo con identidad, procedencia, causa, estado y ruta de reingreso. El patrón más explicable evalúa reglas una vez, añade un mapa o array de violaciones y divide el DataFrame en válido e inválido. La fila de cuarentena conserva payload original, campos parseados, fuente, coordenada, versión de contrato, tiempo de detección y nombres de reglas. Una reparación produce una nueva versión o evento, no modifica silenciosamente la evidencia. El reingreso utiliza la misma clave de negocio e idempotencia para evitar duplicar el target. Los datos pueden contener PII, por lo que cuarentena necesita controles de acceso y retención al menos tan estrictos como producción. Sus métricas revelan deuda: volumen entrante, edad, porcentaje corregido y causas recurrentes con owner.

Provenance

Metadatos que identifican origen, posición, versión y transformación mediante los cuales una fila llegó a la decisión de cuarentena.

Permite reproducir el fallo, localizar productores responsables y demostrar que una corrección corresponde al registro exacto.
Estado de remediación

Ciclo explícito de una anomalía, por ejemplo pendiente, investigada, corregida, descartada o reingresada con referencia a evidencia.

Convierte cuarentena en una cola gobernada y permite medir deuda y cumplimiento de tiempos de resolución.
Reingreso idempotente

Proceso que devuelve una fila corregida al flujo canónico sin crear más de una entidad o aplicar dos veces el mismo cambio.

Evita que solucionar calidad introduzca duplicados y conserva una historia auditable de la reparación.
Resumen

Puntos clave

  • Válidos e inválidos deben derivarse de una única clasificación para evitar huecos.
  • Cada fila inválida conserva evidencia suficiente para diagnóstico y replay.
  • La remediación vuelve a entrar por una frontera gobernada y evita doble conteo.

Evita

  • Definir filtros independientes para válidos e inválidos y dejar filas que no caen en ninguna rama.
  • Guardar solo un conteo de errores, sin payload ni procedencia para repararlos.
03
Operación

EXPECT OR DROP

El pipeline event log es la fuente estructurada para progreso, calidad, linaje y errores, y debe consultarse usando campos documentados.

Los eventos `flow_progress` incluyen estado, métricas y `data_quality` dentro de `details`. `origin` identifica pipeline, update y flow. La función `event_log(TABLE(...))` permite consultar el log asociado a una tabla del pipeline desde SQL y construir tendencias por expectation.

No todo campo interno es contrato público. Las consultas operativas seleccionan campos documentados y toleran ausencia de métricas en eventos que no sean de progreso. Conservar `update_id`, nombre de flow y timestamp permite relacionar una caída de calidad con despliegue y lote.

SQLInspección de calidad en el event log
SELECT
  timestamp,
  origin.update_id AS update_id,
  origin.flow_name AS flow_name,
  details:flow_progress.status::STRING AS status,
  details:flow_progress.data_quality:expectations AS expectations
FROM event_log(TABLE(main.silver.orders_silver))
WHERE event_type = 'flow_progress'
ORDER BY timestamp DESC
LIMIT 100;

Normaliza el array de expectations en una vista operativa para calcular tasas por regla y actualización.

¿Qué evento contiene normalmente las métricas de expectations?

Profundiza

El pipeline event log es la bitácora estructurada de una ejecución declarativa. Registra eventos de actualización, flows, progreso, calidad, linaje, configuración y errores con un schema documentado y campos JSON para detalles. No es una tabla de negocio ni conviene depender de campos internos no documentados, porque pueden cambiar. La lectura comienza por identificar update y flow, ordenar por la secuencia del evento y extraer únicamente estructuras soportadas. Una fila aislada rara vez cuenta la historia completa: se correlacionan inicio, progreso y final de la misma actualización. Publicar el event log como tabla de Unity Catalog facilita permisos, retención y consultas cross-pipeline. Las expectations `warn` y `drop` producen métricas consultables; una violación `fail` aborta y puede no registrar contadores equivalentes, por lo que se combina el error del flow con datos de entrada y logs.

Update

Instancia identificable de actualización del pipeline que agrupa planificación, ejecución de flows y resultado final bajo una misma operación.

Es la unidad correcta para correlacionar eventos y evitar mezclar métricas de ejecuciones concurrentes o sucesivas.
Event sequence

Metadato estructurado que permite ordenar y relacionar eventos distribuidos del pipeline más fiablemente que un timestamp aislado.

Hace posible reconstruir causalidad durante fallos y distinguir progreso anterior de mensajes posteriores de cierre.
Campo documentado

Atributo del schema que Databricks declara apto para consumo de clientes y cuya semántica está publicada oficialmente.

Reduce roturas al evitar dashboards dependientes de detalles internos que pueden cambiar sin contrato público.
Resumen

Puntos clave

  • Filtra `event_type = 'flow_progress'` antes de interpretar métricas de flow.
  • El JSON `details` se analiza con rutas documentadas y tipos explícitos.
  • Event log complementa, no sustituye, una reconciliación de negocio.

Evita

  • Consultar cualquier evento como si contuviera `flow_progress` y generar nulos difíciles de interpretar.
  • Construir dependencias permanentes sobre campos internos no documentados del JSON.
04
Diagnóstico

EXPECT OR FAIL

Un contrato de calidad define dimensión, expresión, acción, umbral, owner y procedimiento de remediación antes de escribir código.

Validez, completitud, unicidad, consistencia, puntualidad y exactitud requieren evidencias diferentes. `amount >= 0` comprueba validez, pero no exactitud frente al sistema de pagos. Una expectation fila a fila tampoco demuestra unicidad global sin una transformación o agregación adecuada.

Las reglas evolucionan como código versionado. Un cambio de umbral debe revisar impacto histórico y despliegue; renombrar una expectation rompe series de métricas. Las reglas críticas pueden agruparse con `expect_all_or_fail`, mientras reglas de observación usan `expect_all` para una configuración compartida.

PythonDiccionario versionable de reglas
structural_rules = {
    "order_id_present": "order_id IS NOT NULL",
    "event_ts_present": "event_ts IS NOT NULL",
    "amount_non_negative": "amount >= 0",
}

business_observations = {
    "known_currency": "currency IN ('EUR', 'USD', 'GBP')",
    "reasonable_amount": "amount <= 100000",
}

Agrupar expresiones facilita reutilización, pero documenta por separado la acción y el owner de cada conjunto.

¿Puede `amount >= 0` demostrar que el importe cobrado es exacto?

Profundiza

Un contrato de calidad conecta significado de negocio con una expresión ejecutable y una respuesta operativa. Por cada regla documenta dimensión —validez, completitud, unicidad, consistencia, puntualidad—, ámbito, expresión, tolerancia, acción, owner, evidencia y procedimiento de remediación. Una expectation solo implementa la parte por fila; una clave única global, reconciliación entre tablas o frescura requieren controles agregados adicionales. Los umbrales deben basarse en riesgo: cero puede ser correcto para identidad primaria, pero absurdo para un atributo opcional con fuente imperfecta. Versionar el contrato permite explicar por qué cambió una tasa y revalidar historia. Antes de escribir código se prueban ejemplos límite, nulls, zonas horarias y evolución de schema. Separar regla de acción posibilita observar primero una nueva constraint, calibrarla y endurecerla sin modificar su significado.

Dimensión de calidad

Categoría semántica que describe qué propiedad se evalúa, como completitud, validez, consistencia, unicidad o puntualidad del dato.

Evita listas inconexas de expresiones y ayuda a comprobar que el producto cubre riesgos relevantes.
Denominador

Población exacta sobre la que se calcula una tasa de cumplimiento, con inclusiones, exclusiones y ventana temporal definidas.

Impide métricas engañosas cuyo porcentaje cambia por mezclar datos de prueba, replays o segmentos no comparables.
Contrato versionado

Especificación identificable de reglas, umbrales, acciones y ownership válida para una versión del producto de datos.

Permite auditar cambios, ejecutar reglas en sombra y atribuir variaciones a datos o a definiciones.
Resumen

Puntos clave

  • Cada expresión debe medir la dimensión que afirma medir.
  • Acción y umbral forman parte del contrato, no son detalles de implementación.
  • Nombres estables permiten comparar calidad entre versiones y actualizaciones.

Evita

  • Llamar 'exactitud' a una comprobación de formato que no compara con una fuente autoritativa.
  • Cambiar nombres de reglas en cada release y perder continuidad del indicador.
05
Decisión de diseño

Quarantine pattern y event log

La operación combina enforcement en el pipeline con Data Quality Monitoring para observar frescura, completitud, distribución y drift en activos de Unity Catalog.

Un conteo absoluto de inválidos confunde crecimiento de volumen con degradación. La tasa `failed / (passed + failed)` por expectation, flow y ventana permite comparar. Expectations observan, descartan o fallan registros dentro del pipeline: son enforcement vinculado al código y al update.

Data Quality Monitoring complementa ese contrato sin modificar las tablas monitorizadas ni añadir trabajo al job productor: anomaly detection aprende patrones históricos para freshness y completeness, y data profiling calcula estadísticas y drift. Se ejecuta en serverless con coste propio. Úsalo para descubrir degradaciones transversales; conserva expectations o reconciliaciones cuando una regla debe bloquear o aislar datos.

SQLTasa de violación por regla
SELECT
  window_start,
  flow_name,
  expectation_name,
  failed_records,
  passed_records,
  failed_records / NULLIF(failed_records + passed_records, 0) AS failure_rate
FROM main.ops.pipeline_expectation_metrics
WHERE window_start >= current_timestamp() - INTERVAL 24 HOURS
ORDER BY failure_rate DESC;

Define el umbral y la duración mínima de incumplimiento en el runbook, no dentro de una consulta ad hoc.

¿Qué diferencia operacional separa Data Quality Monitoring de una expectation crítica?

Profundiza

Operar calidad significa convertir eventos y reglas en decisiones sostenidas. Los contadores de una expectation se transforman en tasas usando passed más failed como denominador, se agregan por update y se comparan con baselines y SLO. Una sola fila inválida puede ser crítica si afecta identidad; millones pueden ser tolerables si pertenecen a un campo opcional durante una migración acordada. Por eso alertas combinan severidad, proporción, volumen absoluto, duración y segmento. El event log muestra qué ocurrió dentro del pipeline; una tabla de calidad estable conserva tendencias, owners y estado de incidente. El ciclo completo incluye detectar, contener, diagnosticar, remediar, reingresar y prevenir recurrencia. Cambiar una expectation de drop a warn para hacer verde un Job no resuelve el problema: consume o redefine un riesgo y requiere aprobación del contrato.

Tasa de violación

Proporción de registros fallidos respecto al total evaluado para una regla, update y población comparables claramente identificados.

Normaliza volúmenes, pero debe acompañarse de conteo absoluto y criticidad para valorar el impacto real.
Presupuesto de calidad

Cantidad acordada de incumplimiento tolerable durante una ventana antes de detener, degradar o escalar el producto.

Hace explícito el equilibrio entre disponibilidad y corrección y evita decisiones improvisadas durante incidentes.
Modo degradado

Estado operativo definido que mantiene parte del servicio mientras etiqueta, limita o retrasa resultados afectados por una anomalía conocida.

Puede preservar utilidad sin presentar datos incompletos como normales, siempre que consumidores comprendan la señal.
Resumen

Puntos clave

  • Expectations aplican acciones; Data Quality Monitoring observa activos.
  • Anomaly detection cubre frescura/completitud; profiling, estadísticas y drift.
  • Cada alerta conserva owner, coste, contexto y criterio de cierre.

Evita

  • Sustituir una regla crítica de enforcement por una anomalía observacional que no bloquea publicación.
  • Activar profiling masivo sin owner, coste serverless o control de acceso a sus tablas métricas.

Fuente revisada · vista externa

Lakeflow Declarative Pipelines examples

Lakeflow Declarative Pipelines examples · commit 1d8b163

python/Retail Sales.py

Lectura en GitHub

Este notebook se abre desde su fuente revisada

El repositorio no permite republicar su contenido dentro de Lakehouse Lab. Conservamos la misma experiencia lateral, la ruta exacta y el commit auditado, y dejamos la lectura en GitHub para respetar la autoría.

Autor
Databricks
Licencia
No verificada
Formato
repository
Ver notebook en GitHub ↗

Módulo 19

Contenido del módulo