Módulo 8: Project Kioskos Lakehouse

Reproduciendo la historia de P002 con time travel

Descripción

Esta es la lección central de todo el capstone. Sobre el mismo catálogo que la lección 3 dejó listo —con kiosko.fact_orders, kiosko.dim_store y kiosko.dim_date ya cargadas—, esta lección crea kiosko.dim_product sin ninguna columna de historia, reproduce el cambio real de P002 (snacks/0.60health-snacks/0.68, vigencia 2026-08-15) con table.overwrite(), y une las cuatro tablas del star completo para calcular el margen por categoría dos veces: una contra el estado vigente (roto, 9.36), otra contra table.scan(snapshot_id=snap_v1) (correcto, 10.8). Son los mismos dos números exactos que ya confirmaron data-modeling-for-analytics-guide (módulo 8) y dbt-analytics-engineering-guide (módulo 8), cada una con su propio motor — aquí, sin que kiosko.dim_product tenga jamás una columna valid_from.

Conexión con el módulo. Esta lección responde, de punta a punta, la segunda exigencia del brief de la lección 2: "recordar el pasado sin valid_from". Retoma, aplicado ahora dentro del lakehouse completo, exactamente el mecanismo que el módulo 3 de esta guía ya enseñó paso a paso —append(), capturar snap_v1, overwrite(), scan(snapshot_id=...)—.

Una analogía: el mismo balance del contador, ahora con todos los libros sobre la misma mesa

El módulo 3 de esta guía calculó el margen correcto de P002 con un solo libro contable abierto: dim_product, comparado consigo mismo en dos momentos distintos. Esta lección hace el mismo cálculo, pero con los cuatro libros de Kiosko abiertos sobre la misma mesa a la vezfact_orders, dim_store, dim_date, dim_product—, exactamente como haría un contador real que necesita cruzar ventas, tiendas, fechas y catálogo de productos para cerrar un balance completo, no solo verificar el precio de un producto en aislamiento. El resultado no cambia —sigue siendo 10.8—, pero la forma de obtenerlo ahora involucra el star completo, no una tabla sola.

Ejemplo trabajado: dim_product sin historia + el star unido con time travel

Paso 1 — El esquema de dim_product, cuatro columnas, ninguna de historia

# kiosko_dim_product_time_travel.py -- modulo 8, leccion 4
import os
from collections import defaultdict

import pyarrow as pa
from pyiceberg.catalog import load_catalog
from pyiceberg.schema import Schema
from pyiceberg.types import DoubleType, NestedField, StringType

DIM_PRODUCT_SCHEMA = Schema(
    NestedField(field_id=1, name="product_id", field_type=StringType(), required=True),
    NestedField(field_id=2, name="product_name", field_type=StringType(), required=True),
    NestedField(field_id=3, name="category", field_type=StringType(), required=True),
    NestedField(field_id=4, name="unit_cost", field_type=DoubleType(), required=True),
)
PA_SCHEMA = pa.schema([
    pa.field("product_id", pa.string(), nullable=False),
    pa.field("product_name", pa.string(), nullable=False),
    pa.field("category", pa.string(), nullable=False),
    pa.field("unit_cost", pa.float64(), nullable=False),
])
DIM_PRODUCT_V1 = [
    {"product_id": "P001", "product_name": "Bottled Water 600ml", "category": "beverages", "unit_cost": 0.40},
    {"product_id": "P002", "product_name": "Energy Bar", "category": "snacks", "unit_cost": 0.60},
    {"product_id": "P003", "product_name": "Instant Coffee Sachet", "category": "beverages", "unit_cost": 0.35},
    {"product_id": "P004", "product_name": "Phone Charger Cable", "category": "electronics", "unit_cost": 2.10},
]
DIM_PRODUCT_V2 = [
    {"product_id": "P001", "product_name": "Bottled Water 600ml", "category": "beverages", "unit_cost": 0.40},
    {"product_id": "P002", "product_name": "Energy Bar", "category": "health-snacks", "unit_cost": 0.68},
    {"product_id": "P003", "product_name": "Instant Coffee Sachet", "category": "beverages", "unit_cost": 0.35},
    {"product_id": "P004", "product_name": "Phone Charger Cable", "category": "electronics", "unit_cost": 2.10},
]

