Módulo 3: Snapshots And Time Travel

Proyecto: el dim_product de Kiosko viajado en el tiempo

Descripción

Este proyecto cierra el módulo 3. Tienes kiosko.dim_product creada sin columnas de historia (lección 2), sabes provocar el cambio de P002 con table.overwrite() (lección 3), sabes por qué nunca hardcodeas un snapshot_id y cómo capturarlo o recuperarlo correctamente (lección 4), sabes viajar en el tiempo con table.scan(snapshot_id=...) (lección 5), sabes recuperar el margen correcto de P002 uniendo esa lectura histórica contra fact_orders (lección 6), y sabes, con evidencia ejecutada, dónde termina esta técnica (lección 7). Falta un solo paso: juntar las seis piezas en un solo script, corrido de punta a punta, con assert automáticos que confirman cada número.

Conexión con el módulo. Este proyecto no introduce ningún concepto nuevo — es la integración final de las siete lecciones anteriores. Retoma, de forma literal, la promesa que abrió este módulo en la lección 1: recuperar el mismo resultado correcto que ya confirmaron data-modeling-for-analytics-guide y dbt-analytics-engineering-guide, con una tabla que tiene cero columnas de historia.

Una analogía: el archivo completo, recorrido de una sola vez

Las lecciones 2 a 7 de este módulo construyeron, una pieza a la vez, el archivo completo de fotos de dim_product: el estante instalado y la primera foto tomada (lección 2), el reabastecimiento que cambió P002 (lección 3), la disciplina de anotar el número de rollo correcto (lección 4), la solicitud de la foto vieja al encargado del archivo (lección 5), el balance del contador hecho con la foto correcta (lección 6), y la advertencia honesta de cuándo ese archivo no basta (lección 7). Este proyecto es el momento de repetir todo el proceso, de punta a punta, en un solo gesto continuo — el mismo tipo de integración que ya hiciste al cerrar el módulo 1.

El material: todo lo que este módulo construyó, en un solo lugar

Necesitas, en un directorio de trabajo nuevo:

kiosko_time_travel/
├── raw_orders.py                              (modulo 1, leccion 6: la semana fija de 40 ordenes)
└── kiosko_time_traveled_dim_product.py         (este proyecto: junta las 7 piezas)

Con PyIceberg instalado en tu entorno (pip install "pyiceberg[sql-sqlite,pyarrow]", módulo 1, lección 4).

La solución de referencia, verificada

# kiosko_time_traveled_dim_product.py -- proyecto de cierre del modulo 3
# dim_product sin ninguna columna de historia + time travel recupera P002 V1
import os
from collections import defaultdict
from datetime import datetime

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

from raw_orders import RAW_ORDERS

DIM_STORE = [
    {"store_id": "S01", "store_name": "Kiosko Centro", "city": "Bogota"},
    {"store_id": "S02", "store_name": "Kiosko Norte", "city": "Lima"},
    {"store_id": "S03", "store_name": "Kiosko Sur", "city": "Santiago"},
]
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},
]

FACT_ORDERS_SCHEMA = Schema(
    NestedField(field_id=1, name="order_id", field_type=StringType(), required=True),
    NestedField(field_id=2, name="store_id", field_type=StringType(), required=True),
    NestedField(field_id=3, name="product_id", field_type=StringType(), required=True),
    NestedField(field_id=4, name="quantity", field_type=IntegerType(), required=True),
    NestedField(field_id=5, name="unit_price", field_type=DoubleType(), required=True),
    NestedField(field_id=6, name="revenue", field_type=DoubleType(), required=True),
    NestedField(field_id=7, name="order_ts", field_type=TimestampType(), required=True),
)
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),
)


