Módulo 6: Freshness Volume And Lineage

Escribiendo un check de freshness para S04

Descripción

Esta es la lección central del módulo: escribe check_freshness(), la primera función de nivel de archivo que corre esta guía completa, y la corre de verdad sobre orders_2026-08-14.csv. El resultado no es una sorpresa —la lección 5 del módulo 1 ya lo anticipó con evidencia—, pero sí es la primera vez que ese resultado sale de una función reusable, con Polars, en vez de un script de una sola vez con aritmética hecha a mano.

Conexión con el módulo. Las lecciones 2 y 3 construyeron los dos argumentos que esta lección junta: freshness se mide sobre el archivo completo (lección 2), usando una referencia de tiempo fija, nunca el reloj real (lección 3). check_freshness(df, run_at, sla_hours) es la síntesis exacta de esos dos argumentos, convertida en código que cualquier pipeline de Kiosko podría importar y reusar.

Una analogía: del cajero que lee la etiqueta a mano, al lector automático

Piensa en dos formas de revisar la fecha de vencimiento de un producto en una tienda. La primera: un cajero, cada vez que alguien compra algo, toma el producto, busca la fecha impresa en la etiqueta, la lee, y la compara a ojo contra la fecha de hoy. Funciona, pero depende de que ese cajero específico se acuerde de mirar, y de que lea bien la fecha cada vez —exactamente lo que hizo el script de la lección 5 del módulo 1: un cálculo hecho a mano, una sola vez, para un caso específico—. La segunda forma: un lector automático en la caja registradora, que escanea el código de barras, extrae la fecha de vencimiento codificada en él, la compara contra la fecha del sistema, y muestra un semáforo —verde o rojo— sin que ningún cajero tenga que leer ni calcular nada.

check_freshness() es ese lector automático. No reemplaza el conocimiento de la lección 5 del módulo 1 —la misma idea de comparar una fecha contra un límite—, lo convierte en una pieza que se puede llamar sobre cualquier DataFrame con una columna de tiempo, sin que nadie tenga que reescribir la aritmética cada vez. Y, como cualquier lector automático bien diseñado, funciona igual sin importar qué producto (o qué archivo) le pongas delante — la lección de hoy lo prueba primero con datos de juguete, después con el archivo real de S04.

Ejemplo trabajado: construir, probar en limpio, correr sobre S04

Paso 1 — la función

# checks.py
from datetime import datetime

import duckdb
import polars as pl

PIPELINE_RUN_AT = "2026-08-16T09:00:00"


def check_freshness(df: pl.DataFrame, run_at: str, sla_hours: int, timestamp_col: str = "order_ts") -> dict:
    """Compara el order_ts mas reciente del DataFrame contra run_at. Falla si supera sla_hours."""
    latest_ts = df.select(pl.col(timestamp_col).max()).item()
    run_at_dt = datetime.fromisoformat(run_at)
    hours_since_latest = (run_at_dt - latest_ts).total_seconds() / 3600
    return {
        "check": "freshness",
        "latest_row_ts": str(latest_ts),
        "run_at": run_at,
        "sla_hours": sla_hours,
        "hours_since_latest": round(hours_since_latest, 2),
        "status": "PASS" if hours_since_latest <= sla_hours else "FAIL",
    }

Cinco líneas de cuerpo, cada una con un rol preciso. df.select(pl.col(timestamp_col).max()).item() es la traducción directa de la lección 2: en vez de recorrer fila por fila, una sola expresión de Polars encuentra el order_ts más reciente de todo el DataFrame, y .item() lo extrae como un valor de Python normal (un datetime.datetime), no como un DataFrame de una celda. datetime.fromisoformat(run_at) convierte la cadena de texto que recibe la función —nunca datetime.now(), la regla completa de la lección 3— en un objeto comparable. La resta y la división por 3600 son la misma aritmética de horas que ya usó freshness_preview.py del módulo 1. Y el dict de salida no es un capricho de estilo: cada campo documenta, en el resultado mismo, con qué datos se calculó el veredicto —latest_row_ts, run_at, sla_hours— sin que quien lo lea tenga que adivinar ni volver al código fuente.

Paso 2 — probarla con datos de juguete, antes de tocar S04

