Módulo 8: Project Kioskos Batch Pipeline

Construyendo la capa bronze

Descripción

Esta lección construye la primera capa real del pipeline: bronze, la zona de aterrizaje donde los datos crudos de Kiosko llegan tal cual, organizados por fecha, sin que nadie decida todavía si son correctos o no. Vas a escribir ingest_bronze(source_folder, bronze_root, week_days), una función que reutiliza extract_orders() del módulo 2 —sin cambiar una sola línea de su lógica— y agrega la pieza que ningún módulo anterior necesitó: escribir cada día como su propia partición en disco, siguiendo el patrón data/bronze/orders/dt=YYYY-MM-DD/orders.csv, y hacerlo de forma que correr la función dos veces produzca exactamente el mismo resultado.

Conexión con el módulo. El brief de la lección 2 exigía un pipeline "seguro de re-correr" — esa exigencia empieza a cumplirse aquí, en la primera capa, con el mismo patrón overwrite-partition que ya trabajaste en el módulo 6: cada partición se sobrescribe completa cada vez, nunca se le agrega nada encima. Bronze es la base sobre la que la lección 4 construye silver — si esta capa tiene un error, se propaga a todo lo que sigue, así que vale la pena construirla con cuidado antes de avanzar.

Una analogía: la bodega de recepción, organizada por fecha de llegada

Piensa en la bodega de recepción de un supermercado grande, no una tienda pequeña. Cuando llega un camión con mercancía, nadie la mezcla con la de camiones anteriores en un montón indistinguible — cada entrega se registra en su propio espacio, etiquetado con la fecha exacta de llegada: "recepción del 3 de agosto", "recepción del 4 de agosto", cada una en su propio pasillo. Si alguien necesita revisar qué llegó un día específico, va directo a ese pasillo, sin tener que revisar toda la bodega. Y si el camión del 5 de agosto llega tarde y hay que volver a descargarlo por algún error de conteo, se vuelve a colocar exactamente en el pasillo del 5 de agosto — reemplazando lo que había, no amontonándolo encima.

ingest_bronze() es esa bodega. Cada partición dt=YYYY-MM-DD es un pasillo con la fecha exacta de esa entrega. Nada se mezcla entre fechas, nada se transforma todavía —la mercancía sigue en su empaque original, tal como llegó—, y si alguna vez hace falta volver a "descargar" el mismo día (correr el pipeline dos veces por accidente, o a propósito porque el archivo fuente cambió), el pasillo se reemplaza completo, nunca se le apila mercancía encima.

Ejemplo trabajado: ingest_bronze(), particionado y verificado dos veces

# bronze.py
import csv
from pathlib import Path

from kiosko import extract_orders

REQUIRED_COLUMNS = ["order_id", "store_id", "product_id", "quantity", "unit_price", "order_ts"]


def write_bronze_partition(rows_for_day: list[dict], bronze_root: str, partition_date: str) -> Path:
    """Write one day of raw rows to data/bronze/orders/dt=YYYY-MM-DD/orders.csv.

    Always overwrites the partition file completely (mode='w') -- this is the
    overwrite-partition pattern: re-running this function for the same
    partition_date produces the exact same file, never a bigger one.
    """
    partition_dir = Path(bronze_root) / "orders" / f"dt={partition_date}"
    partition_dir.mkdir(parents=True, exist_ok=True)
    target = partition_dir / "orders.csv"
    with open(target, "w", newline="") as f:
        writer = csv.DictWriter(f, fieldnames=REQUIRED_COLUMNS)
        writer.writeheader()
        writer.writerows(rows_for_day)
    return target


def ingest_bronze(source_folder: str, bronze_root: str, week_days: list[str]) -> dict[str, int]:
    """Extract each day of the week from the source folder and land it as its own bronze partition."""
    counts: dict[str, int] = {}
    for day in week_days:
        day_rows = extract_orders(source_folder, day, day)
        write_bronze_partition(day_rows, bronze_root, day)
        counts[day] = len(day_rows)
    return counts

Y ahora, corrida sobre la semana completa, dos veces seguidas:

if __name__ == "__main__":
    WEEK_DAYS = [
        "2026-08-03", "2026-08-04", "2026-08-05", "2026-08-06",
        "2026-08-07", "2026-08-08", "2026-08-09",
    ]

    print("=== Kiosko: building the bronze layer ===\n")
    counts = ingest_bronze(".", "data/bronze", WEEK_DAYS)
    for day in WEEK_DAYS:
        print(f"data/bronze/orders/dt={day}/orders.csv: {counts[day]:2} rows")
    print(f"\nTotal bronze rows: {sum(counts.values())}")

    print("\n=== Running ingest_bronze() a second time (same source, same days) ===")
    counts_again = ingest_bronze(".", "data/bronze", WEEK_DAYS)
    assert counts == counts_again, "bronze counts changed between two identical runs"
    print("Verificacion: counts primera corrida == counts segunda corrida -> OK")
    print(f"Total bronze rows (segunda corrida): {sum(counts_again.values())}")

Qué esperar. Al correr python3 bronze.py en la carpeta con los siete archivos de la semana, la salida es exactamente esta:

=== Kiosko: building the bronze layer ===

data/bronze/orders/dt=2026-08-03/orders.csv:  8 rows
data/bronze/orders/dt=2026-08-04/orders.csv:  6 rows
data/bronze/orders/dt=2026-08-05/orders.csv:  2 rows
data/bronze/orders/dt=2026-08-06/orders.csv:  5 rows
data/bronze/orders/dt=2026-08-07/orders.csv:  7 rows
data/bronze/orders/dt=2026-08-08/orders.csv:  9 rows
data/bronze/orders/dt=2026-08-09/orders.csv:  3 rows

Total bronze rows: 40

=== Running ingest_bronze() a second time (same source, same days) ===
Verificacion: counts primera corrida == counts segunda corrida -> OK
Total bronze rows (segunda corrida): 40

Cuarenta filas, en siete particiones, exactamente los mismos conteos por día que ya viste en los módulos 1 y 2 (ocho el lunes, tres el domingo). La segunda corrida es el momento que importa de verdad: ingest_bronze() se llamó otra vez, sobre los mismos siete archivos fuente, y el resultado —counts_again— es idéntico byte a byte al de la primera corrida. Si write_bronze_partition() hubiera usado modo "a" (append) en vez de "w" (overwrite) al abrir cada archivo, la segunda corrida habría duplicado cada partición — dieciséis filas el lunes en vez de ocho, ochenta filas en total en vez de cuarenta —, y el assert lo habría detectado de inmediato.

Diagrama: el árbol de particiones bronze resultante

data/bronze/orders/
├── dt=2026-08-03/orders.csv    (8 filas)
├── dt=2026-08-04/orders.csv    (6 filas)
├── dt=2026-08-05/orders.csv    (2 filas)
├── dt=2026-08-06/orders.csv    (5 filas)
├── dt=2026-08-07/orders.csv    (7 filas)
├── dt=2026-08-08/orders.csv    (9 filas)
└── dt=2026-08-09/orders.csv    (3 filas)
                                 ─────────
                                 40 filas totales
flowchart TD
    A["orders_2026-08-03.csv .. orders_2026-08-09.csv\n(7 archivos fuente, modulos 1-2)"] --> B["ingest_bronze(source, bronze_root, week_days)"]
    B --> C["extract_orders(source, day, day)\npor cada dia (modulo 2)"]
    C --> D["write_bronze_partition()\nmode='w', overwrite completo"]
    D --> E["data/bronze/orders/dt=YYYY-MM-DD/orders.csv\n7 particiones, 40 filas"]

Profundización: por qué cada día es su propio archivo, no un CSV gigante