def build_fact_orders_table() -> pa.Table:
    store_ids = {s["store_id"] for s in DIM_STORE}
    product_ids = {p["product_id"] for p in DIM_PRODUCT_V1}
    cols = {"order_id": [], "store_id": [], "product_id": [], "quantity": [],
            "unit_price": [], "revenue": [], "order_ts": []}
    for order_id, store_id, product_id, quantity, unit_price, ts in RAW_ORDERS:
        if store_id not in store_ids:
            raise ValueError(f"unknown store_id: {store_id}")
        if product_id not in product_ids:
            raise ValueError(f"unknown product_id: {product_id}")
        cols["order_id"].append(order_id)
        cols["store_id"].append(store_id)
        cols["product_id"].append(product_id)
        cols["quantity"].append(quantity)
        cols["unit_price"].append(unit_price)
        cols["revenue"].append(round(quantity * unit_price, 10))
        cols["order_ts"].append(datetime.fromisoformat(ts))
    schema = pa.schema([
        pa.field("order_id", pa.string(), nullable=False),
        pa.field("store_id", pa.string(), nullable=False),
        pa.field("product_id", pa.string(), nullable=False),
        pa.field("quantity", pa.int32(), nullable=False),
        pa.field("unit_price", pa.float64(), nullable=False),
        pa.field("revenue", pa.float64(), nullable=False),
        pa.field("order_ts", pa.timestamp("us"), nullable=False),
    ])
    return pa.Table.from_arrays(
        [
            pa.array(cols["order_id"], type=pa.string()),
            pa.array(cols["store_id"], type=pa.string()),
            pa.array(cols["product_id"], type=pa.string()),
            pa.array(cols["quantity"], type=pa.int32()),
            pa.array(cols["unit_price"], type=pa.float64()),
            pa.array(cols["revenue"], type=pa.float64()),
            pa.array(cols["order_ts"], type=pa.timestamp("us")),
        ],
        schema=schema,
    )


def dim_product_pa_table(rows: list[dict]) -> pa.Table:
    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),
    ])
    return pa.Table.from_pylist(rows, schema=schema)


def margin_by_category(dim_rows: list[dict], fact_rows: list[dict]) -> tuple[dict, dict]:
    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


def main() -> None:
    print("=== Kiosko: dim_product sin columnas de historia + time travel ===\n")

    warehouse_path = os.path.abspath("kiosko_warehouse")
    catalog_db_path = os.path.abspath("kiosko_catalog.db")
    os.makedirs(warehouse_path, exist_ok=True)
    catalog = load_catalog(
        "kiosko", type="sql",
        uri=f"sqlite:///{catalog_db_path}", warehouse=f"file://{warehouse_path}",
    )
    catalog.create_namespace("kiosko")
    print(f"Paso 1/7 -- catalogo '{catalog.name}' y namespace 'kiosko' listos")

    fact_orders = catalog.create_table("kiosko.fact_orders", schema=FACT_ORDERS_SCHEMA)
    fact_orders.append(build_fact_orders_table())
    print(f"Paso 2/7 -- kiosko.fact_orders cargada: {fact_orders.scan().to_arrow().num_rows} filas")

    dim_product = catalog.create_table("kiosko.dim_product", schema=DIM_PRODUCT_SCHEMA)
    print("Paso 3/7 -- kiosko.dim_product creada, columnas:",
          [f.name for f in dim_product.schema().fields], "(sin valid_from/valid_to/is_current)")

    dim_product.append(dim_product_pa_table(DIM_PRODUCT_V1))
    snap_v1 = dim_product.current_snapshot().snapshot_id
    print("Paso 4/7 -- V1 cargada con table.append(), snap_v1 capturado (P002=snacks/0.60)")

    dim_product.overwrite(dim_product_pa_table(DIM_PRODUCT_V2))
    print("Paso 5/7 -- V2 cargada con table.overwrite() (P002=health-snacks/0.68)")

    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 6/7 -- estado actual P002: {current_p002['category']}/{current_p002['unit_cost']}"
          f" | AS OF snap_v1 P002: {v1_p002['category']}/{v1_p002['unit_cost']}")

    fact_rows = fact_orders.scan().to_arrow().to_pylist()
    revenue_broken, margin_broken = margin_by_category(current_rows, fact_rows)
    revenue_correct, margin_correct = margin_by_category(v1_rows, fact_rows)
    print("Paso 7/7 -- margen por categoria calculado, roto vs. correcto\n")

    print("=== Verificacion final ===\n")
    print("ROTO (sin time travel, 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 (AS OF 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}")

    total_revenue = round(sum(f["revenue"] for f in fact_rows), 2)
    print(f"\nRevenue total (identico en ambos casos): {total_revenue}")

    assert dim_product.scan().to_arrow().num_rows == 4, "dim_product deberia tener 4 filas (una por producto)"
    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
    assert len(dim_product.schema().fields) == 4, "dim_product debe tener exactamente 4 columnas, ninguna de historia"
    print("\nTodas las verificaciones pasaron: 4 columnas en dim_product (cero de historia), "
          "margin roto=9.36, margin correcto=10.8 via time travel, revenue total=106.15.")