Antes de confiar esta función al archivo real, confírmala con un caso mínimo, construido a mano, donde ya sabes de antemano cuál debería ser la respuesta:

# checks.py -- continuacion
if __name__ == "__main__":
    from datetime import datetime as dt

    toy_df = pl.DataFrame({
        "order_id": ["T1", "T2"],
        "order_ts": [dt(2026, 1, 1, 10, 0, 0), dt(2026, 1, 1, 12, 0, 0)],
    })

    print("=== check_freshness sobre datos de juguete ===")
    print("Caso PASS (2 horas de retraso, SLA de 24):")
    print(f"  {check_freshness(toy_df, run_at='2026-01-01T14:00:00', sla_hours=24)}")

    print("\nCaso FAIL (50 horas de retraso, SLA de 24):")
    print(f"  {check_freshness(toy_df, run_at='2026-01-03T14:00:00', sla_hours=24)}")

Qué esperar. Al correr python3 checks.py hasta este punto, la salida es exactamente esta:

=== check_freshness sobre datos de juguete ===
Caso PASS (2 horas de retraso, SLA de 24):
  {'check': 'freshness', 'latest_row_ts': '2026-01-01 12:00:00', 'run_at': '2026-01-01T14:00:00', 'sla_hours': 24, 'hours_since_latest': 2.0, 'status': 'PASS'}

Caso FAIL (50 horas de retraso, SLA de 24):
  {'check': 'freshness', 'latest_row_ts': '2026-01-01 12:00:00', 'run_at': '2026-01-03T14:00:00', 'sla_hours': 24, 'hours_since_latest': 50.0, 'status': 'FAIL'}

Dos filas de juguete, T1 a las 10:00 y T2 a las 12:00latest_row_ts toma la más reciente de las dos (12:00:00), tal como exige la lección 2: no importa cuántas filas tenga el DataFrame, solo importa la más nueva. Con run_at dos horas después de esa fila más reciente, PASS; con run_at dos días después, FAIL. La función se comportó exactamente como se esperaba en ambos casos extremos, antes de arriesgar ninguna conclusión sobre el archivo real de S04.

Paso 3 — corrida sobre orders_2026-08-14.csv, de verdad

# checks.py -- continuacion
    con = duckdb.connect("kiosko.duckdb")
    con.execute("""
        CREATE OR REPLACE TABLE orders_s04 AS
        SELECT * FROM read_csv('orders_2026-08-14.csv', header=True,
            columns={
                'order_id': 'VARCHAR', 'store_id': 'VARCHAR', 'product_id': 'VARCHAR',
                'quantity': 'BIGINT', 'unit_price': 'DOUBLE', 'order_ts': 'TIMESTAMP'
            })
    """)
    df = con.sql("SELECT * FROM orders_s04").pl()
    print(f"\norders_s04: {df.height} filas\n")

    freshness_result = check_freshness(df, run_at=PIPELINE_RUN_AT, sla_hours=24)
    print("=== check_freshness(df, run_at=PIPELINE_RUN_AT, sla_hours=24) ===")
    for k, v in freshness_result.items():
        print(f"  {k}: {v}")

Qué esperar. Con orders_2026-08-14.csv (las doce líneas exactas del módulo 1) en la misma carpeta, la salida adicional es exactamente esta:

orders_s04: 12 filas

=== check_freshness(df, run_at=PIPELINE_RUN_AT, sla_hours=24) ===
  check: freshness
  latest_row_ts: 2026-08-14 09:25:00
  run_at: 2026-08-16T09:00:00
  sla_hours: 24
  hours_since_latest: 47.58
  status: FAIL

Ahí está: la primera vez, en esta guía completa, que un check de nivel de archivo corre de verdad y falla — no calculado en prosa, no una predicción, un dict real, devuelto por una función real. latest_row_ts: 2026-08-14 09:25:00 es ORD-9511, la venta más tardía de las doce que trae el archivo. hours_since_latest: 47.58 es el tiempo transcurrido entre esa venta y PIPELINE_RUN_AT, casi el doble del SLA de 24 horas pactado. status: FAIL.

Por qué este número (47.58) no es el mismo que el de la lección 5 del módulo 1 (33.0)

