Módulo 3: Structured Logging For Pipelines

Logging de una corrida del pipeline, de principio a fin

Descripción

Esta es la lección central del módulo: junta todo lo que construiste en las lecciones 3, 4 y 5 —niveles, JSON, contexto vinculado— en un solo archivo permanente, logging_config.py, y lo conecta de verdad a pipeline.py y a __init__.py. Al final de esta lección, kiosko_pipeline no va a tener ni un solo print() — cada paso de cada una de las siete particiones de la semana, desde la extracción hasta la carga, va a quedar registrado como un evento JSON, con run_id y partition_date vinculados automáticamente. Y vas a confirmar algo que no es obvio de entrada: que la salida de una corrida real, con logging estructurado, puede ser tan reproducible byte por byte como cualquier test de pytest — pero solo si prestas atención a un detalle concreto que esta lección te muestra fallando primero, en vivo.

Conexión con el módulo. Esta lección resuelve, con código ejecutado de punta a punta, el problema completo que abrió el módulo: qué pasó exactamente en la última corrida de kiosko_pipeline, sin depender de que alguien estuviera mirando la terminal en el momento.

Una analogía: el itinerario completo de un vuelo con escalas, no solo el reporte de llegada

Un vuelo con varias escalas genera, para los sistemas de control aéreo, un registro por cada tramo: despegue de origen, llegada a la primera escala, despegue de la escala, llegada a destino — no solo un mensaje final de "vuelo completado". Si algo sale mal en el segundo tramo, ese registro por tramo permite reconstruir exactamente hasta dónde llegó el vuelo, sin depender de un único reporte al final que, si el vuelo nunca llega, simplemente nunca se genera.

kiosko_pipeline, hasta el final del módulo 2, solo generaba el equivalente al reporte final de llegada —nueve líneas de print(), después de que las siete particiones ya habían terminado—. Esta lección construye el registro por tramo: un evento por cada paso —extracción, calidad, transformación, carga— de cada una de las siete particiones, para que, sin importar en qué punto de la semana se detenga una corrida real, quede un rastro completo de hasta dónde llegó.

Ejemplo trabajado, parte 1: logging_config.py completo

Crea un archivo nuevo, src/kiosko_pipeline/logging_config.py:

# src/kiosko_pipeline/logging_config.py
"""Configure structlog so every pipeline event becomes one JSON object per line on stdout.

Call configure_logging() once, before the pipeline runs. Every logger obtained
afterwards with get_pipeline_logger() renders through the same processor chain and
picks up whatever context -- run_id, partition_date -- gets bound with
structlog.contextvars, without any log call having to pass that context explicitly.
"""

from __future__ import annotations

import logging
import sys

import structlog


def build_deterministic_run_id(week_days: list[str]) -> str:
    """Derive a run_id from the data itself -- never from a clock or uuid4().

    Two runs over the same week_days produce the exact same run_id, on purpose.
    In production you would typically reach for uuid4() or an orchestrator-assigned
    run id (Airflow's run_id, for example) -- something that tells two runs of the
    SAME data apart. This guide needs the opposite: every "Que esperar" block has
    to be reproducible byte for byte, so the id is derived from the partition range
    instead of anything that changes between runs.
    """
    return f"kiosko-{week_days[0]}-{week_days[-1]}"


def _fixed_timestamper(fixed_timestamp: str):
    """A processor that stamps every event with the same literal timestamp."""

    def add_fixed_timestamp(logger, method_name, event_dict):
        event_dict["timestamp"] = fixed_timestamp
        return event_dict

    return add_fixed_timestamp


def configure_logging(level: int = logging.INFO, fixed_timestamp: str | None = None) -> None:
    """Wire structlog on top of the stdlib logging module: JSON out, one event per line.

    fixed_timestamp is a teaching knob, not a production one: pass a literal string
    to pin every event of a run to the same timestamp, which is what makes the
    "Que esperar" blocks in this guide reproducible no matter when you run them.
    Leave it as None for a real run -- structlog then falls back to its own
    TimeStamper, which reads the actual system clock.
    """
    timestamp_processor = (
        _fixed_timestamper(fixed_timestamp)
        if fixed_timestamp is not None
        else structlog.processors.TimeStamper(fmt="iso", utc=True)
    )

    logging.basicConfig(format="%(message)s", stream=sys.stdout, level=level)

    structlog.configure(
        processors=[
            structlog.contextvars.merge_contextvars,
            structlog.processors.add_log_level,
            timestamp_processor,
            structlog.processors.StackInfoRenderer(),
            structlog.processors.format_exc_info,
            structlog.processors.JSONRenderer(sort_keys=True),
        ],
        wrapper_class=structlog.make_filtering_bound_logger(level),
        logger_factory=structlog.stdlib.LoggerFactory(),
        cache_logger_on_first_use=True,
    )


