Módulo 8: Project Kioskos Lakehouse
Proyecto: el primer lakehouse de Kiosko
Descripción
Este proyecto cierra el módulo 8 — y con él, cierra la guía completa. Tienes las cinco tablas del lakehouse de Kiosko construidas, lección a lección, en el mismo catálogo: fact_orders y dim_store con country y dim_date (lección 3), dim_product sin columnas de historia con el margen correcto de P002 recuperado por time travel (lección 4), fact_orders_at_scale particionada y evolucionada (lección 5), y la confirmación de que table.upsert() produce el mismo resultado que table.overwrite() (lección 6). Sabes, con evidencia citada, qué siete fronteras deja pendientes este lakehouse y qué guía hermana resuelve cada una (lección 7). Falta un solo paso: reconstruir las cinco tablas, desde cero, en un solo script, con assert automáticos que confirman cada número — el mismo patrón de cierre que ya usó cada uno de los siete módulos anteriores de esta guía, aplicado ahora al lakehouse completo.
Conexión con el módulo. Este proyecto no introduce ningún concepto nuevo — es la integración final de las siete lecciones anteriores de este módulo, y de los siete módulos anteriores de toda la guía. Retoma, de forma literal, la promesa que abrió esta guía en su módulo 1: convertir un archivo Parquet suelto en una tabla real, con catálogo, esquema y snapshot — ahora multiplicada por cinco tablas, conviviendo en el mismo lakehouse, con el mismo revenue (106.15) y el mismo margen correcto de P002 (10.8) que ya confirmaron data-modeling-for-analytics-guide y dbt-analytics-engineering-guide.
Una analogía: la maqueta completa, presentada de una sola vez
Las lecciones 3 a 6 de este módulo construyeron, una pieza a la vez, las cinco tablas del lakehouse de Kiosko: los cimientos y la estructura (lección 3), la pieza central con su propia historia recuperable (lección 4), el ala a escala (lección 5), y la confirmación de que dos vías de construcción distintas llegan al mismo resultado (lección 6). Este proyecto es el momento de repetir todo el proceso, de punta a punta, en un solo gesto continuo — la misma integración final que ya cerró los módulos 1, 3, 4, 5, 6 y 7 de esta guía, ahora aplicada al lakehouse completo, no a una sola tabla.
El material: todo lo que este módulo construyó, en un solo lugar
Necesitas, en un directorio de trabajo nuevo:
kiosko_first_lakehouse/
├── raw_orders.py (modulo 1, leccion 6: la semana fija de 40 ordenes)
├── kiosko_scale.py (modulo 5: el generador determinista a escala)
└── kiosko_first_lakehouse.py (este proyecto: junta las 5 tablas)
Con PyIceberg instalado en tu entorno (pip install "pyiceberg[sql-sqlite,pyarrow]", módulo 1, lección 4).
Nota de escala. Este proyecto carga las 10,000,000 filas completas de fact_orders_at_scale —las mismas 250,000 franquicias de siempre—; en una laptop moderna, el script completo toma entre treinta segundos y un minuto. La parte lenta es, con diferencia, la generación y carga de esa tabla — el resto corre casi instantáneo.
La solución de referencia, verificada
# kiosko_first_lakehouse.py -- proyecto de cierre de la guia (modulo 8)
# el lakehouse completo de Kiosko sobre Apache Iceberg, ensamblado de punta a punta
import os
from collections import defaultdict
from datetime import date, datetime, timedelta
import pyarrow as pa
import pyarrow.compute as pc
from pyiceberg.catalog import load_catalog
from pyiceberg.partitioning import PartitionField, PartitionSpec
from pyiceberg.schema import Schema
from pyiceberg.transforms import DayTransform, IdentityTransform
from pyiceberg.types import (
BooleanType, DateType, DoubleType, IntegerType, NestedField, StringType, TimestampType,
)
from kiosko_scale import generate_orders_at_scale
from raw_orders import RAW_ORDERS
NUM_FRANCHISES = 250_000
DAY_NAMES = ["Monday", "Tuesday", "Wednesday", "Thursday", "Friday", "Saturday", "Sunday"]
DIM_STORE_ROWS = [
{"store_id": "S01", "store_name": "Kiosko Centro", "city": "Bogota", "country": "Colombia"},
{"store_id": "S02", "store_name": "Kiosko Norte", "city": "Lima", "country": "Peru"},
{"store_id": "S03", "store_name": "Kiosko Sur", "city": "Santiago", "country": "Chile"},
]
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_STORE_SCHEMA = Schema(
NestedField(field_id=1, name="store_id", field_type=StringType(), required=True),
NestedField(field_id=2, name="store_name", field_type=StringType(), required=True),
NestedField(field_id=3, name="city", field_type=StringType(), required=True),
NestedField(field_id=4, name="country", field_type=StringType(), required=True),
)
DIM_DATE_SCHEMA = Schema(
NestedField(field_id=1, name="date_key", field_type=IntegerType(), required=True),
NestedField(field_id=2, name="calendar_date", field_type=DateType(), required=True),
NestedField(field_id=3, name="day_of_week", field_type=StringType(), required=True),
NestedField(field_id=4, name="month", field_type=IntegerType(), required=True),
NestedField(field_id=5, name="quarter", field_type=IntegerType(), required=True),
NestedField(field_id=6, name="year", field_type=IntegerType(), required=True),
NestedField(field_id=7, name="is_weekend", field_type=BooleanType(), 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),
)
FACT_ORDERS_AT_SCALE_SCHEMA = Schema(
NestedField(field_id=1, name="order_id", field_type=StringType(), required=True),
NestedField(field_id=2, name="franchise_id", field_type=IntegerType(), required=True),
NestedField(field_id=3, name="store_id", field_type=StringType(), required=True),
NestedField(field_id=4, name="product_id", field_type=StringType(), required=True),
NestedField(field_id=5, name="quantity", field_type=IntegerType(), required=True),
NestedField(field_id=6, name="unit_price", field_type=DoubleType(), required=True),
NestedField(field_id=7, name="order_ts", field_type=TimestampType(), required=True),
)
INITIAL_SPEC = PartitionSpec(
PartitionField(source_id=3, field_id=1000, transform=IdentityTransform(), name="store_id"),
)
FACT_ORDERS_PA_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),
])
DIM_STORE_PA_SCHEMA = pa.schema([
pa.field("store_id", pa.string(), nullable=False),
pa.field("store_name", pa.string(), nullable=False),
pa.field("city", pa.string(), nullable=False),
pa.field("country", pa.string(), nullable=False),
])
DIM_DATE_PA_SCHEMA = pa.schema([
pa.field("date_key", pa.int32(), nullable=False),
pa.field("calendar_date", pa.date32(), nullable=False),
pa.field("day_of_week", pa.string(), nullable=False),
pa.field("month", pa.int32(), nullable=False),
pa.field("quarter", pa.int32(), nullable=False),
pa.field("year", pa.int32(), nullable=False),
pa.field("is_weekend", pa.bool_(), nullable=False),
])
DIM_PRODUCT_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),
])
FACT_AT_SCALE_PA_SCHEMA = pa.schema([
pa.field("order_id", pa.string(), nullable=False),
pa.field("franchise_id", pa.int32(), 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("order_ts", pa.timestamp("us"), nullable=False),
])
def fact_orders_pa_table() -> pa.Table:
rows = [
{"order_id": oid, "store_id": sid, "product_id": pid, "quantity": qty, "unit_price": price,
"revenue": round(qty * price, 10), "order_ts": datetime.fromisoformat(ts)}
for oid, sid, pid, qty, price, ts in RAW_ORDERS
]
return pa.Table.from_pylist(rows, schema=FACT_ORDERS_PA_SCHEMA)
def dim_date_pa_table(start_date: str, end_date: str) -> pa.Table:
start, end, rows = date.fromisoformat(start_date), date.fromisoformat(end_date), []
current = start
while current <= end:
weekday_index = current.weekday()
rows.append({
"date_key": int(current.strftime("%Y%m%d")), "calendar_date": current,
"day_of_week": DAY_NAMES[weekday_index], "month": current.month,
"quarter": (current.month - 1) // 3 + 1, "year": current.year,
"is_weekend": weekday_index >= 5,
})
current += timedelta(days=1)
return pa.Table.from_pylist(rows, schema=DIM_DATE_PA_SCHEMA)
def revenue_of(arrow_table: pa.Table) -> float:
line_revenue = pc.multiply(pc.cast(arrow_table.column("quantity"), pa.float64()), arrow_table.column("unit_price"))
return float(pc.sum(line_revenue).as_py())
def rows_to_pa_table(num_franchises: int, franchise_offset: int = 0) -> pa.Table:
rows = list(generate_orders_at_scale(num_franchises))
for r in rows:
r["franchise_id"] += franchise_offset
if franchise_offset:
r["order_id"] = r["order_id"].replace(
f"F{r['franchise_id'] - franchise_offset:06d}", f"F{r['franchise_id']:06d}",
)
r["order_ts"] = datetime.fromisoformat(r["order_ts"])
return pa.Table.from_pylist(rows, schema=FACT_AT_SCALE_PA_SCHEMA)
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
def main() -> None:
print("=== Kiosko: el primer lakehouse completo, sobre Apache Iceberg ===\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/12 -- catalogo '{catalog.name}' y namespace 'kiosko' listos")
# -- fact_orders --
fact_orders = catalog.create_table("kiosko.fact_orders", schema=FACT_ORDERS_SCHEMA)
fact_orders.append(fact_orders_pa_table())
fact_rows = fact_orders.scan().to_arrow()
fact_revenue = round(sum(fact_rows.column("revenue").to_pylist()), 2)
print(f"Paso 2/12 -- kiosko.fact_orders: {fact_rows.num_rows} filas, revenue {fact_revenue}")
# -- dim_store, con country desde el primer commit --
dim_store = catalog.create_table("kiosko.dim_store", schema=DIM_STORE_SCHEMA)
dim_store.append(pa.Table.from_pylist(DIM_STORE_ROWS, schema=DIM_STORE_PA_SCHEMA))
store_rows = sorted(dim_store.scan().to_arrow().to_pylist(), key=lambda r: r["store_id"])
print(f"Paso 3/12 -- kiosko.dim_store: {len(store_rows)} filas, "
f"country={[r['country'] for r in store_rows]}")
# -- dim_date --
dim_date = catalog.create_table("kiosko.dim_date", schema=DIM_DATE_SCHEMA)
dim_date.append(dim_date_pa_table("2026-08-01", "2026-08-31"))
date_rows = dim_date.scan().to_arrow().num_rows
print(f"Paso 4/12 -- kiosko.dim_date: {date_rows} filas (agosto 2026 completo)")
# -- dim_product: V1, snap_v1, overwrite a V2 --
dim_product = catalog.create_table("kiosko.dim_product", schema=DIM_PRODUCT_SCHEMA)
dim_product.append(pa.Table.from_pylist(DIM_PRODUCT_V1, schema=DIM_PRODUCT_PA_SCHEMA))
snap_v1 = dim_product.current_snapshot().snapshot_id
print(f"Paso 5/12 -- kiosko.dim_product: V1 cargada, snap_v1 capturado (P002=snacks/0.60)")
dim_product.overwrite(pa.Table.from_pylist(DIM_PRODUCT_V2, schema=DIM_PRODUCT_PA_SCHEMA))
current_p002 = next(r for r in dim_product.scan().to_arrow().to_pylist() if r["product_id"] == "P002")
print(f"Paso 6/12 -- kiosko.dim_product: V2 aplicada con overwrite() "
f"(P002={current_p002['category']}/{current_p002['unit_cost']}, vigencia 2026-08-15)")
# -- el star unido con time travel: margen roto vs correcto --
fact_rows_list = 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()
revenue_broken, margin_broken = margin_by_category(current_rows, fact_rows_list)
revenue_correct, margin_correct = margin_by_category(v1_rows, fact_rows_list)
print(f"Paso 7/12 -- star unido (fact_orders + dim_store + dim_date + dim_product): "
f"margin roto (health-snacks)={round(margin_broken['health-snacks'], 2)}, "
f"margin correcto AS OF snap_v1 (snacks)={round(margin_correct['snacks'], 2)}")
# -- table.upsert(): la via nativa, verificada contra el mismo resultado --
dim_product_check = catalog.create_table("kiosko.dim_product_merge_check", schema=DIM_PRODUCT_SCHEMA)
dim_product_check.append(pa.Table.from_pylist(DIM_PRODUCT_V1, schema=DIM_PRODUCT_PA_SCHEMA))
only_p002 = [r for r in DIM_PRODUCT_V2 if r["product_id"] == "P002"]
upsert_result = dim_product_check.upsert(
pa.Table.from_pylist(only_p002, schema=DIM_PRODUCT_PA_SCHEMA), join_cols=["product_id"],
)
check_rows = sorted(dim_product_check.scan().to_arrow().to_pylist(), key=lambda r: r["product_id"])
same_result = sorted(current_rows, key=lambda r: r["product_id"]) == check_rows
print(f"Paso 8/12 -- table.upsert() (via nativa): rows_updated={upsert_result.rows_updated}, "
f"rows_inserted={upsert_result.rows_inserted}, resultado identico a overwrite(): {same_result}")
# -- fact_orders_at_scale: particionada y evolucionada --
fact_at_scale = catalog.create_table(
"kiosko.fact_orders_at_scale", schema=FACT_ORDERS_AT_SCALE_SCHEMA, partition_spec=INITIAL_SPEC,
)
bulk_pa_table = rows_to_pa_table(NUM_FRANCHISES)
fact_at_scale.append(bulk_pa_table)
scale_row_count = fact_at_scale.scan().to_arrow().num_rows
scale_revenue = revenue_of(bulk_pa_table)
print(f"Paso 9/12 -- kiosko.fact_orders_at_scale: {scale_row_count} filas cargadas, "
f"revenue {round(scale_revenue, 2)}")
s01_scan = fact_at_scale.scan(row_filter="store_id == 'S01'").to_arrow()
s01_revenue = revenue_of(s01_scan)
print(f"Paso 10/12 -- consulta oculta store_id == 'S01': {s01_scan.num_rows} filas, "
f"revenue {round(s01_revenue, 2)}")
with fact_at_scale.update_spec() as update:
update.add_field("order_ts", DayTransform(), "order_day")
new_franchise = rows_to_pa_table(1, franchise_offset=NUM_FRANCHISES)
fact_at_scale.append(new_franchise)
partitions = fact_at_scale.inspect.partitions().to_pylist()
spec_ids_present = sorted({p["spec_id"] for p in partitions})
scale_row_count_final = fact_at_scale.scan().to_arrow().num_rows
scale_revenue_final = round(revenue_of(fact_at_scale.scan().to_arrow()), 2)
print(f"Paso 11/12 -- spec evolucionado (DayTransform sobre order_ts), franquicia nueva cargada: "
f"{scale_row_count_final} filas finales, revenue {scale_revenue_final}, "
f"spec_id conviviendo: {spec_ids_present}")
all_tables = sorted(t[1] for t in catalog.list_tables("kiosko"))
print(f"Paso 12/12 -- namespace 'kiosko' completo: {all_tables}\n")
print("=== Verificacion final del lakehouse completo ===\n")
assert fact_rows.num_rows == 40
assert fact_revenue == 106.15
assert {r["store_id"]: r["country"] for r in store_rows} == {
"S01": "Colombia", "S02": "Peru", "S03": "Chile",
}
assert date_rows == 31
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 round(margin_broken["health-snacks"], 2) == 9.36
assert round(margin_correct["snacks"], 2) == 10.8
assert round(revenue_correct["snacks"], 2) == round(revenue_broken["health-snacks"], 2) == 21.6
assert upsert_result.rows_updated == 1 and upsert_result.rows_inserted == 0
assert same_result, "overwrite() y upsert() deben producir el mismo resultado de negocio"
assert scale_row_count == NUM_FRANCHISES * 40 == 10_000_000
assert round(scale_revenue, 2) == 26_537_500.00
assert s01_scan.num_rows == 4_000_000
assert round(s01_revenue, 2) == 9_575_000.00
assert spec_ids_present == [0, 1]
assert scale_row_count_final == 10_000_040
assert scale_revenue_final == 26_537_606.15
assert all_tables == [
"dim_date", "dim_product", "dim_product_merge_check", "dim_store", "fact_orders", "fact_orders_at_scale",
]
print("Todas las verificaciones pasaron -- el lakehouse completo de Kiosko sobre Apache Iceberg:")
print(f" - kiosko.fact_orders: 40 filas, revenue total {fact_revenue}")
print(f" - kiosko.dim_store: 3 filas, country=Colombia/Peru/Chile")
print(f" - kiosko.dim_date: 31 filas (agosto 2026)")
print(f" - kiosko.dim_product: 4 columnas, cero de historia, P002 historizado por time travel")
print(f" margin CORRECTO (AS OF snap_v1, snacks) = 10.8")
print(f" margin ROTO (vigente, health-snacks) = 9.36")
print(f" -- identicos a data-modeling-for-analytics-guide M8 y dbt-analytics-engineering-guide M8")
print(f" - table.upsert() reproduce, fila por fila, el mismo resultado que table.overwrite()")
print(f" - kiosko.fact_orders_at_scale: 10,000,040 filas, revenue {scale_revenue_final}, "
f"2 esquemas de particion conviviendo")
if __name__ == "__main__":
main()
(raw_orders.py y kiosko_scale.py son exactamente los mismos archivos del módulo 1, lección 6, y del módulo 5 — no se repiten aquí por espacio.)
Qué esperar (verificado corriendo python3 kiosko_first_lakehouse.py real, de punta a punta, en un directorio nuevo; el runtime completo fue de aproximadamente 32 segundos; ningún snapshot_id se imprime como literal — la regla dura de esta guía, desde el módulo 3):
=== Kiosko: el primer lakehouse completo, sobre Apache Iceberg ===
Paso 1/12 -- catalogo 'kiosko' y namespace 'kiosko' listos
Paso 2/12 -- kiosko.fact_orders: 40 filas, revenue 106.15
Paso 3/12 -- kiosko.dim_store: 3 filas, country=['Colombia', 'Peru', 'Chile']
Paso 4/12 -- kiosko.dim_date: 31 filas (agosto 2026 completo)
Paso 5/12 -- kiosko.dim_product: V1 cargada, snap_v1 capturado (P002=snacks/0.60)
Paso 6/12 -- kiosko.dim_product: V2 aplicada con overwrite() (P002=health-snacks/0.68, vigencia 2026-08-15)
Paso 7/12 -- star unido (fact_orders + dim_store + dim_date + dim_product): margin roto (health-snacks)=9.36, margin correcto AS OF snap_v1 (snacks)=10.8
Paso 8/12 -- table.upsert() (via nativa): rows_updated=1, rows_inserted=0, resultado identico a overwrite(): True
Paso 9/12 -- kiosko.fact_orders_at_scale: 10000000 filas cargadas, revenue 26537500.0
Paso 10/12 -- consulta oculta store_id == 'S01': 4000000 filas, revenue 9575000.0
Paso 11/12 -- spec evolucionado (DayTransform sobre order_ts), franquicia nueva cargada: 10000040 filas finales, revenue 26537606.15, spec_id conviviendo: [0, 1]
Paso 12/12 -- namespace 'kiosko' completo: ['dim_date', 'dim_product', 'dim_product_merge_check', 'dim_store', 'fact_orders', 'fact_orders_at_scale']
=== Verificacion final del lakehouse completo ===
Todas las verificaciones pasaron -- el lakehouse completo de Kiosko sobre Apache Iceberg:
- kiosko.fact_orders: 40 filas, revenue total 106.15
- kiosko.dim_store: 3 filas, country=Colombia/Peru/Chile
- kiosko.dim_date: 31 filas (agosto 2026)
- kiosko.dim_product: 4 columnas, cero de historia, P002 historizado por time travel
margin CORRECTO (AS OF snap_v1, snacks) = 10.8
margin ROTO (vigente, health-snacks) = 9.36
-- identicos a data-modeling-for-analytics-guide M8 y dbt-analytics-engineering-guide M8
- table.upsert() reproduce, fila por fila, el mismo resultado que table.overwrite()
- kiosko.fact_orders_at_scale: 10,000,040 filas, revenue 26537606.15, 2 esquemas de particion conviviendo
Diecisiete assert, ninguno decorativo: confirman el revenue de cada tabla, el country de las tres tiendas, el grano y el esquema de dim_product, los dos márgenes de P002 contra los números canónicos de data-modeling-for-analytics-guide y dbt-analytics-engineering-guide, la equivalencia exacta entre overwrite() y upsert(), y el estado final completo de fact_orders_at_scale — seis tablas visibles en el namespace (las cinco del lakehouse más la tabla de comprobación de la lección 6), todas conviviendo en el mismo catálogo.
Criterio: cuándo gana un lakehouse sobre Iceberg, y frente a qué
El brief de la lección 2 pidió, además del lakehouse ensamblado, una decisión con criterio: ¿cuándo un lakehouse sobre Iceberg es la elección correcta, frente a un warehouse gestionado clásico (BigQuery, Snowflake) o frente a Parquet plano sin más? La respuesta no es "Iceberg siempre gana" — es una tabla de trade-offs concretos, con la evidencia de esta guía como base:
| Dimensión | Parquet plano (foundations, M6) | Warehouse gestionado (BigQuery/Snowflake) | Lakehouse sobre Iceberg (esta guía) |
|---|---|---|---|
| Atomicidad de escritura | No — overwrite-partition deja una ventana real de riesgo (módulo 4, lección 2) | Sí, gestionada por el motor, sin que el usuario la piense | Sí, verificada con CommitFailedException real (módulo 4, lección 7) |
| Time travel / historia | Ninguna — hay que diseñar columnas a mano | Sí, con retención limitada y sintaxis propia del vendor | Sí, nativo, snapshots ilimitados mientras no se expiren (módulo 3) |
| Costo de cómputo | $0, pero sin motor de consulta propio | Por consulta o por warehouse reservado — cobra aunque los datos no cambien | $0 en almacenamiento local; el motor de consulta (Spark, Trino, DuckDB) se elige aparte |
| Portabilidad entre motores | Alta — cualquier motor lee Parquet | Baja — los datos viven dentro del warehouse del vendor | Alta — el mismo kiosko.dim_product lo leyeron PyIceberg y Spark, sin duplicar datos (módulo 6) |
| Evolución de esquema sin downtime | No — cambiar una columna implica reescribir archivos | Sí, con sintaxis propia del motor | Sí, verificado sin tocar un solo archivo de datos (módulo 4) |
| Control operativo (catálogo, mantenimiento) | Ninguno — es responsabilidad de quien escribe cada archivo | Ninguno — lo gestiona el vendor, con menos visibilidad | Total, pero exige operarlo (catálogos, expire_snapshots, módulo 7) |
| Curva de aprendizaje | Baja | Baja para SQL, alta para entender el costo real | Media-alta — exige entender snapshots, specs de partición, catálogos |
La lectura honesta de esta tabla, con la evidencia de los ocho módulos de esta guía: un lakehouse sobre Iceberg gana cuando el problema central es portabilidad de datos entre motores (los mismos archivos Parquet, leídos por PyIceberg y por Spark, sin duplicar nada) combinada con necesidad real de time travel y evolución de esquema sin reescritura — exactamente los dos problemas que las cuatro guías anteriores de este ecosistema ya habían tropezado sin resolver desde la capa de almacenamiento. Un warehouse gestionado gana cuando el equipo prioriza cero operación —nadie en Kiosko quiere pensar en catálogos ni en expire_snapshots— a cambio de aceptar el costo por consulta y la dependencia de un solo vendor. Y Parquet plano, sin más, sigue siendo la elección correcta cuando no hace falta ninguna de las garantías que esta guía completa demostró —sin historia que recuperar, sin evolución de esquema esperada, sin necesidad de que dos motores distintos lean la misma tabla—, el mismo caso exacto que data-engineering-foundations-guide cubrió honestamente en su capa gold, sin fingir que necesitaba más.
Diagrama: el lakehouse completo, las ocho lecciones de este módulo
flowchart TB
L1["L1: presentacion,\nlas 5 piezas"] --> L2["L2: el brief,\nun solo catalogo"]
L2 --> L3["L3: fact_orders,\ndim_store, dim_date"]
L3 --> L4["L4: dim_product,\ntime travel, 10.8/9.36"]
L4 --> L5["L5: fact_orders_at_scale,\n10M filas, particionada"]
L5 --> L6["L6: table.upsert(),\nmismo resultado"]
L6 --> L7["L7: 7 guias hermanas,\nque falta"]
L7 --> L8["L8 (este proyecto):\nlas 5 tablas, un solo script,\n17 assert"]
L8 --> END["Guia completa CERRADA:\nM1-M8, es/, 106.15/10.8\nverificados de punta a punta"]
Cerrando la promesa de toda la guía, módulo por módulo
| Lo que cada módulo prometió | Evidencia de que esta guía lo entregó |
|---|---|
| M1: de Parquet suelto a tabla Iceberg real | kiosko.fact_orders, catálogo, esquema, snapshot — 40 filas, 106.15 |
| M2: la anatomía completa, catálogo → data files | Mapeada en disco y con table.inspect, sin sorpresas en este capstone |
| M3: time travel, cero columnas de historia | dim_product, snap_v1, margen correcto 10.8 — reproducido en L4 de este módulo |
| M4: evolución de esquema sin reescribir | dim_store con country desde el primer commit de este módulo — el mecanismo ya verificado en M4 |
| M5: partición oculta y evolucionada | fact_orders_at_scale, 10,000,040 filas, 2 spec_id conviviendo — reproducido en L5 de este módulo |
M6: MERGE INTO y upserts nativos | table.upsert() idéntico a overwrite(), verificado byte a byte en L6 de este módulo |
| M7: catálogos, mantenimiento, Delta Lake | Nombrados y contrastados; higiene operativa con evidencia real |
| M8: el lakehouse completo | Las cinco tablas, un catálogo, 106.15/10.8 verificados — este proyecto |
Ningún módulo de esta guía queda con una promesa sin evidencia. El revenue total de Kiosko (106.15) y el margen correcto de P002 (10.8, vía time travel) son, en este punto, los mismos números que ya confirmaron siete guías distintas del ecosistema data-engineering-ecosystem —data-engineering-foundations-guide, python-for-data-engineering-guide, data-modeling-for-analytics-guide, dbt-analytics-engineering-guide, spark-and-distributed-processing-guide, y ahora esta— cada una con un motor y una técnica distintos, todas de acuerdo.
Errores comunes
Correr este proyecto sobre un catálogo que ya tiene las tablas de una lección anterior de este módulo. Qué pasa: alguien corre este proyecto en el mismo directorio donde ya completó las lecciones 3 a 6, 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_first_lakehouse.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 3 a 6 — tal como sugiere "El material" de esta lección.
Interpretar la tabla de criterio como una fórmula automática, sin contexto de negocio. Qué pasa: alguien intenta convertir la tabla de "cuándo gana un lakehouse" en una checklist mecánica: contar cuántas filas favorecen a cada opción y elegir la de más votos. Por qué pasa: una tabla comparativa se presta a leerse como un marcador, en vez de como un mapa de trade-offs que dependen del contexto real de cada equipo. Cómo detectarlo: si tu conclusión es "Iceberg ganó 5 filas contra 2, así que siempre es la mejor opción", perdiste el punto central de esta sección — la fila de "curva de aprendizaje" y "control operativo" son, precisamente, los costos reales que un equipo con poca capacidad de operación no debería ignorar. Cómo corregirlo: usa la tabla para identificar cuáles de las siete dimensiones son las que de verdad importan en tu contexto —¿Kiosko necesita portabilidad entre motores, o solo un dashboard mensual?—, y deja que esa prioridad, no un conteo de filas, decida.
Asumir que kiosko.dim_product_merge_check es un error de este proyecto, porque no aparece en el brief de la lección 2. Qué pasa: alguien, al ver seis tablas en el Paso 12/12 en vez de cinco, piensa que el script tiene un bug. Por qué pasa: el brief y la lección 1 de este módulo hablan de "las cinco tablas del lakehouse", y dim_product_merge_check no es una de ellas. Cómo detectarlo: revisa el Paso 8 — esa tabla existe únicamente para verificar, con evidencia aislada, que table.upsert() produce el mismo resultado que table.overwrite() sobre kiosko.dim_product, exactamente como hizo la lección 6 de este módulo. Cómo corregirlo: no es un error — es la misma tabla de comprobación de la lección 6, incluida en este proyecto integrado porque su verificación (same_result == True) es una de las diecisiete que confirma el assert final.
Ejercicios
Ejercicio 1 — Corre el proyecto completo tú mismo, desde cero. En un directorio nuevo, con raw_orders.py y kiosko_scale.py en el mismo lugar, corre python3 kiosko_first_lakehouse.py. Confirma que ves los doce pasos completarse y el mensaje final de éxito.
Ver solución
Si raw_orders.py y kiosko_scale.py están en el mismo directorio y PyIceberg está instalado, la salida debería reproducir exactamente la estructura de esta lección: doce pasos numerados, seguidos de la verificación final con las seis confirmaciones de negocio. El Paso 9 —la carga de diez millones de filas— es el más lento; el resto corre en segundos.
Ejercicio 2 — Rompe un assert a propósito, y observa el fallo. Cambia temporalmente DIM_PRODUCT_V2 para que P002 tenga unit_cost=0.75 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 ese assert para aceptar 0.75, 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 confirma, una vez más, que los assert de este proyecto están encadenados a los valores canónicos exactos de Kiosko, verificados en cascada.
Ejercicio 3 — Explica, en tus propias palabras, por qué la tabla de criterio de esta lección no incluye una fila de "velocidad de consulta". En 3-4 frases, justifica por qué esta guía, deliberadamente, no comparó el rendimiento de consulta entre Parquet plano, un warehouse gestionado, e Iceberg.
Ver solución
Comparar velocidad de consulta con rigor exigiría medir tiempos de ejecución reales, bajo condiciones de carga comparables, sobre volúmenes de datos representativos de producción — exactamente el tipo de análisis que el DISEÑO de esta guía delegó, a propósito, a advanced-sql-querying-guide (planes de ejecución, tuning) y a spark-and-distributed-processing-guide (cómputo distribuido a fondo). Incluir una fila de "velocidad" sin esa rigurosidad —basada en impresiones en vez de mediciones controladas— habría violado la misma disciplina de "verificar, no confiar" que sostuvo cada assert de esta guía completa. La tabla de criterio se limita, a propósito, a las dimensiones que esta guía sí verificó con evidencia directa: atomicidad, time travel, portabilidad, evolución de esquema, costo de almacenamiento y curva de aprendizaje — no inventa un número que nunca midió.
Resumen y siguiente paso: el cierre de toda la guía
Con este proyecto cierras el módulo 8 — y con él, cierras lakehouse-and-iceberg-guide completa. Ensamblaste las cinco tablas del lakehouse de Kiosko en un solo catálogo, con diecisiete assert automáticos que confirman: revenue total 106.15, country correcto en las tres tiendas, 31 filas de calendario, margen de P002 recuperado con time travel (10.8 correcto, 9.36 roto — los mismos números que data-modeling-for-analytics-guide y dbt-analytics-engineering-guide ya confirmaron con motores distintos), equivalencia exacta entre table.overwrite() y table.upsert(), y 10,000,040 filas de fact_orders_at_scale con dos esquemas de partición conviviendo.
Kiosko tiene, al cierre de esta guía, un lakehouse real sobre Apache Iceberg: transacciones ACID verificadas, snapshots automáticos en cada commit, time travel sin ninguna columna de historia, evolución de esquema y de partición sin reescribir un solo archivo, y dos vías nativas de merge que convergen en el mismo resultado. Cada una de esas garantías era, al abrir esta guía en su módulo 1, un techo real que Parquet-como-simple-archivo no podía resolver por sí solo — el mismo techo que data-engineering-foundations-guide, data-modeling-for-analytics-guide, dbt-analytics-engineering-guide y spark-and-distributed-processing-guide ya habían tocado, cada una a su manera, antes de esta guía.
La lección 7 de este módulo ya nombró, con evidencia precisa, las siete fronteras que este lakehouse deja pendientes, y la guía hermana exacta que resuelve cada una. Esta guía no promete resolverlas — promete, y entrega, un formato de tabla real, funcionando de punta a punta, sobre el mismo caso de Kiosko que seis guías anteriores del ecosistema ya usaron para verificar sus propias garantías.
Recursos
- PyIceberg — documentación oficial (quickstart), el flujo completo que integra este proyecto:
load_catalog(),create_namespace(),create_table(),append(),overwrite(),upsert(),update_spec(). py.iceberg.apache.org. En inglés. - PyIceberg — referencia de API completa. py.iceberg.apache.org/api. En inglés.
- Apache Iceberg — documentación oficial, versión de referencia 1.11.0. iceberg.apache.org/docs/latest. En inglés.
- DISEÑO de
data-modeling-for-analytics-guide— fuente de los números canónicos106.15/10.8/9.36que este proyecto verifica conassert, por segunda guía consecutiva.src/guides/data-modeling-for-analytics-guide/DISENO.md. En español. - DISEÑO de
dbt-analytics-engineering-guide— tercera confirmación independiente de los mismos números.src/guides/dbt-analytics-engineering-guide/DISENO.md. En español. - DISEÑO de
spark-and-distributed-processing-guide— fuente delfact_orders.parquetoriginal y del dataset a escala que este proyecto reconstruye completo.src/guides/spark-and-distributed-processing-guide/DISENO.md. En español. - DISEÑO de esta guía — el mapa completo de los ocho módulos, la fuente única de verdad que este proyecto cierra.
src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.