Si comparaste este resultado contra el de freshness_preview.py, del módulo 1, vas a notar una diferencia real: ese script reportó 33.0 horas por encima del SLA; check_freshness() reporta 47.58 horas desde la fila más reciente. Los dos números son correctos — miden preguntas distintas, a propósito. El script del módulo 1 comparaba PIPELINE_RUN_AT contra SLA_DEADLINE (EXPECTED_ARRIVAL + 24 horas = 2026-08-15T00:00:00), una fecha calculada de antemano, basada en cuándo alguien en Kiosko esperaba que llegara el archivo — un dato que solo un humano, con contexto de negocio, podía saber por adelantado. check_freshness(), en cambio, no necesita que nadie le diga de antemano "cuándo se esperaba" nada: solo mira los datos que el archivo de verdad trae, encuentra el más reciente (09:25:00 del 14 de agosto), y lo compara directo contra run_at.

Esta diferencia no es un error — es, precisamente, lo que hace que check_freshness() sea reusable más allá de S04. El cálculo del módulo 1 solo funciona si alguien mantiene, a mano, un calendario de "cuándo se espera cada archivo" para cada tienda de Kiosko. check_freshness() funciona sobre cualquier tabla con una columna de tiempo, sin ese calendario — la misma definición, verificada de forma independiente, que usa dbt source freshness en producción: comparar el dato más reciente disponible contra el momento de la revisión, sin necesitar saber de antemano cuándo "debería" haber llegado. Los dos números —33.0 y 47.58— cuentan la misma historia (el archivo llegó tarde), con dos varas de medir distintas y ambas legítimas.

Diagrama: de la lección 1 del módulo 1 a esta lección

flowchart LR
    A["Modulo 1, L5:\nfreshness_preview.py\nEXPECTED_ARRIVAL + SLA fijos\na mano, un solo uso"] --> B["33.0 horas sobre el SLA\n(prosa, no reusable)"]
    C["Modulo 6, L2-L3:\nfreshness es de archivo,\nel reloj tiene que ser fijo"] --> D["check_freshness(df, run_at, sla_hours)\nfuncion real, reusable"]
    D --> E["Probada en datos de juguete:\nPASS y FAIL, ambos correctos"]
    E --> F["Corrida sobre S04 real:\nlatest_row_ts=09:25:00\nhours_since_latest=47.58\nFAIL"]

Profundización: qué significa, en la práctica, que este check "falle"

Vale la pena ser preciso sobre qué implica, en este punto de la guía, que check_freshness() devuelva status: FAIL. No significa que el archivo se descarte, ni que las doce filas dejen de procesarse — eso sería repetir el error de foundations M7, el rechazo total que ya criticó la lección 5 del módulo 1. Significa, con precisión, que Kiosko ahora tiene evidencia estructurada —un dict, no una sensación— de que este archivo violó un acuerdo de puntualidad, disponible para que un sistema de decisión posterior (el módulo 7 de esta guía, con quarantine() y raise_alert()) decida qué hacer con esa información. check_freshness(), por sí sola, solo detecta y reporta — la misma división de responsabilidades que ya sostuvo cada una de las cinco herramientas anteriores de esta guía: una función que decide "¿pasa o no pasa?", separada de otra función que decide "¿y ahora qué hacemos con eso?".

Errores comunes

Usar .min() en vez de .max() para encontrar el order_ts de referencia. Qué pasa: alguien, sin pensarlo con cuidado, usa pl.col(timestamp_col).min() en vez de .max(), y obtiene un resultado de freshness mucho peor de lo real (comparando contra la venta más antigua del archivo, no la más reciente). Por qué pasa: min/max son intercambiables sintácticamente, y sin conectar el código con la pregunta de negocio, es fácil elegir el equivocado por error de tecleo. Cómo detectarlo: si tu latest_row_ts es la fila más antigua del archivo (para S04, sería ORD-9501 a las 08:05:00, no ORD-9511 a las 09:25:00), invertiste el criterio. Cómo corregirlo: recuerda la pregunta exacta que responde freshness — "¿qué tan al día está el archivo respecto a su dato más nuevo?" — siempre .max(), nunca .min(), sobre la columna de tiempo.