def get_pipeline_logger(name: str = "kiosko_pipeline") -> structlog.stdlib.BoundLogger:
    """A thin wrapper so every module in the package asks for its logger the same way."""
    return structlog.get_logger(name)

Tres piezas que ya conoces, ensambladas. build_deterministic_run_id() retoma la lección 5: en vez de un uuid4() o un timestamp del reloj —lo que usarías en un pipeline real, para distinguir corridas entre sí—, esta guía deriva el run_id de los datos mismos (week_days[0] y week_days[-1]), porque cada "Qué esperar" de este módulo tiene que reproducirse exactamente igual sin importar cuándo lo corras. _fixed_timestamper() hace lo mismo con el timestamp: en una corrida real, dejarías fixed_timestamp=None y structlog.processors.TimeStamper(fmt="iso", utc=True) —la propia herramienta de structlog, que lee el reloj del sistema— haría el trabajo; para esta guía, un timestamp literal fijo hace que la salida sea idéntica hoy, en un año, y en cualquier máquina. Y configure_logging() arma la cadena completa de processors —contexto, nivel, timestamp, JSON— en un solo lugar, llamado una sola vez, exactamente como recomendó la profundización de la lección 3.

Ejemplo trabajado, parte 2: instrumentando pipeline.py

Ahora, conecta ese archivo a pipeline.py. La lógica de negocio —qué se extrae, qué se valida, cómo se calcula revenue— no cambia ni una línea; lo único que se agrega son los eventos de log, en los puntos exactos donde algo termina:

# src/kiosko_pipeline/pipeline.py
"""Chain extract -> validate -> transform -> load for Kiosko's full week, over one connection."""

import sqlite3
from dataclasses import dataclass
from pathlib import Path

import structlog

from kiosko_pipeline.extract import extract_orders, parse_order
from kiosko_pipeline.load import load_to_warehouse
from kiosko_pipeline.logging_config import get_pipeline_logger
from kiosko_pipeline.quality import validate_orders
from kiosko_pipeline.transform import DIM_PRODUCT, DIM_STORE, transform_fact_orders

WEEK_DAYS = [
    "2026-08-03", "2026-08-04", "2026-08-05", "2026-08-06",
    "2026-08-07", "2026-08-08", "2026-08-09",
]

log = get_pipeline_logger(__name__)


@dataclass
class PipelineResult:
    rows_extracted: int
    rows_valid: int
    rows_rejected: int
    rows_loaded: int
    week_revenue: float
    revenue_by_store: dict[str, float]
    revenue_by_product: dict[str, float]


