Módulo 7: Partitioning And Orchestration
Encadenando los pasos del pipeline
Descripción
Tienes, repartidas en varios módulos distintos, las piezas que Kiosko necesita para procesar un día de ventas: extract_orders() (módulo 2), validate_orders() (módulo 5), parse_order() (módulo 1) para convertir cada fila válida en un Order tipado, transform_fact_orders() (módulo 4), y write_partition() (la de la lección anterior, sobre el patrón del módulo 6). Esta lección las junta, por primera vez, en una sola función: run_pipeline(date). Es, literalmente, el corazón de este módulo — la "orquestación mínima" que anunció la lección 1.
Conexión con el módulo. Fíjate en el orden exacto: extract → validate → parse → transform → load. No es el orden en que construiste las piezas a lo largo de la guía —el módulo 4 (transformar) vino antes que el módulo 5 (validar)—, sino el orden en que deberían ejecutarse en un pipeline real: nunca vale la pena transformar una fila que ya sabes que está rota. Validar antes de transformar evita gastar trabajo —el join contra DIM_STORE y DIM_PRODUCT— sobre datos que de todas formas vas a descartar. Y el parseo va entre los dos: validate_orders() trabaja sobre filas crudas (dict, texto sin convertir, tal como las definió el módulo 5), pero transform_fact_orders() espera objetos Order ya tipados (la firma exacta que declaró el módulo 4) — por eso cada row de valid_rows pasa por parse_order() antes de llegar a la transformación, nunca directo.
Una analogía: la receta, ahora escrita en código
En la lección 1 comparaste run_pipeline() con una receta de cocina: una lista de pasos, en orden, que alguien sigue sin necesitar coordinación externa. Esta lección escribe esa receta, literal, como una función de Python. Una receta bien escrita no repite la explicación completa de "cómo picar una cebolla" cada vez que la necesita — dice "pica la cebolla" y confía en que ya sabes hacerlo, o remite a una técnica ya explicada antes en el recetario. run_pipeline() hace lo mismo: no vuelve a escribir cómo se lee un CSV, ni cómo se valida una fila, ni cómo se hace el join con las dimensiones — llama a las funciones que ya resuelven cada una de esas técnicas, en el orden correcto, y confía en que cada una hace bien su parte.
Ejemplo trabajado: run_pipeline(), la primera versión
# kiosko.py (agrega esto al archivo que vienes construyendo desde el modulo 1)
from dataclasses import dataclass
@dataclass
class PipelineResult:
partition_date: str
status: str
rows_extracted: int
rows_valid: int
rows_rejected: int
rows_loaded: int
partition_path: str
def run_pipeline(partition_date: str, folder: str = ".") -> PipelineResult:
"""Encadena extract -> validate -> parse -> transform -> load para una sola fecha."""
raw_rows = extract_orders(folder, partition_date, partition_date)
valid_rows, rejected_rows = validate_orders(raw_rows)
orders = [parse_order(row) for row in valid_rows]
fact_rows = transform_fact_orders(orders, DIM_STORE, DIM_PRODUCT)
partition_path = write_partition(fact_rows, partition_date)
return PipelineResult(
partition_date=partition_date,
status="ok",
rows_extracted=len(raw_rows),
rows_valid=len(valid_rows),
rows_rejected=len(rejected_rows),
rows_loaded=len(fact_rows),
partition_path=partition_path,
)
Y así se ve correrla, sobre un solo día conocido:
# chain_monday.py
from kiosko import run_pipeline
result = run_pipeline("2026-08-03")
print(result)
Qué esperar. Al correr python3 chain_monday.py en una carpeta con los siete archivos de la semana, la salida es exactamente esta:
PipelineResult(partition_date='2026-08-03', status='ok', rows_extracted=8, rows_valid=8, rows_rejected=0, rows_loaded=8, partition_path='data/fact_orders/dt=2026-08-03')
Una sola línea, pero contiene todo lo que necesitas saber sobre esa corrida: la fecha que procesó, que terminó bien (status='ok'), cuántas filas entraron en cada paso —8 extraídas, las 8 pasaron la validación, las 8 se cargaron—, y dónde quedó el resultado. Cinco funciones, escritas en cuatro módulos distintos de esta guía —extract_orders(), validate_orders(), parse_order(), transform_fact_orders(), write_partition()—, corrieron en el orden correcto con una sola llamada.
Diagrama: la cadena completa
flowchart LR
A["extract_orders(folder, date, date)"] --> B["validate_orders(raw_rows)"]
B -->|"valid_rows"| P["parse_order() por fila"]
P --> C["transform_fact_orders(orders, DIM_STORE, DIM_PRODUCT)"]
B -->|"rejected_rows"| E["(se cuentan, no se cargan)"]
C --> D["write_partition(fact_rows, date)"]
D --> F["PipelineResult"]
Profundización: por qué devolver un objeto, no solo imprimir
Fíjate en que run_pipeline() no imprime nada — construye y devuelve un PipelineResult. Es tentador, mientras se prueba una función nueva, meter un par de print() adentro para ver qué está pasando. Pero una función que imprime en vez de devolver un resultado estructurado tiene un problema real: quien la llama no puede usar esa información, solo verla pasar por la terminal. La lección 8 de este módulo va a llamar a run_pipeline() siete veces seguidas, una por cada día de la semana, y va a necesitar comparar los resultados entre sí —¿cuántas filas cargó cada día?, ¿algún día falló?—. Eso solo es posible si cada llamada devuelve algo que el código que sigue puede leer, no algo que solo un humano leyendo la terminal puede interpretar.
@dataclass —la misma herramienta que usaste para Order en el módulo 1— es la elección natural para PipelineResult: te da campos con nombre (result.rows_loaded, no result[3]), un __repr__ legible de gratis (la línea completa que viste en "Qué esperar"), y la garantía de que todo resultado de una corrida tiene exactamente la misma forma, sin importar qué día procesó. Compáralo con devolver un dict suelto: funcionaría, pero result["rows_loadde"] —un typo en la llave— fallaría en silencio con un KeyError solo cuando alguien intentara leer ese campo, mientras que result.rows_loadde con una dataclass falla de inmediato, en el momento exacto del typo, con un error claro.
Errores comunes
Meter lógica de negocio dentro de run_pipeline() en vez de llamar a las funciones existentes. Qué pasa: alguien, al escribir run_pipeline(), reimplementa parte de la validación o la transformación directamente adentro de la función, en vez de llamar a validate_orders() o transform_fact_orders(). Por qué pasa: mientras se escribe una función nueva, es tentador resolver todo ahí mismo en vez de acordarse de las piezas que ya existen en otros módulos. Cómo detectarlo: si run_pipeline() tiene más de cuatro o cinco líneas de lógica propia —más allá de encadenar llamadas y armar el PipelineResult—, probablemente estás reimplementando algo que ya construiste antes. Cómo corregirlo: run_pipeline() no debería saber cómo se valida una fila o cómo se calcula revenue — solo debería saber en qué orden llamar a las funciones que sí lo saben. Es el mismo principio de responsabilidad única que ya viste en el módulo 2, al mantener separados extract_orders() y parse_order().
Hardcodear la fecha dentro de la función en vez de recibirla como parámetro. Qué pasa: alguien escribe una primera versión de run_pipeline() con "2026-08-03" escrito directamente en el cuerpo de la función, en vez de usar el parámetro partition_date. Por qué pasa: mientras se prueba con un solo día conocido, es fácil escribir el valor literal y "arreglarlo después". Cómo detectarlo: si llamar a run_pipeline("2026-08-04") sigue procesando el 2026-08-03, la fecha está hardcodeada en vez de usarse de verdad. Cómo corregirlo: cada uso de una fecha dentro de la función debe venir del parámetro partition_date, nunca de un valor escrito a mano — es exactamente la misma disciplina que ya aplicaste con extract_orders(folder, start_date, end_date) en el módulo 2.
Olvidar que validate_orders() devuelve una tupla, no una lista. Qué pasa: alguien escribe valid_rows = validate_orders(raw_rows) con una sola variable, en vez de valid_rows, rejected_rows = validate_orders(raw_rows), y valid_rows termina siendo la tupla completa (valid, rejected) en vez de solo la lista de filas válidas. Por qué pasa: es fácil olvidar la firma exacta de una función construida en un módulo anterior. Cómo detectarlo: si transform_fact_orders(valid_rows, ...) lanza un error extraño, o si len(valid_rows) te da 2 en vez del número de filas que esperabas, revisa si valid_rows en verdad contiene una tupla de dos listas en vez de una sola lista. Cómo corregirlo: siempre que llames a validate_orders(), desempaca las dos partes con nombre — valid_rows, rejected_rows = validate_orders(rows) —, tal como hace el ejemplo trabajado de esta lección.
Pasar valid_rows sin parsear directo a transform_fact_orders(). Qué pasa: alguien escribe transform_fact_orders(valid_rows, DIM_STORE, DIM_PRODUCT), saltándose parse_order(), porque valid_rows "ya pasó la compuerta de calidad" y se siente listo para usarse. Por qué pasa: validate_orders() (módulo 5) devuelve filas crudas —dict, texto sin convertir— que simplemente resultaron ser válidas; "válido" no es lo mismo que "parseado", y es fácil confundir las dos cosas. Cómo detectarlo: si tu pipeline revienta con AttributeError: 'dict' object has no attribute 'store_id' en el paso de transformación, es exactamente esta causa — transform_fact_orders(rows: list[Order], ...) accede a order.store_id por atributo, algo que un dict no tiene. Cómo corregirlo: convierte siempre valid_rows con orders = [parse_order(row) for row in valid_rows] antes de pasarlo a transform_fact_orders() — es el mismo parse_order() del módulo 1, aplicado aquí como el puente obligatorio entre la compuerta de calidad (que trabaja con dict crudos) y la transformación (que trabaja con Order tipados).
Ejercicios
Ejercicio 1 — Corre la cadena sobre otro día. Usando run_pipeline() del ejemplo trabajado, córrela sobre el sábado (2026-08-08, el día más ocupado de la semana). Antes de correr el código, predice los valores de rows_extracted y rows_loaded usando lo que ya sabes del módulo 2.
Ver solución
from kiosko import run_pipeline
result = run_pipeline("2026-08-08")
print(result)
Salida esperada:
PipelineResult(partition_date='2026-08-08', status='ok', rows_extracted=9, rows_valid=9, rows_rejected=0, rows_loaded=9, partition_path='data/fact_orders/dt=2026-08-08')
Nueve filas en los cuatro campos —extraídas, válidas, cargadas—, coincidiendo con el conteo del sábado que ya conocías del módulo 2: el día con más órdenes de toda la semana de Kiosko.
Ejercicio 2 — Accede a un solo campo del resultado. Sin imprimir el objeto PipelineResult completo, escribe código que llame a run_pipeline("2026-08-05") (miércoles) e imprima solo una frase con el número de filas cargadas, usando el atributo rows_loaded directamente.
Ver solución
from kiosko import run_pipeline
result = run_pipeline("2026-08-05")
print(f"El {result.partition_date} se cargaron {result.rows_loaded} filas.")
Salida esperada:
El 2026-08-05 se cargaron 2 filas.
Este es, precisamente, el punto de usar una dataclass en vez de imprimir texto suelto: puedes leer cualquier campo del resultado por separado —result.rows_loaded, result.status, result.partition_path—, sin tener que analizar (parsear) una cadena de texto para sacar el número que te interesa.
Ejercicio 3 — Argumenta por qué el orden de los cuatro pasos importa. run_pipeline() llama a las funciones en el orden extract → validate → transform → load. Explica en 2-3 frases qué problema aparecería si el orden fuera, en cambio, extract → transform → validate → load — es decir, transformando antes de validar.
Ver solución
Si transform_fact_orders() corriera antes que validate_orders(), el pipeline haría el trabajo de unir cada fila con su tienda y su producto —y calcular revenue— incluso sobre filas que después van a ser descartadas por estar rotas, desperdiciando ese trabajo. Peor aún: si una fila rota tuviera, por ejemplo, un store_id vacío o mal formado, transform_fact_orders() podría fallar al intentar el join contra DIM_STORE antes de que validate_orders() tuviera la oportunidad de detectar y quitar esa fila del camino de forma controlada — convirtiendo lo que debería ser una fila quarentenada, con un reporte claro, en un error inesperado que interrumpe el pipeline. Validar primero es, en el fondo, la misma idea que "cruzar la aduana antes de cocinar" del módulo 1: nunca proceses algo que todavía no confirmaste que es válido.
Resumen y siguiente paso
En esta lección construiste run_pipeline(date), la función que encadena extract_orders(), validate_orders(), parse_order(), transform_fact_orders() y write_partition() en el orden correcto, y devuelve un PipelineResult —una dataclass, no texto impreso— con todo lo que necesitas saber sobre esa corrida. La corriste sobre el lunes y el sábado de Kiosko, y confirmaste que los módulos anteriores de esta guía ahora trabajan juntos con una sola llamada.
Antes de avanzar deberías poder: explicar el orden extract → validate → transform → load y por qué es ese y no otro; llamar a run_pipeline() sobre cualquier día de la semana de Kiosko; y explicar por qué devolver un objeto estructurado es mejor que imprimir texto dentro de la función.
run_pipeline() funciona, pero es silenciosa por dentro: si la corres sin guardar el resultado, o si la programas para correr sola cada noche sin que nadie mire la terminal en vivo, no queda ningún rastro de qué pasó paso a paso. La lección 5 resuelve exactamente eso con logging, la herramienta de la librería estándar diseñada para pipelines que corren sin supervisión.
Recursos
- Python — documentación oficial de
dataclasses, reutilizada aquí paraPipelineResult. docs.python.org/3/library/dataclasses.html. En inglés. - Joe Reis & Matt Housley, Fundamentals of Data Engineering (O'Reilly, 2022) — el capítulo de transformación describe por qué el orden de las etapas de un pipeline no es arbitrario. oreilly.com/library/view/fundamentals-of-data/9781098108298. En inglés.
- Maxime Beauchemin, Functional Data Engineering — describe pipelines como funciones que encadenan pasos idempotentes, el mismo espíritu de
run_pipeline(). maximebeauchemin.medium.com/functional-data-engineering. En inglés.