Fíjate en una decisión de diseño concreta: write_bronze_partition() escribe un archivo orders.csv por cada día, en vez de acumular las cuarenta filas en un único archivo data/bronze/orders.csv. Esa decisión no es arbitraria — es lo que te permite reprocesar un solo día sin tocar los demás. Si el martes 2026-08-04 llegara un archivo corregido (por ejemplo, porque el punto de venta de esa tienda tuvo un error y volvió a exportar), ingest_bronze() podría correr solo sobre ese día —extract_orders(source_folder, "2026-08-04", "2026-08-04")— y sobrescribir únicamente dt=2026-08-04/orders.csv, sin tocar las otras seis particiones. Con un solo archivo gigante, cualquier corrección obligaría a reescribir la semana completa, incluso los seis días que no cambiaron.

Este es, en el fondo, el mismo argumento que ya viste en el módulo 6 con el patrón overwrite-partition, y en el módulo 7 con el particionamiento por fecha — este módulo no lo repite por repetir, lo aplica ahora a la primera capa de un pipeline de tres, donde el beneficio de granularidad se multiplica: bronze, silver y gold, cada uno particionado, significa que puedes reprocesar un solo día a través de las tres capas sin tocar el resto de la semana en ninguna de ellas.

Errores comunes

Usar open(target, "a") en vez de "w" al escribir la partición. Qué pasa: alguien, pensando en "no perder nada", abre el archivo de la partición en modo append ("a") en vez de write ("w"), asumiendo que así es más seguro. Por qué pasa: la palabra "append" suena, intuitivamente, más cuidadosa que "sobrescribir" — nadie quiere "perder" datos. Cómo detectarlo: si corres ingest_bronze() dos veces seguidas y el conteo de filas por partición crece en vez de mantenerse igual, tienes exactamente este error. Cómo corregirlo: en el patrón overwrite-partition, "seguro" significa lo contrario de lo que la intuición sugiere — cada corrida debe reemplazar completamente la partición, no sumarle nada. open(target, "w"), como hace el ejemplo de esta lección, es la forma correcta.

Olvidar partition_dir.mkdir(parents=True, exist_ok=True) y fallar en la primera corrida. Qué pasa: alguien escribe write_bronze_partition() sin crear primero el directorio dt=YYYY-MM-DD/, y el open() falla con FileNotFoundError porque esa carpeta todavía no existe. Por qué pasa: es fácil olvidar que, a diferencia de un archivo, un directorio no se crea automáticamente al intentar escribir dentro de él. Cómo detectarlo: si tu error es exactamente FileNotFoundError: [Errno 2] No such file or directory: 'data/bronze/orders/dt=2026-08-03/orders.csv', te falta esta línea. Cómo corregirlo: Path(...).mkdir(parents=True, exist_ok=True) antes de cualquier open() en modo escritura — parents=True crea también los directorios padre que falten (data/, data/bronze/, data/bronze/orders/), y exist_ok=True evita un error si la carpeta ya existe de una corrida anterior.

Verificar la idempotencia comparando solo el conteo total, no por partición. Qué pasa: alguien confirma que sum(counts.values()) es igual en las dos corridas (40 == 40), y da por probada la idempotencia, sin comparar partición por partición. Por qué pasa: el total es el número más visible, y "40 == 40" se siente como suficiente evidencia. Cómo detectarlo: un bug que perdiera 2 filas del lunes y ganara 2 filas de más el martes produciría el mismo total (40), aunque el resultado por día esté mal en dos particiones distintas. Cómo corregirlo: el assert counts == counts_again de esta lección compara los diccionarios completos, día por día, no solo la suma — así un error localizado en una sola partición no puede esconderse detrás de un total que "cuadra por casualidad".

Ejercicios

Ejercicio 1 — Reprocesa un solo día. Sin volver a correr ingest_bronze() sobre toda la semana, escribe el código que reprocese únicamente la partición del sábado (2026-08-08), y confirma que las otras seis particiones no cambiaron (puedes verificarlo comparando Path("data/bronze/orders/dt=2026-08-08/orders.csv").stat().st_mtime antes y después, o simplemente confiando en que write_bronze_partition() solo tocó esa ruta).

Ver solución
from bronze import write_bronze_partition
from kiosko import extract_orders

saturday_rows = extract_orders(".", "2026-08-08", "2026-08-08")
write_bronze_partition(saturday_rows, "data/bronze", "2026-08-08")
print(f"Particion del 2026-08-08 reprocesada: {len(saturday_rows)} filas")