def run_pipeline(
    folder: str, week_days: list[str], db_path: str = "data/warehouse.db", run_id: str = "unbound-run"
) -> PipelineResult:
    """Extract -> validate -> parse -> transform -> load, one connection, the whole week."""
    structlog.contextvars.bind_contextvars(run_id=run_id)
    log.info("pipeline_run_started", week_days=week_days, db_path=db_path)

    Path(db_path).parent.mkdir(parents=True, exist_ok=True)
    con = sqlite3.connect(db_path)

    total_extracted = total_valid = total_rejected = total_loaded = 0
    for day in week_days:
        structlog.contextvars.bind_contextvars(partition_date=day)

        raw_rows = extract_orders(folder, day, day)
        log.info("extract_completed", row_count=len(raw_rows))

        valid_rows, rejected_rows = validate_orders(raw_rows)
        log.info(
            "quality_gate_completed",
            valid_count=len(valid_rows),
            rejected_count=len(rejected_rows),
        )

        orders = [parse_order(row) for row in valid_rows]
        fact_rows = transform_fact_orders(orders, DIM_STORE, DIM_PRODUCT)
        partition_revenue = round(sum(row["revenue"] for row in fact_rows), 2)
        log.info("transform_completed", row_count=len(fact_rows), partition_revenue=partition_revenue)

        loaded = load_to_warehouse(con, fact_rows, day)
        log.info("load_completed", rows_loaded=loaded)

        total_extracted += len(raw_rows)
        total_valid += len(valid_rows)
        total_rejected += len(rejected_rows)
        total_loaded += loaded

        structlog.contextvars.unbind_contextvars("partition_date")

    week_revenue = con.execute("SELECT ROUND(SUM(revenue), 2) FROM fact_orders").fetchone()[0]
    revenue_by_store = dict(con.execute(
        "SELECT store_id, ROUND(SUM(revenue), 2) FROM fact_orders GROUP BY store_id ORDER BY store_id"
    ).fetchall())
    revenue_by_product = dict(con.execute(
        "SELECT product_id, ROUND(SUM(revenue), 2) FROM fact_orders GROUP BY product_id ORDER BY product_id"
    ).fetchall())
    con.close()

    log.info(
        "pipeline_run_completed",
        rows_extracted=total_extracted,
        rows_valid=total_valid,
        rows_rejected=total_rejected,
        rows_loaded=total_loaded,
        week_revenue=week_revenue,
    )
    structlog.contextvars.unbind_contextvars("run_id")

    return PipelineResult(
        rows_extracted=total_extracted, rows_valid=total_valid,
        rows_rejected=total_rejected, rows_loaded=total_loaded,
        week_revenue=week_revenue,
        revenue_by_store=revenue_by_store, revenue_by_product=revenue_by_product,
    )

Compáralo contra la versión del módulo 2: la firma de run_pipeline() gana un parámetro nuevo, run_id: str = "unbound-run" (con un valor por defecto, para que el código que ya lo llamaba sin ese argumento —incluida toda la suite de tests del módulo 2— siga funcionando sin cambios); el cuerpo del bucle gana un bind_contextvars(partition_date=day) al principio y un unbind_contextvars("partition_date") al final de cada iteración; y cuatro llamadas a log.info(...) marcan el final de cada paso —extracción, calidad, transformación, carga— con el conteo que corresponde a ese paso específico. Ni una fórmula, ni un SELECT, ni el orden de ninguna operación de negocio cambia.

Ejemplo trabajado, parte 3: main() sin ningún print()

Por último, reescribe __init__.py para que main() configure el logging y reporte el resultado final como un evento más, no como texto formateado para humanos:

# src/kiosko_pipeline/__init__.py
"""kiosko_pipeline: Kiosko's batch pipeline, packaged with uv."""

from kiosko_pipeline.logging_config import build_deterministic_run_id, configure_logging, get_pipeline_logger
from kiosko_pipeline.pipeline import WEEK_DAYS, run_pipeline
from kiosko_pipeline.transform import DIM_PRODUCT, DIM_STORE

STORE_NAMES = {s["store_id"]: s["store_name"] for s in DIM_STORE}
PRODUCT_NAMES = {p["product_id"]: p["product_name"] for p in DIM_PRODUCT}


def main() -> None:
    configure_logging(fixed_timestamp="2026-08-12T09:00:00Z")
    log = get_pipeline_logger(__name__)

    run_id = build_deterministic_run_id(WEEK_DAYS)
    result = run_pipeline("data", WEEK_DAYS, run_id=run_id)

    log.info(
        "kiosko_weekly_report",
        run_id=run_id,
        rows_extracted=result.rows_extracted,
        rows_valid=result.rows_valid,
        rows_rejected=result.rows_rejected,
        rows_loaded=result.rows_loaded,
        week_revenue=result.week_revenue,
        revenue_by_store=result.revenue_by_store,
        revenue_by_product=result.revenue_by_product,
    )