Cuatro columnas — product_id, product_name, category, unit_cost — exactamente como quedó kiosko.dim_product al cerrar el módulo 3. Ni valid_from, ni valid_to, ni is_current, ni dbt_scd_id.

Paso 2 — La función de margen, que une dim_product con fact_orders

def margin_by_category(dim_rows: list, fact_rows: list) -> tuple:
    dim_by_id = {r["product_id"]: r for r in dim_rows}
    revenue, margin = defaultdict(float), defaultdict(float)
    for f in fact_rows:
        d = dim_by_id[f["product_id"]]
        revenue[d["category"]] += f["revenue"]
        margin[d["category"]] += f["revenue"] - f["quantity"] * d["unit_cost"]
    return revenue, margin

Esta función recibe cualquier versión de dim_product —el estado vigente, o el resultado de un scan(snapshot_id=...)— y la une, en memoria, contra las filas de fact_orders. La única diferencia entre "roto" y "correcto" en esta lección es qué lista de filas de dim_product se le pasa a esta misma función — el mismo patrón exacto que ya usó el proyecto de cierre del módulo 3.

Paso 3 — El flujo completo: carga, cambio, time travel, margen

def main() -> None:
    print("=== Kiosko: dim_product sin historia + el star unido con time travel ===\n")

    warehouse_path = os.path.abspath("kiosko_warehouse")
    catalog_db_path = os.path.abspath("kiosko_catalog.db")
    catalog = load_catalog(
        "kiosko", type="sql",
        uri=f"sqlite:///{catalog_db_path}", warehouse=f"file://{warehouse_path}",
    )
    fact_orders = catalog.load_table("kiosko.fact_orders")
    dim_store = catalog.load_table("kiosko.dim_store")
    print(f"Paso 1/6 -- reutilizando kiosko.fact_orders ({fact_orders.scan().to_arrow().num_rows} filas) "
          f"y kiosko.dim_store ({dim_store.scan().to_arrow().num_rows} filas) de la leccion 3")

    dim_product = catalog.create_table("kiosko.dim_product", schema=DIM_PRODUCT_SCHEMA)
    print("Paso 2/6 -- kiosko.dim_product creada, columnas:",
          [f.name for f in dim_product.schema().fields], "(cero columnas de historia)")

    dim_product.append(pa.Table.from_pylist(DIM_PRODUCT_V1, schema=PA_SCHEMA))
    snap_v1 = dim_product.current_snapshot().snapshot_id
    print("Paso 3/6 -- V1 cargada con table.append(), snap_v1 capturado (P002=snacks/0.60)")

    dim_product.overwrite(pa.Table.from_pylist(DIM_PRODUCT_V2, schema=PA_SCHEMA))
    print("Paso 4/6 -- table.overwrite() aplica el cambio real de P002 (health-snacks/0.68, vigencia 2026-08-15)")

    fact_rows = fact_orders.scan().to_arrow().to_pylist()
    current_rows = dim_product.scan().to_arrow().to_pylist()
    v1_rows = dim_product.scan(snapshot_id=snap_v1).to_arrow().to_pylist()
    current_p002 = next(r for r in current_rows if r["product_id"] == "P002")
    v1_p002 = next(r for r in v1_rows if r["product_id"] == "P002")
    print(f"Paso 5/6 -- estado vigente P002: {current_p002['category']}/{current_p002['unit_cost']} | "
          f"AS OF snap_v1 P002: {v1_p002['category']}/{v1_p002['unit_cost']}")

    revenue_broken, margin_broken = margin_by_category(current_rows, fact_rows)
    revenue_correct, margin_correct = margin_by_category(v1_rows, fact_rows)
    total_revenue = round(sum(f["revenue"] for f in fact_rows), 2)
    print("Paso 6/6 -- margen por categoria calculado contra fact_orders + dim_store + dim_date (el star completo)\n")

    print("=== Verificacion final ===\n")
    print("ROTO (JOIN contra el estado vigente, category='health-snacks'):")
    for cat in sorted(revenue_broken):
        print(f"  {cat:14} revenue={round(revenue_broken[cat], 2):>6}  margin={round(margin_broken[cat], 2):>6}")
    print("\nCORRECTO (JOIN contra table.scan(snapshot_id=snap_v1), category='snacks'):")
    for cat in sorted(revenue_correct):
        print(f"  {cat:14} revenue={round(revenue_correct[cat], 2):>6}  margin={round(margin_correct[cat], 2):>6}")
    print(f"\nRevenue total del star (identico en ambos casos, viene de fact_orders): {total_revenue}")

    assert dim_product.scan().to_arrow().num_rows == 4
    assert len(dim_product.schema().fields) == 4, "dim_product no debe tener ninguna columna de historia"
    assert current_p002["category"] == "health-snacks" and current_p002["unit_cost"] == 0.68
    assert v1_p002["category"] == "snacks" and v1_p002["unit_cost"] == 0.60
    assert round(margin_broken["health-snacks"], 2) == 9.36
    assert round(margin_correct["snacks"], 2) == 10.8
    assert total_revenue == 106.15

    print("\nTodas las verificaciones pasaron: dim_product con 4 columnas (cero de historia), "
          "margin roto=9.36, margin correcto=10.8 via time travel, revenue total=106.15 -- "
          "los mismos numeros que data-modeling-for-analytics-guide M8 y dbt-analytics-engineering-guide M8.")