Salida esperada:

Particion del 2026-08-08 reprocesada: 9 filas

Solo se llamó a write_bronze_partition() con partition_date="2026-08-08" — la función nunca toca ninguna otra ruta que la de esa partición específica, así que las otras seis quedan exactamente como estaban. Esa es, en la práctica, la ventaja concreta de particionar: reprocesar un día no implica tocar los demás.

Ejercicio 2 — Rompe la idempotencia a propósito, para entenderla. Modifica temporalmente write_bronze_partition() para que abra el archivo con open(target, "a") en vez de "w" (sin escribir de nuevo el encabezado, para simplificar). Corre ingest_bronze() dos veces seguidas sobre el mismo día y observa qué pasa con el conteo de filas.

Ver solución
def write_bronze_partition_broken(rows_for_day, bronze_root, partition_date):
    partition_dir = Path(bronze_root) / "orders" / f"dt={partition_date}"
    partition_dir.mkdir(parents=True, exist_ok=True)
    target = partition_dir / "orders.csv"
    with open(target, "a", newline="") as f:  # bug: append en vez de write
        writer = csv.DictWriter(f, fieldnames=REQUIRED_COLUMNS)
        writer.writerows(rows_for_day)  # tambien le falta writeheader(), a proposito simplificado

Al correr ingest_bronze() dos veces seguidas con esta versión rota, cada partición terminaría con el doble de filas de las que debería tener — el lunes pasaría de 8 a 16, la semana completa de 40 a 80. El assert counts == counts_again del ejemplo trabajado de esta lección seguiría funcionando bien (porque counts se calcula a partir de extract_orders() sobre el archivo fuente original, no leyendo la partición bronze ya corrompida) — pero si inspeccionaras directamente el archivo data/bronze/orders/dt=2026-08-03/orders.csv con wc -l, verías el problema de inmediato. Deshaz el cambio antes de seguir.

Ejercicio 3 — Argumenta el siguiente paso. Bronze ya tiene las 40 filas de la semana, particionadas y aterrizadas de forma segura. Pero si alguien intentara sumar revenue directamente sobre estas particiones bronze en este momento, ¿qué encontraría que se lo impide? Usando lo que ya sabes del módulo 3 (ELT) y del módulo 4 (modelado), responde en 2-3 frases sin escribir código todavía.

Ver solución

Bronze no tiene ninguna columna revenue — solo tiene las seis columnas originales de orders, exactamente como llegaron, sin ningún cálculo aplicado, igual que load_raw() en el módulo 3. Además, todos los valores siguen siendo texto (quantity y unit_price como cadenas, no como número o decimal), así que ni siquiera se podría multiplicar directamente sin antes convertir los tipos. Calcular revenue requiere el mismo paso que ya construiste en el módulo 4: transform_fact_orders(), que hace el join contra las dimensiones y calcula la columna derivada — exactamente lo que la lección 4 de este módulo (silver) construye a continuación.

Resumen y siguiente paso

En esta lección construiste la capa bronze del pipeline de Kiosko: ingest_bronze() reutilizó extract_orders() del módulo 2 sin cambiar su lógica, y agregó el particionamiento por fecha con el patrón overwrite-partition (data/bronze/orders/dt=YYYY-MM-DD/orders.csv). Confirmaste, corriendo la función dos veces seguidas, que el resultado es idéntico ambas veces —cuarenta filas, en siete particiones, sin duplicar nada—, la primera de las tres exigencias del brief de la lección 2 ya cumplida.

Antes de avanzar deberías poder: explicar por qué cada día es su propia partición, no un archivo gigante; identificar el error concreto de usar "a" en vez de "w" al escribir una partición; y reprocesar un solo día sin tocar los demás.

La lección 4 toma estas particiones bronze —crudas, sin validar— y construye la capa silver: la compuerta de calidad del módulo 5 separa lo confiable de lo roto, y transform_fact_orders() del módulo 4 convierte lo confiable en fact_orders, con revenue calculado por fin.

Recursos