Nota una decisión de diseño, a propósito: revenue_by_store y revenue_by_product se loguean con las claves originales (S01, P001), no con los nombres legibles de STORE_NAMES/PRODUCT_NAMES que sí usaba el print() del módulo 2 ("Kiosko Centro", "Bottled Water 600ml"). No es un descuido — es la prioridad correcta para un log estructurado: store_id es el identificador estable que cualquier consumidor posterior —un dashboard, una alerta, otra parte del pipeline— va a usar para filtrar o agrupar; el nombre legible es información de presentación, útil para un humano leyendo directamente, pero no es lo que un programa debería tener que parsear para encontrar la fila que le importa. STORE_NAMES y PRODUCT_NAMES siguen definidos en el archivo, disponibles si en el futuro alguna parte del pipeline los necesita —el CLI del módulo 7, por ejemplo, podría usarlos para una salida legible en la terminal, aparte del log—.

Ejemplo trabajado, parte 4: la corrida completa

Corre el pipeline, exactamente como siempre:

uv run python -m kiosko_pipeline

Qué esperar (treinta y una líneas, una por evento; se muestran las primeras cinco, un salto, y las dos últimas — la lección 8 incluye la corrida completa):

{"db_path": "data/warehouse.db", "event": "pipeline_run_started", "level": "info", "run_id": "kiosko-2026-08-03-2026-08-09", "timestamp": "2026-08-12T09:00:00Z", "week_days": ["2026-08-03", "2026-08-04", "2026-08-05", "2026-08-06", "2026-08-07", "2026-08-08", "2026-08-09"]}
{"event": "extract_completed", "level": "info", "partition_date": "2026-08-03", "row_count": 8, "run_id": "kiosko-2026-08-03-2026-08-09", "timestamp": "2026-08-12T09:00:00Z"}
{"event": "quality_gate_completed", "level": "info", "partition_date": "2026-08-03", "rejected_count": 0, "run_id": "kiosko-2026-08-03-2026-08-09", "timestamp": "2026-08-12T09:00:00Z", "valid_count": 8}
{"event": "transform_completed", "level": "info", "partition_date": "2026-08-03", "partition_revenue": 15.85, "row_count": 8, "run_id": "kiosko-2026-08-03-2026-08-09", "timestamp": "2026-08-12T09:00:00Z"}
{"event": "load_completed", "level": "info", "partition_date": "2026-08-03", "rows_loaded": 8, "run_id": "kiosko-2026-08-03-2026-08-09", "timestamp": "2026-08-12T09:00:00Z"}
...
{"event": "pipeline_run_completed", "level": "info", "rows_extracted": 40, "rows_loaded": 40, "rows_rejected": 0, "rows_valid": 40, "run_id": "kiosko-2026-08-03-2026-08-09", "timestamp": "2026-08-12T09:00:00Z", "week_revenue": 106.15}
{"event": "kiosko_weekly_report", "level": "info", "revenue_by_product": {"P001": 33.55, "P002": 21.6, "P003": 10.5, "P004": 40.5}, "revenue_by_store": {"S01": 38.3, "S02": 38.8, "S03": 29.05}, "rows_extracted": 40, "rows_loaded": 40, "rows_rejected": 0, "rows_valid": 40, "run_id": "kiosko-2026-08-03-2026-08-09", "timestamp": "2026-08-12T09:00:00Z", "week_revenue": 106.15}

Confirma los números de siempre en la última línea: rows_extracted=40, week_revenue=106.15, y el mismo desglose por tienda y producto que ya conoces del módulo 1 — instrumentar el pipeline con logging no cambió ni un decimal del resultado, exactamente la misma promesa que ya viste cumplirse en los dos módulos anteriores.

Un detalle que casi rompe la reproducibilidad: el orden de las claves

Antes de cerrar la lección, vale la pena mostrar un problema real que aparece la primera vez que alguien arma esta cadena de processors, porque tiene una solución de una sola palabra que es fácil pasar por alto. Quita, temporalmente, el argumento sort_keys=True de structlog.processors.JSONRenderer() en logging_config.py:

structlog.processors.JSONRenderer(),  # sin sort_keys, temporalmente

Corre el pipeline dos veces seguidas, guardando cada salida en un archivo:

uv run python -m kiosko_pipeline > /tmp/run1.jsonl 2>/dev/null
uv run python -m kiosko_pipeline > /tmp/run2.jsonl 2>/dev/null
diff /tmp/run1.jsonl /tmp/run2.jsonl

Qué esperar (con el bug, sin sort_keys; se muestra la primera de las veintiocho líneas que difieren, marcadas con 2,29c2,29 porque el bloque completo va de la línea 2 a la 29):