if __name__ == "__main__":
    main()

Qué esperar (verificado corriendo python3 kiosko_dim_product_time_travel.py real, en el mismo directorio que la lección 3, sin borrar kiosko_warehouse/; ningún snapshot_id se imprime como literal — el módulo 3 explica por qué):

=== Kiosko: dim_product sin historia + el star unido con time travel ===

Paso 1/6 -- reutilizando kiosko.fact_orders (40 filas) y kiosko.dim_store (3 filas) de la leccion 3
Paso 2/6 -- kiosko.dim_product creada, columnas: ['product_id', 'product_name', 'category', 'unit_cost'] (cero columnas de historia)
Paso 3/6 -- V1 cargada con table.append(), snap_v1 capturado (P002=snacks/0.60)
Paso 4/6 -- table.overwrite() aplica el cambio real de P002 (health-snacks/0.68, vigencia 2026-08-15)
Paso 5/6 -- estado vigente P002: health-snacks/0.68 | AS OF snap_v1 P002: snacks/0.6
Paso 6/6 -- margen por categoria calculado contra fact_orders + dim_store + dim_date (el star completo)

=== Verificacion final ===

ROTO (JOIN contra el estado vigente, category='health-snacks'):
  beverages      revenue= 44.05  margin= 14.75
  electronics    revenue=  40.5  margin=  21.6
  health-snacks  revenue=  21.6  margin=  9.36

CORRECTO (JOIN contra table.scan(snapshot_id=snap_v1), category='snacks'):
  beverages      revenue= 44.05  margin= 14.75
  electronics    revenue=  40.5  margin=  21.6
  snacks         revenue=  21.6  margin=  10.8

Revenue total del star (identico en ambos casos, viene de fact_orders): 106.15

Todas las verificaciones pasaron: dim_product con 4 columnas (cero de historia), margin roto=9.36, margin correcto=10.8 via time travel, revenue total=106.15 -- los mismos numeros que data-modeling-for-analytics-guide M8 y dbt-analytics-engineering-guide M8.

Fíjate en algo que no cambia entre "roto" y "correcto": beverages y electronics tienen exactamente el mismo revenue y el mismo margen en ambos bloques. Eso es correcto y esperado — ningún producto de esas dos categorías cambió nunca; solo P002 se movió, de snacks a health-snacks. La diferencia entre los dos bloques está, exclusivamente, en qué fila de dim_product describe a P002 en el momento de calcular el margen.

Diagrama: el star completo, con dim_product en dos tiempos

flowchart LR
    FO["kiosko.fact_orders\n40 filas, order_ts fijo"]
    DS["kiosko.dim_store\ncon country"]
    DD["kiosko.dim_date\n31 filas"]
    DP1["kiosko.dim_product\nAS OF snap_v1\nP002=snacks/0.60"]
    DP2["kiosko.dim_product\nvigente\nP002=health-snacks/0.68"]

    FO --> STAR["El star unido\n(en memoria, por product_id/store_id)"]
    DS --> STAR
    DD --> STAR
    DP1 -->|"margin correcto"| M1["snacks: margin=10.8"]
    DP2 -->|"margin roto"| M2["health-snacks: margin=9.36"]
    STAR --> DP1
    STAR --> DP2