Olvidar que run_at es una cadena de texto ISO, y pasarle directamente un objeto datetime. Qué pasa: alguien llama a check_freshness(df, run_at=datetime(2026, 8, 16, 9, 0, 0), sla_hours=24), pasando un objeto datetime en vez de la cadena "2026-08-16T09:00:00" que la función espera, y obtiene un AttributeError dentro de datetime.fromisoformat(run_at) (que espera una cadena, no otro datetime). Por qué pasa: PIPELINE_RUN_AT se ve, a simple vista, como si pudiera ser cualquiera de los dos tipos, y la función no valida explícitamente cuál recibió. Cómo detectarlo: si el error menciona fromisoformat y el argumento que le pasaste no es una cadena de texto, revisa el tipo exacto de tu run_at. Cómo corregirlo: respeta la firma exacta de esta lección — run_at: str, siempre una cadena en formato ISO 8601, exactamente como está definida PIPELINE_RUN_AT desde el módulo 1. Si en tu propio código prefieres trabajar con objetos datetime directamente, el Ejercicio 2 de la lección 3 de este módulo ya mostró cómo adaptar la función para ese caso.

Confundir sla_hours=24 con "el archivo tiene 24 horas para llegar desde que se genera". Qué pasa: alguien interpreta el parámetro sla_hours como si describiera cuánto tiempo, en total, puede tardar un archivo entre que se genera y que llega a Kiosko — cuando en realidad describe algo más específico: cuánto tiempo puede pasar entre el dato más reciente dentro del archivo y el momento en que alguien lo revisa. Por qué pasa: las dos interpretaciones suenan parecidas en prosa suelta. Cómo detectarlo: revisa qué compara exactamente hours_since_latestrun_at (cuándo se revisó) contra latest_row_ts (la venta más reciente que el archivo contiene), no contra ningún concepto de "cuándo se generó el archivo como objeto" (un CSV no tiene una marca de tiempo propia más allá de los datos que contiene). Cómo corregirlo: piensa en sla_hours siempre como "cuánto tiempo, como máximo, puede pasar entre la venta más reciente que conocemos y el momento en que decidimos revisar el archivo" — la definición exacta que usa esta lección, y la misma que usa dbt source freshness en producción.

Ejercicios

Ejercicio 1 — Corre check_freshness() con un sla_hours mucho más laxo, y confirma que S04 pasaría. Usando el df real de S04, llama a check_freshness(df, run_at=PIPELINE_RUN_AT, sla_hours=72) (un SLA de tres días en vez de uno). ¿Pasa o falla?

Ver solución
result_laxo = check_freshness(df, run_at=PIPELINE_RUN_AT, sla_hours=72)
print(result_laxo)

Salida esperada:

{'check': 'freshness', 'latest_row_ts': '2026-08-14 09:25:00', 'run_at': '2026-08-16T09:00:00', 'sla_hours': 72, 'hours_since_latest': 47.58, 'status': 'PASS'}

PASShours_since_latest no cambia (47.58, siempre la misma diferencia real entre latest_row_ts y run_at), pero el umbral contra el que se compara sí cambia, de 24 a 72. Este ejercicio confirma algo importante: el SLA no es un hecho técnico fijo de los datos, es una decisión de negocio —exactamente lo que ya advirtió el Ejercicio 1 de la lección 5 del módulo 1—. Los mismos datos, la misma función, dos veredictos distintos, según qué tan estricto decida ser Kiosko.

Ejercicio 2 — Confirma que el orden de las filas en el DataFrame no afecta el resultado. Vuelve a cargar orders_s04 con SELECT * FROM orders_s04 ORDER BY order_id (orden alfabético por order_id, distinto del orden original del CSV) y corre check_freshness() de nuevo. ¿Cambia algún valor del resultado?

Ver solución
df_reordenado = con.sql("SELECT * FROM orders_s04 ORDER BY order_id").pl()
print(check_freshness(df_reordenado, run_at=PIPELINE_RUN_AT, sla_hours=24))

Salida esperada: exactamente el mismo resultado que la corrida original —latest_row_ts: 2026-08-14 09:25:00, hours_since_latest: 47.58, status: FAIL—. pl.col(timestamp_col).max() es una operación de agregación: recorre todas las filas del DataFrame sin importar en qué orden estén físicamente almacenadas, y siempre encuentra el mismo valor máximo. Este comportamiento es una propiedad general de las funciones de agregación (max, min, sum, mean) en cualquier motor de DataFrame o SQL — nunca dependen del orden de llegada de las filas, a diferencia de operaciones como head() o first(), que sí lo hacen.