2,29c2,29
< {"row_count": 8, "event": "extract_completed", "run_id": "kiosko-2026-08-03-2026-08-09", "partition_date": "2026-08-03", "level": "info", "timestamp": "2026-08-12T09:00:00Z"}
< {"valid_count": 8, "rejected_count": 0, "event": "quality_gate_completed", "run_id": "kiosko-2026-08-03-2026-08-09", "partition_date": "2026-08-03", "level": "info", "timestamp": "2026-08-12T09:00:00Z"}
...
---
> {"row_count": 8, "event": "extract_completed", "partition_date": "2026-08-03", "run_id": "kiosko-2026-08-03-2026-08-09", "level": "info", "timestamp": "2026-08-12T09:00:00Z"}
> {"valid_count": 8, "rejected_count": 0, "event": "quality_gate_completed", "partition_date": "2026-08-03", "run_id": "kiosko-2026-08-03-2026-08-09", "level": "info", "timestamp": "2026-08-12T09:00:00Z"}
...

Lee esto con atención: los dos objetos JSON son semánticamente idénticos —mismas claves, mismos valores, row_count: 8, partition_date: "2026-08-03", todo igual—, pero el orden en el que aparecen las claves cambió entre una corrida y la otra: run_id antes de partition_date en la primera corrida, partition_date antes de run_id en la segunda. Esto pasa porque el contexto que vincula merge_contextvars viene de un diccionario interno que Python arma en un orden que depende del historial exacto de operaciones de esa ejecución específica del proceso —no es aleatorio en el sentido de random, pero tampoco es algo que este módulo pueda dejar sin controlar, porque rompe la promesa central de esta guía: que cada "Qué esperar" sea reproducible byte por byte.

La solución es la que ya tenías en el archivo original: structlog.processors.JSONRenderer(sort_keys=True) — el parámetro sort_keys, que JSONRenderer reenvía directamente a json.dumps(), ordena las claves alfabéticamente sin importar en qué orden llegaron al diccionario. Restaura esa línea, y repite la comparación:

uv run python -m kiosko_pipeline > /tmp/run1.jsonl 2>/dev/null
uv run python -m kiosko_pipeline > /tmp/run2.jsonl 2>/dev/null
diff /tmp/run1.jsonl /tmp/run2.jsonl && echo "IDENTICAL"

Qué esperar (con sort_keys=True):

IDENTICAL

Sin ninguna diferencia. Esta guía usa sort_keys=True precisamente por esto — y, como beneficio adicional, un JSON con las claves siempre en el mismo orden alfabético es más fácil de comparar visualmente entre dos líneas de log distintas, aunque no dependieras de la reproducibilidad exacta para nada más.

Diagrama: un evento por paso, por partición

flowchart TD
    START["pipeline_run_started\n(run_id vinculado)"] --> LOOP{"por cada uno de\nlos 7 dias"}
    LOOP --> BIND["bind partition_date"]
    BIND --> EXT["extract_completed\n(row_count)"]
    EXT --> QUAL["quality_gate_completed\n(valid_count, rejected_count)"]
    QUAL --> TRANS["transform_completed\n(row_count, partition_revenue)"]
    TRANS --> LOAD["load_completed\n(rows_loaded)"]
    LOAD --> UNBIND["unbind partition_date"]
    UNBIND --> LOOP
    LOOP -->|"7 dias completos"| DONE["pipeline_run_completed\n(totales de la semana)"]
    DONE --> REPORT["kiosko_weekly_report\n(desglose por tienda/producto)"]

Treinta y una líneas en total: un pipeline_run_started, cuatro eventos por cada uno de los siete días (28 en total), un pipeline_run_completed, y un kiosko_weekly_report final. Si el pipeline se hubiera roto en el día cuatro —como en el fallo de la lección 2—, tendrías, como mínimo, pipeline_run_started más los eventos completos de los primeros tres días (doce líneas) antes del corte — una reconstrucción precisa de cuánto se alcanzó a procesar, algo que la versión de print() nunca pudo darte.

Errores comunes