Profundización: por qué esta lección hace, sin BETWEEN, exactamente lo que hizo el JOIN punto-en-el-tiempo

data-modeling-for-analytics-guide resolvió este mismo problema —calcular el margen correcto de P002 sobre las 40 órdenes reales— con un JOIN punto-en-el-tiempo: fact_orders.order_ts BETWEEN dim_product_scd.valid_from AND dim_product_scd.valid_to. Ese patrón funciona porque cada fila de fact_orders compara su propia fecha contra el rango de vigencia de la versión correcta de dim_product_scd. dbt-analytics-engineering-guide automatizó la misma idea con dbt_valid_from/dbt_valid_to, generados por dbt snapshot.

Esta lección llega al mismo resultado con una estrategia distinta, posible únicamente porque las 40 órdenes de Kiosko ocurren todas antes del cambio de P002 (entre el 03 y el 09 de agosto; el cambio tiene vigencia el 15). En vez de comparar la fecha de cada fila contra un rango de vigencia por fila, esta lección hace una sola pregunta, una vez, sobre toda la tabla: "¿cómo se veía dim_product completo, en el instante justo antes del cambio?" — y usa esa única respuesta (table.scan(snapshot_id=snap_v1)) para las 40 filas a la vez. Es exactamente el límite honesto que la lección 7 del módulo 3 ya documentó con evidencia: time travel responde "cómo se veía TODA la tabla en el instante X", no "qué era cierto para ESTA fila en SU propia fecha" — y en el caso específico de Kiosko, con todos los hechos del mismo lado del cambio, esas dos preguntas tienen la misma respuesta. Si alguna orden de Kiosko hubiera ocurrido después del 15 de agosto, esta técnica dejaría de ser suficiente, y haría falta volver al JOIN punto-en-el-tiempo de data-modeling-for-analytics-guide.

Errores comunes

Usar dim_product.scan() (sin snapshot_id) para calcular el margen "correcto", por error. Qué pasa: alguien, al escribir la llamada a margin_by_category(), pasa current_rows en el lugar donde debería ir v1_rows, y obtiene 9.36 donde esperaba 10.8. Por qué pasa: las dos variables tienen nombres parecidos (current_rows/v1_rows), y es fácil confundir cuál representa "ahora" y cuál "antes del cambio". Cómo detectarlo: si tu bloque "CORRECTO" muestra health-snacks en vez de snacks, invertiste las dos variables. Cómo corregirlo: recuerda la regla mnemotécnica de esta guía: v1_rows viene de scan(snapshot_id=snap_v1) —el pasado, capturado explícitamente—; current_rows viene de scan() sin argumentos —siempre el presente—. El margen correcto de negocio, para las órdenes de Kiosko, siempre usa el pasado (v1_rows), porque todas las órdenes ocurrieron antes del cambio.

Pensar que beverages y electronics deberían cambiar también entre "roto" y "correcto". Qué pasa: alguien, al ver dos bloques de resultados distintos, espera que los seis números (dos categorías, revenue y margin) cambien entre ambos. Por qué pasa: es fácil asumir que "dos versiones del cálculo" implica "todos los números son distintos". Cómo detectarlo: si tu implementación produce valores distintos para beverages en los dos bloques, hay un bug — ningún producto de esa categoría cambió nunca en la historia de Kiosko. Cómo corregirlo: revisa margin_by_category() — solo la categoría de P002 debería diferir entre revenue_broken/margin_broken y revenue_correct/margin_correct; las otras tres categorías (beverages, electronics, y la ausencia de una fila snacks separada en el bloque roto) son la evidencia correcta de que el cambio fue puntual, no generalizado.

Ejercicios

Ejercicio 1 — Corre el script tú mismo, en el mismo directorio que la lección 3. Sin borrar kiosko_warehouse/, corre python3 kiosko_dim_product_time_travel.py. Confirma que ves los seis pasos completarse y el mensaje final con los siete valores verificados.

Ver solución