if __name__ == "__main__":
    main()

(raw_orders.py es exactamente el mismo archivo con las cuarenta órdenes fijas del módulo 1, lección 6 — no se repite aquí por espacio.)

Qué esperar (verificado corriendo python3 kiosko_time_traveled_dim_product.py real, de punta a punta, en un directorio nuevo; ningún snapshot_id se imprime como literal — la lección 4 de este módulo explica por qué):

=== Kiosko: dim_product sin columnas de historia + time travel ===

Paso 1/7 -- catalogo 'kiosko' y namespace 'kiosko' listos
Paso 2/7 -- kiosko.fact_orders cargada: 40 filas
Paso 3/7 -- kiosko.dim_product creada, columnas: ['product_id', 'product_name', 'category', 'unit_cost'] (sin valid_from/valid_to/is_current)
Paso 4/7 -- V1 cargada con table.append(), snap_v1 capturado (P002=snacks/0.60)
Paso 5/7 -- V2 cargada con table.overwrite() (P002=health-snacks/0.68)
Paso 6/7 -- estado actual P002: health-snacks/0.68 | AS OF snap_v1 P002: snacks/0.6
Paso 7/7 -- margen por categoria calculado, roto vs. correcto

=== Verificacion final ===

ROTO (sin time travel, 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 (AS OF 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 (identico en ambos casos): 106.15

Todas las verificaciones pasaron: 4 columnas en dim_product (cero de historia), margin roto=9.36, margin correcto=10.8 via time travel, revenue total=106.15.

Fíjate en los siete assert del final: no son decorativos. Confirman, en una sola pasada automática, cada una de las afirmaciones centrales de este módulo — que dim_product tiene el grano correcto (cuatro filas), que el estado vigente y el histórico son los que corresponden, que los dos márgenes coinciden exactamente con los números ya verificados por data-modeling y dbt, que el revenue total no se movió un centavo, y que el esquema de la tabla nunca tuvo más de cuatro columnas.

Diagrama: de dónde venías, a dónde llegaste

flowchart LR
    A["Modulo 1-2:\nfact_orders cargada,\nanatomia inspeccionada"] --> B["Leccion 2:\ndim_product creada,\nV1 cargada, snap_v1"]
    B --> C["Leccion 3:\ntable.overwrite(V2)\nP002 cambia"]
    C --> D["Leccion 4-5:\nsnap_v1 capturado bien,\ntime travel ejecutado"]
    D --> E["Leccion 6:\nmargen correcto 10.8\nrecuperado, 0 columnas"]
    E --> F["Leccion 7:\nel limite honesto,\nprobado con evidencia"]
    F --> G["Este proyecto:\nlas 7 piezas, un solo script,\nassert automatico"]
    G --> H["Modulo 4:\nevolucion de esquema\n(dim_store + country)"]

Cerrando la promesa del módulo, punto por punto

Lo que la lección 1 prometióEvidencia de que este módulo lo entregó
Cada escritura es un snapshot nuevoLección 2 (append), lección 3 (overwrite, que reveló dos snapshots internos: delete + append)
Recuperar P002 V1 sin ninguna columna de historiakiosko.dim_product, cuatro columnas, verificado en el assert final de este proyecto
El snapshot-id nunca se hardcodeaLección 4: evidencia de dos corridas con valores completamente distintos, y la técnica de recuperación por contenido
Time travel AS OF un snapshot-idLección 5: table.scan(snapshot_id=snap_v1), ejecutado, de solo lectura
El margen correcto de P002 (10.8), igual que data-modeling y dbtLección 6 y este proyecto: margin_correct["snacks"] == 10.8, verificado con assert
El límite honesto: qué NO resuelve el time travelLección 7: el experimento toy_price, con evidencia de que ningún snapshot_id único sirve a dos ventas con costos distintos a la vez

Este módulo no resolvió el caso general de una dimensión con múltiples cambios y hechos repartidos entre ellos — esa frontera quedó declarada, con evidencia, en la lección 7. Lo que este módulo entrega es exactamente lo que prometió: el caso favorable de Kiosko, resuelto con time travel puro, sin ninguna columna de historia, y verificado número por número contra las dos guías anteriores del ecosistema que ya resolvieron el mismo problema con técnicas distintas.

Errores comunes

Correr este proyecto sobre un catálogo que ya tiene kiosko.fact_orders o kiosko.dim_product de una lección anterior. Qué pasa: alguien corre este proyecto en el mismo directorio donde ya completó las lecciones 2 a 7, y catalog.create_table(...) falla porque las tablas ya están registradas. Por qué pasa: este proyecto repite, a propósito, la creación completa desde cero, para que sea autocontenido y reproducible sin depender del estado exacto que dejaron las lecciones anteriores. Cómo detectarlo: si ves TableAlreadyExistsError al correr kiosko_time_traveled_dim_product.py, ya tienes un catálogo con esas tablas registradas en el mismo directorio. Cómo corregirlo: corre este proyecto en un directorio de trabajo nuevo, separado de donde hiciste las lecciones 2 a 7 — tal como sugiere "El material" de esta lección.

Interpretar el assert round(margin_broken["health-snacks"], 2) == 9.36 como un error del script. Qué pasa: alguien, revisando el código, ve que el script hace un assert sobre un número que su propia sección llama "ROTO", y se pregunta si eso es un error en el proyecto. Por qué pasa: es fácil asumir que un assert siempre verifica "lo correcto", y confundirse al ver que este verifica, a propósito, el resultado que la propia guía llama incorrecto para el negocio. Cómo detectarlo: si dudas de por qué el script verifica un número "roto", relee el propósito: el assert no dice que 9.36 sea el margen correcto de negocio — dice que el cálculo sin time travel produce, de forma consistente y reproducible, ese número específico, exactamente como predicen data-modeling y dbt. Cómo corregirlo: nada que corregir — verificar el número roto es tan importante como verificar el correcto, porque confirma que el experimento completo —con y sin time travel— se comporta exactamente como esta guía predijo, no solo la mitad "bonita" del resultado.

Adaptar este proyecto para un caso real sin haber leído la lección 7 primero. Qué pasa: alguien toma el patrón de este proyecto —crear una tabla sin columnas de historia, confiar en time travel para recuperar versiones anteriores— y lo aplica directamente a un problema de producción, sin verificar si ese problema cumple la condición de la lección 7 (todos los hechos relevantes del mismo lado de cada cambio de dimensión). Por qué pasa: el resultado de este proyecto es limpio y convincente, y es tentador copiar el patrón sin revisar sus condiciones de aplicabilidad. Cómo detectarlo: si tu caso real tiene una dimensión que cambia más de una vez, o hechos con fechas repartidas a ambos lados de un cambio, y planeas usar solo table.scan(snapshot_id=...) para reconstruir la historia, estás en riesgo de este error. Cómo corregirlo: antes de aplicar este patrón fuera de esta guía, revisa la lección 7 completa y confirma explícitamente que tu caso cumple la condición favorable — si no la cumple, necesitas valid_from/valid_to a nivel de fila, no time travel puro.

Ejercicios

Ejercicio 1 — Corre el proyecto completo tú mismo, desde cero. En un directorio nuevo, con solo raw_orders.py y kiosko_time_traveled_dim_product.py, corre python3 kiosko_time_traveled_dim_product.py. Confirma que ves los siete pasos completarse y el mensaje final "Todas las verificaciones pasaron".

Ver solución

Si raw_orders.py está en el mismo directorio y PyIceberg está instalado, la salida debería reproducir exactamente la estructura de esta lección: siete pasos numerados, seguidos de la verificación final con los cuatro márgenes correctos, 106.15 de revenue total, y el mensaje de éxito con las siete condiciones confirmadas.

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.

Ver solución

El primer assert en fallar debería ser 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—. Si corrigieras solo ese assert para que acepte 0.70, el siguiente en fallar sería assert round(margin_broken["health-snacks"], 2) == 9.36, porque el margen roto cambiaría con el nuevo costo. Este ejercicio demuestra que los assert de este proyecto están encadenados a los valores canónicos exactos de Kiosko — cualquier desviación se detecta de inmediato, en el primer punto donde deja de cumplirse.

Ejercicio 3 — Explica, en tus propias palabras, por qué este proyecto verifica siete condiciones distintas en vez de solo el margen final. En 3-4 frases, justifica por qué el assert final de este script no se limita a verificar margin_correct["snacks"] == 10.8, sino que también verifica el grano, el esquema, y el resultado roto.

Ver solución

Verificar solo el número final —10.8— confirmaría que el resultado es correcto, pero no confirmaría por qué es correcto: podría llegar a ese número por casualidad, con una tabla de grano equivocado, o con un esquema que sí tiene columnas de historia escondidas. Verificar las siete condiciones —el grano (4 filas), el esquema (4 columnas exactas), el estado roto (9.36) y el correcto (10.8), además del revenue total— confirma que el resultado correcto llega por el camino correcto: una tabla con el diseño exacto que este módulo se propuso construir, no solo un número que coincide al final. Es la misma disciplina de "verificar, no confiar" que ya viste en el proyecto del módulo 1, aplicada aquí a un experimento con más piezas móviles.

Resumen y siguiente paso: el cierre de este módulo

Con este proyecto cierras el módulo 3. Integraste las siete lecciones anteriores —creación de dim_product sin columnas de historia, dos escrituras, la disciplina de captura de snapshot_id, el viaje en el tiempo, el margen correcto recuperado, y el límite honesto de la técnica— en un solo script, corrido de punta a punta, con assert automáticos que confirman cada número contra data-modeling-for-analytics-guide y dbt-analytics-engineering-guide.

Kiosko tiene, por primera vez en este ecosistema, una dimensión con historia recuperable sin haber diseñado ni una sola columna para eso. kiosko.dim_product sigue teniendo exactamente cuatro columnas de negocio — el mecanismo que hace posible recuperar P002 V1 vive enteramente en el snapshot, no en el esquema.

Hacia dónde sigues. El módulo 4 —Evolución de esquema sin reescritura— toma kiosko.dim_store, la tabla de las tres tiendas de Kiosko, y le agrega una columna nueva —country, derivada de city— sin reescribir ni un solo archivo de datos existente. Vas a ver por qué esa garantía —agregar, renombrar o borrar una columna sin tocar los archivos Parquet ya escritos— depende exactamente del mismo field_id que ya conociste en el módulo 1, y qué garantiza "ACID" en este contexto preciso.

Recursos

  • PyIceberg — documentación oficial (quickstart), el flujo completo de catálogo, tabla, append() y overwrite() que integra este proyecto. py.iceberg.apache.org. En inglés.
  • PyIceberg — referencia de API, table.scan(snapshot_id=...), table.current_snapshot(), table.history(). py.iceberg.apache.org/api. En inglés.
  • DISEÑO de data-modeling-for-analytics-guide — fuente de los números canónicos 9.36/10.8 que este proyecto verifica con assert. src/guides/data-modeling-for-analytics-guide/DISENO.md. En español.
  • DISEÑO de dbt-analytics-engineering-guide — fuente de dbt snapshot, la segunda confirmación independiente de los mismos números. src/guides/dbt-analytics-engineering-guide/DISENO.md. En español.
  • DISEÑO de esta guía — el mapa completo de los ocho módulos, incluido el módulo 4 que sigue. src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.