Reconfigurar el logging a mitad de una sesión de Python, y no ver el cambio reflejado en loggers que ya existían. Qué pasa: alguien llama configure_logging() una vez, obtiene un logger con get_pipeline_logger(), y más adelante en la misma sesión de Python llama configure_logging() de nuevo —por ejemplo, con un fixed_timestamp distinto— esperando que el logger que ya tenía cambie de comportamiento. No cambia: sigue usando la configuración de cuando se creó. Por qué pasa: cache_logger_on_first_use=True en structlog.configure() —una opción de rendimiento deliberada, no un descuido— hace que cada logger, la primera vez que se usa, guarde una copia de la cadena de processors vigente en ese momento; reconfigurar después no actualiza esa copia ya cacheada, solo afecta a loggers nuevos, obtenidos después del cambio. Cómo detectarlo: si reconfiguraste el logging y algunos eventos parecen "atrasados" respecto al cambio mientras otros lo reflejan correctamente, revisa si esos loggers específicos se obtuvieron antes o después de la reconfiguración. Cómo corregirlo: llama configure_logging() una sola vez, al principio absoluto del programa —exactamente el patrón de main() en esta lección—, y nunca a mitad de una ejecución. Si de verdad necesitas reconfigurar (algo poco común fuera de tests), obtén loggers nuevos después del cambio, no reutilices los que ya tenías de antes.

Dejar un print() residual sin migrar, y romper el parseo línea por línea de todo el log. Qué pasa: alguien migra casi todo kiosko_pipeline a logging estructurado, pero se olvida de una línea suelta —un print() de depuración que quedó de una sesión anterior, en cualquier archivo del paquete— y no la elimina. El resultado: una línea de texto libre, mezclada entre las líneas JSON, que rompe cualquier programa que intente hacer json.loads() de cada línea del log sin excepción. Por qué pasa: es fácil revisar __init__.py y pipeline.py —los dos archivos que esta lección modifica explícitamente— y no notar un print() olvidado en algún otro punto del código, sobre todo si se agregó temporalmente durante alguna sesión de depuración anterior y nunca se quitó. Cómo detectarlo: corre el pipeline, guarda la salida completa en un archivo, e intenta parsear cada línea, sin excepción, con json.loads() — cualquier línea que falle señala, con el número de línea exacto, dónde quedó el residuo de texto libre. Cómo corregirlo: después de terminar esta lección, busca en todo el paquete (grep -rn "print(" src/kiosko_pipeline/) cualquier print() que haya quedado — al cierre de este módulo, el resultado de esa búsqueda debería ser una lista vacía.

Olvidar pasar run_id explícitamente, y no notar "unbound-run" hasta mucho después. Qué pasa: alguien llama run_pipeline("data", WEEK_DAYS) directamente —sin pasar el argumento run_id, quizás desde un test o un script de prueba rápida— y todos los eventos de esa corrida quedan etiquetados con el valor por defecto, "unbound-run", sin que nada avise que faltó un argumento. Por qué pasa: run_id: str = "unbound-run" tiene un valor por defecto explícito, a propósito, para que el código del módulo 2 que llama run_pipeline() sin ese argumento —o cualquier test futuro— siga funcionando sin romperse; pero ese mismo valor por defecto puede esconder, en una corrida real, el hecho de que alguien olvidó pasar un run_id genuino. Cómo detectarlo: si ves "run_id": "unbound-run" en cualquier log que debería representar una corrida real de producción, es una señal clara de que ese run_id nunca se generó ni se pasó explícitamente — revisa el punto exacto donde se llamó run_pipeline(). Cómo corregirlo: en cualquier código que no sea un test aislado, pasa siempre un run_id explícito y significativo —como hace main() en esta lección, con build_deterministic_run_id(WEEK_DAYS)—; reserva el valor por defecto "unbound-run" para el caso genuino de llamar la función de forma aislada, donde el contexto de una corrida completa no aplica.

Ejercicios

Ejercicio 1 — Reproduce el bug del orden de claves tú mismo. Quita sort_keys=True de JSONRenderer(), corre el pipeline tres veces seguidas guardando cada salida en un archivo distinto, y compara los tres archivos entre sí con diff. Confirma que al menos dos de las tres corridas difieren en el orden de las claves de al menos una línea. Después, restaura sort_keys=True y confirma que las tres corridas vuelven a ser idénticas.

Ver solución