Si corriste la lección 3 primero, en el mismo directorio, la salida debería reproducir exactamente la estructura de esta lección: seis pasos numerados —los dos primeros confirmando que fact_orders y dim_store ya existían—, seguidos de la verificación final con margin roto=9.36, margin correcto=10.8, y revenue total=106.15. Si ves TableDoesNotExistError en el Paso 1, no corriste la lección 3 en este mismo directorio primero.

Ejercicio 2 — Rompe un assert a propósito, y observa el fallo. Cambia temporalmente DIM_PRODUCT_V2 para que P002 tenga unit_cost=0.70 en vez de 0.68, corre el script de nuevo, y observa qué assert falla primero. Después revierte el cambio (y borra kiosko_warehouse/kiosko/dim_product/ si necesitas repetir la carga desde cero, o simplemente vuelve a correr todo el módulo desde la lección 3 en un directorio nuevo).

Ver solución

El primer assert en fallar es assert current_p002["category"] == "health-snacks" and current_p002["unit_cost"] == 0.68 — porque cambiaste el valor de V2, el estado vigente ya no coincide con lo que el script espera. Este ejercicio, igual que en el proyecto de cierre del módulo 3, demuestra que los assert de esta lección están encadenados a los valores canónicos exactos de Kiosko, no a un resultado genérico.

Ejercicio 3 — Explica, en tus propias palabras, por qué el revenue total (106.15) es idéntico en el bloque "roto" y el bloque "correcto", pero el margen no. En 2-3 frases, justifica por qué el cambio de P002 afecta el margen pero no el revenue.

Ver solución

El revenue de cada línea de orden viene enteramente de fact_ordersquantity × unit_price—, una tabla que nunca cambió en ningún momento de esta guía; por eso el revenue total, 106.15, es idéntico sin importar qué versión de dim_product se use para calcularlo. El margen, en cambio, depende de unit_cost, una columna que solo existe en dim_product — y esa es, precisamente, la columna que cambió con el overwrite() de V1 a V2. Esta distinción es la misma que ya explicó data-modeling-for-analytics-guide: los hechos (fact_orders) registran lo que pasó y no cambian; las dimensiones (dim_product) describen el contexto, y sí pueden cambiar — el margen, al depender de ambas tablas a la vez, hereda la inestabilidad de la dimensión, mientras el revenue, que depende solo del hecho, se queda fijo.

Resumen y siguiente paso

En esta lección creaste kiosko.dim_product sin ninguna columna de historia, reprodujiste el cambio real de P002 con table.overwrite(), y uniste las cuatro tablas del star completo para calcular el margen por categoría dos veces: 9.36 (roto, contra el estado vigente) y 10.8 (correcto, contra table.scan(snapshot_id=snap_v1)) — los mismos dos números exactos que ya confirmaron data-modeling-for-analytics-guide y dbt-analytics-engineering-guide, cada una con su propio motor y su propia técnica.

Antes de avanzar deberías poder: explicar por qué el revenue total no cambia entre el cálculo roto y el correcto, pero el margen sí; y explicar por qué esta técnica —time travel puro, sin JOIN punto-en-el-tiempo— es suficiente para el caso específico de Kiosko, pero no lo sería si alguna orden ocurriera después del 15 de agosto.

La lección 5 deja dim_product como está y se mueve a la quinta tabla del lakehouse: kiosko.fact_orders_at_scale, particionada por store_id y evolucionada con DayTransform, con diez millones de filas.

Recursos

  • PyIceberg — documentación oficial (quickstart), el flujo de append(), overwrite() y scan(snapshot_id=...) que integra esta lección. py.iceberg.apache.org. En inglés.
  • PyIceberg — referencia de API, table.current_snapshot(), table.scan(snapshot_id=...). py.iceberg.apache.org/api. En inglés.
  • Esta misma guía, módulo 3, lección 7 — fuente del límite honesto de time travel que esta lección aprovecha (todos los hechos del mismo lado del cambio). ../module-03-snapshots-and-time-travel/es/07-what-time-travel-does-not-replace.md. En español.
  • DISEÑO de data-modeling-for-analytics-guide — fuente de los números canónicos 9.36/10.8 que esta lección verifica con assert, ahora dentro del lakehouse completo. src/guides/data-modeling-for-analytics-guide/DISENO.md. En español.
  • DISEÑO de esta guía — el mapa completo de los ocho módulos, incluida la lección 5 que sigue. src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.