Ejercicio 3 — Argumenta si check_freshness() seguiría funcionando correctamente si orders_s04 tuviera alguna fila con order_ts nulo. En 2-3 frases, considerando cómo se comporta pl.col(timestamp_col).max() frente a valores nulos en Polars, explica qué pasaría si, por ejemplo, ORD-9503 tuviera también un order_ts vacío (no solo el unit_price vacío que ya tiene).

Ver solución

Por defecto, las funciones de agregación de Polars —incluida .max()— ignoran los valores nulos al calcular su resultado, así que un order_ts nulo en una fila no impediría que check_freshness() encontrara correctamente el máximo entre las filas que sí tienen fecha. El resultado seguiría siendo técnicamente correcto sobre las filas con dato, pero valdría la pena preguntarse si eso es lo que Kiosko realmente quiere: una fila sin fecha de venta es, en sí misma, un problema de completeness (ya cubierto por el módulo 2 de esta guía) que probablemente debería atraparse antes de que el archivo llegue a check_freshness(), no ser ignorada silenciosamente por una función que solo mide freshness. Este es un buen ejemplo de por qué los checks de esta guía están pensados para combinarse en un pipeline completo (el proyecto del módulo 8), no para usarse aislados uno del otro.

Resumen y siguiente paso

En esta lección escribiste check_freshness(), la síntesis directa de los dos argumentos de las lecciones 2 y 3: una función que mide el archivo completo (usando el order_ts máximo, nunca fila por fila) contra una referencia de tiempo fija (run_at, nunca el reloj real). La probaste primero en datos de juguete, confirmando que produce tanto PASS como FAIL correctamente en casos donde ya sabías la respuesta de antemano, y la corriste, por fin, sobre orders_2026-08-14.csv real: latest_row_ts: 2026-08-14 09:25:00, hours_since_latest: 47.58, status: FAIL — la sexta y última de las seis dimensiones de calidad de datos de esta guía, con evidencia ejecutada.

Antes de avanzar deberías poder: explicar, sin ver el código, por qué check_freshness() usa .max() sobre la columna de tiempo; reproducir el resultado FAIL de esta lección corriendo el código tú mismo; y explicar por qué su número (47.58 horas) es distinto, y no obstante también correcto, del 33.0 que calculó la lección 5 del módulo 1.

Freshness ya tiene su check. Queda una pieza más antes de cerrar el diagnóstico de nivel de archivo: la lección 5 construye check_volume() — y, a diferencia de esta lección, el resultado sobre S04 va a ser PASS, el contraste deliberado de que no todo en este archivo está roto.

Recursos

  • Polars — documentación oficial, expresiones de agregación (max, min, y su comportamiento frente a valores nulos, relevante para el Ejercicio 3 de esta lección). docs.pola.rs. En inglés.
  • dbt Labs — documentación oficial, "Add freshness checks to sources" (la definición de freshness verificada de forma independiente que confirma el diseño de check_freshness(): comparar el timestamp más reciente disponible contra el momento de la revisión). docs.getdbt.com/reference/resource-properties/freshness. En inglés.
  • Módulo 1, lección 5, de esta misma guía — fuente de freshness_preview.py y del número 33.0 horas, contrastado en esta lección contra 47.58. src/guides/data-reliability-and-governance-guide/workbook/module-01-when-green-does-not-mean-correct/es/05-meet-s04-kioskos-fourth-store.md. En español.
  • Módulo 2, lección 4, de esta misma guía — fuente exacta del puente CREATE OR REPLACE TABLE orders_s04 + .pl() reutilizado en esta lección. src/guides/data-reliability-and-governance-guide/workbook/module-02-declarative-data-quality-tests-with-pandera/es/04-installing-pandera-and-bridging-duckdb-to-polars.md. En español.
  • DISEÑO de esta guía — el mandato exacto de check_freshness(df, run_at=PIPELINE_RUN_AT, sla_hours=24). src/guides/data-reliability-and-governance-guide/DISENO.md. En español.