Sin sort_keys=True, el orden de las claves dentro de cada objeto JSON depende del orden interno en que merge_contextvars combina el contexto vinculado con los campos del evento — un orden que puede variar entre ejecuciones distintas del proceso de Python, aunque el contenido semántico (las claves y sus valores) sea idéntico. diff entre dos o tres corridas debería mostrar diferencias línea por línea, típicamente en la posición de run_id y partition_date dentro de cada objeto. Con sort_keys=True restaurado, las claves de cada objeto JSON quedan siempre en el mismo orden alfabético, sin importar el orden interno en que se construyó el diccionario — y las tres corridas deberían ser idénticas, confirmado por un diff sin ninguna salida.

Ejercicio 2 — Cuenta las líneas exactas de un fallo simulado. Reintroduce, temporalmente, el mismo dato malformado de la lección 2 (ORD-4003 con store_id="S09" en orders_2026-08-06.csv), y corre el pipeline instrumentado de esta lección. Cuenta cuántas líneas de log JSON alcanzas a ver antes del Traceback, y confirma que corresponden exactamente a los primeros tres días completos (doce eventos: pipeline_run_started más cuatro eventos por cada uno de los tres primeros días). Deshaz el cambio al terminar.

Ver solución

Con el dato malformado reintroducido, deberías ver trece líneas de log JSON antes del Traceback: un pipeline_run_started, y cuatro eventos completos (extract_completed, quality_gate_completed, transform_completed, load_completed) por cada uno de los tres primeros días (2026-08-03, 2026-08-04, 2026-08-05), doce eventos en total, más el inicial — trece. El cuarto día (2026-08-06) alcanza a generar extract_completed y quality_gate_completed (la fila con S09 pasa la compuerta de calidad sin problema, porque no es un error de tipo ni de rango), pero se detiene antes de transform_completed, exactamente donde transform_fact_orders() lanza el ValueError. Comparado con el Traceback desnudo de la lección 2 —que no daba ninguna de estas trece confirmaciones—, esta es, con precisión, la diferencia que construyó este módulo.

Ejercicio 3 — Explica la decisión de loguear IDs en vez de nombres legibles. Sin mirar la lección otra vez, explica en 2-3 frases por qué main() loguea revenue_by_store con claves como "S01" en vez de "Kiosko Centro", a pesar de que el print() del módulo 2 sí mostraba el nombre legible.

Ver solución

Un log estructurado está pensado, en primer lugar, para que un programa —no una persona leyendo directamente— lo consuma: un dashboard, una alerta, o cualquier sistema posterior que necesite filtrar o agrupar por tienda va a comparar contra un identificador estable como store_id, no contra un nombre legible que podría escribirse de formas ligeramente distintas ("Kiosko Centro" contra "kiosko centro" contra "Centro (Kiosko)") sin que eso represente ningún cambio real de negocio. STORE_NAMES y PRODUCT_NAMES siguen definidos en el paquete, disponibles para cualquier parte del código que sí necesite presentar información a un humano —una salida de terminal legible, un reporte—, pero el log en sí prioriza lo que un sistema puede usar sin ambigüedad.

Resumen y siguiente paso

En esta lección construiste logging_config.py, el archivo permanente que combina niveles (lección 3), formato JSON (lección 4) y contexto vinculado con structlog (lección 5), y lo conectaste de verdad a pipeline.py (un evento por paso, por partición) y a __init__.py (main() sin ningún print()). Corriste el pipeline completo y confirmaste, con diff, que la salida es reproducible byte por byte entre corridas — después de encontrar y corregir, en vivo, un problema real de orden de claves que casi rompe esa promesa.

Antes de avanzar deberías poder: explicar qué evento se genera en cada paso del pipeline, y con qué campos; explicar por qué sort_keys=True es necesario para la reproducibilidad de este módulo; y nombrar la diferencia entre un run_id genuino y el valor por defecto "unbound-run".

El pipeline ahora registra todo lo que hace — pero "todo" es una palabra peligrosa cuando se trata de datos. La lección 7 se detiene, antes de seguir construyendo, para responder una pregunta de seguridad que este módulo no puede saltarse: de todo lo que el pipeline sabe sobre cada orden, ¿qué es seguro loguear, y qué nunca debería aparecer en un archivo de texto plano?

Recursos