Módulo 8: Project Kioskos Lakehouse

Evolucionando y particionando a escala

Descripción

Esta lección agrega la quinta y última tabla del lakehouse de Kiosko: kiosko.fact_orders_at_scale, con diez millones de filas, particionada por store_id con partición oculta, y evolucionada hacia adelante con DayTransform sobre order_ts — exactamente el mismo mecanismo que el módulo 5 de esta guía ya demostró, ejecutado ahora dentro del catálogo único de este capstone, junto a las otras cuatro tablas que las lecciones 3 y 4 ya cargaron.

Conexión con el módulo. Esta lección responde la tercera exigencia del brief de la lección 2: "particionar y evolucionar sin reescribir". A diferencia de kiosko.dim_product en la lección 4, kiosko.fact_orders_at_scale no tiene ninguna relación de negocio con el cambio de P002 — es un hilo independiente, heredado de spark-and-distributed-processing-guide, que este capstone integra al mismo catálogo por ser, igual que las otras cuatro, una tabla real del lakehouse de Kiosko.

Una analogía: la quinta pieza, fabricada al mismo ritmo que las otras cuatro

Si las lecciones 3 y 4 fabricaron los cimientos, la estructura y las divisiones interiores de la maqueta, esta lección fabrica el ala nueva del edificio — una construida a otra escala (diez millones de filas contra apenas cuarenta), con su propio sistema de organización interna (partición oculta), pero anclada al mismo terreno: el catálogo kiosko que ya sostiene las otras cuatro piezas. No hace falta que el ala nueva "converse" con las otras cuatro en el sentido de compartir datos —de hecho, no comparte ninguna fila con fact_orders, dim_store, dim_date ni dim_product—; hace falta que sea, físicamente, parte del mismo edificio: el mismo kiosko_catalog.db, el mismo namespace kiosko.

Ejemplo trabajado: fact_orders_at_scale, particionada y evolucionada, en el mismo catálogo

Paso 1 — El generador determinista, idéntico al del módulo 5

# kiosko_scale.py -- el generador determinista, identico al del modulo 5
from typing import Any, Dict, Iterator

from raw_orders import RAW_ORDERS

KIOSKO_WEEK = [
    {"order_id": order_id, "store_id": store_id, "product_id": product_id,
     "quantity": quantity, "unit_price": unit_price, "order_ts": ts}
    for order_id, store_id, product_id, quantity, unit_price, ts in RAW_ORDERS
]


def generate_orders_at_scale(num_franchises: int) -> Iterator[Dict[str, Any]]:
    for franchise_id in range(num_franchises):
        for row in KIOSKO_WEEK:
            yield {
                "order_id": f"F{franchise_id:06d}-{row['order_id']}",
                "franchise_id": franchise_id,
                "store_id": row["store_id"],
                "product_id": row["product_id"],
                "quantity": row["quantity"],
                "unit_price": row["unit_price"],
                "order_ts": row["order_ts"],
            }

Paso 2 — El esquema, el spec inicial, y la carga completa de 250,000 franquicias

# kiosko_fact_orders_at_scale.py -- modulo 8, leccion 5
import os
from datetime import datetime

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 DoubleType, IntegerType, NestedField, StringType, TimestampType

from kiosko_scale import generate_orders_at_scale

NUM_FRANCHISES = 250_000

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"),
)
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 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=PA_SCHEMA)

Fíjate en field_id=1000 del PartitionField: el mismo rango alto, separado de los field_id 1-7 de las columnas, que ya usó el módulo 5 — evita cualquier colisión entre el identificador de una columna y el de un campo de partición.

Paso 3 — El flujo completo: carga, consulta oculta, evolución, verificación

def main() -> None:
    print("=== Kiosko: fact_orders_at_scale, particionada y evolucionada ===\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}",
    )
    table = catalog.create_table(
        "kiosko.fact_orders_at_scale", schema=FACT_ORDERS_AT_SCALE_SCHEMA, partition_spec=INITIAL_SPEC,
    )
    print(f"Paso 1/6 -- kiosko.fact_orders_at_scale creada con spec inicial: {table.spec()}")

    bulk_pa_table = rows_to_pa_table(NUM_FRANCHISES)
    table.append(bulk_pa_table)
    row_count = table.scan().to_arrow().num_rows
    total_revenue = revenue_of(bulk_pa_table)
    print(f"Paso 2/6 -- {row_count} filas cargadas ({NUM_FRANCHISES} franquicias x 40), "
          f"revenue total {round(total_revenue, 2)}")

    s01_scan = table.scan(row_filter="store_id == 'S01'").to_arrow()
    s01_revenue = revenue_of(s01_scan)
    files_touched_s01 = len(list(table.scan(row_filter="store_id == 'S01'").plan_files()))
    files_touched_all = len(list(table.scan().plan_files()))
    print(f"Paso 3/6 -- consulta oculta store_id == 'S01': {s01_scan.num_rows} filas, "
          f"revenue {round(s01_revenue, 2)}, toco {files_touched_s01}/{files_touched_all} archivos")

    with table.update_spec() as update:
        update.add_field("order_ts", DayTransform(), "order_day")
    print(f"Paso 4/6 -- spec evolucionado: {table.spec()}")

    new_franchise_pa_table = rows_to_pa_table(1, franchise_offset=NUM_FRANCHISES)
    table.append(new_franchise_pa_table)
    print(f"Paso 5/6 -- franquicia {NUM_FRANCHISES} ({new_franchise_pa_table.num_rows} filas) "
          f"aterriza bajo el spec nuevo (con order_day)")

    partitions = table.inspect.partitions().to_pylist()
    spec_ids_present = sorted({p["spec_id"] for p in partitions})
    rows_by_spec = {
        spec_id: sum(p["record_count"] for p in partitions if p["spec_id"] == spec_id)
        for spec_id in spec_ids_present
    }
    row_count_final = table.scan().to_arrow().num_rows
    revenue_final = revenue_of(table.scan().to_arrow())
    print(f"Paso 6/6 -- estado final: {row_count_final} filas, revenue {round(revenue_final, 2)}, "
          f"spec_id presentes: {spec_ids_present}, filas por spec: {rows_by_spec}\n")

    print("=== Verificacion final ===\n")
    assert row_count == NUM_FRANCHISES * 40 == 10_000_000
    assert round(total_revenue, 2) == 26_537_500.00
    assert s01_scan.num_rows == 4_000_000
    assert round(s01_revenue, 2) == 9_575_000.00
    assert files_touched_s01 < files_touched_all, "la consulta filtrada debe tocar menos archivos que el scan completo"
    assert spec_ids_present == [0, 1], "ambos esquemas de particion deben convivir en la misma tabla"
    assert rows_by_spec[0] == 10_000_000
    assert rows_by_spec[1] == 40
    assert row_count_final == 10_000_040
    assert round(revenue_final, 2) == 26_537_606.15

    print("Todas las verificaciones pasaron:")
    print("  - 10,000,000 filas bajo IdentityTransform sobre store_id, S01 = 9,575,000.00, poda real de archivos")
    print("  - update_spec().add_field(DayTransform) no reescribio ningun archivo existente")
    print("  - inspect.partitions() confirma los dos esquemas de particion conviviendo (spec_id 0 y 1)")
    print(f"  - estado final: 10,000,040 filas, revenue {round(revenue_final, 2)}")


if __name__ == "__main__":
    main()

Qué esperar (verificado corriendo python3 kiosko_fact_orders_at_scale.py real, en el mismo directorio que las lecciones 3 y 4, sin borrar kiosko_warehouse/; en una laptop moderna, la generación y carga de las diez millones de filas toma alrededor de medio minuto):

=== Kiosko: fact_orders_at_scale, particionada y evolucionada ===

Paso 1/6 -- kiosko.fact_orders_at_scale creada con spec inicial: [
  1000: store_id: identity(3)
]
Paso 2/6 -- 10000000 filas cargadas (250000 franquicias x 40), revenue total 26537500.0
Paso 3/6 -- consulta oculta store_id == 'S01': 4000000 filas, revenue 9575000.0, toco 1/3 archivos
Paso 4/6 -- spec evolucionado: [
  1000: store_id: identity(3)
  1001: order_day: day(7)
]
Paso 5/6 -- franquicia 250000 (40 filas) aterriza bajo el spec nuevo (con order_day)
Paso 6/6 -- estado final: 10000040 filas, revenue 26537606.15, spec_id presentes: [0, 1], filas por spec: {0: 10000000, 1: 40}

=== Verificacion final ===

Todas las verificaciones pasaron:
  - 10,000,000 filas bajo IdentityTransform sobre store_id, S01 = 9,575,000.00, poda real de archivos
  - update_spec().add_field(DayTransform) no reescribio ningun archivo existente
  - inspect.partitions() confirma los dos esquemas de particion conviviendo (spec_id 0 y 1)
  - estado final: 10,000,040 filas, revenue 26537606.15

El revenue final (26,537,606.15) es la suma de las diez millones de filas originales (26,537,500.00) más la franquicia nueva (106.15, la misma semana de Kiosko, una vez más) — la misma aritmética exacta que ya verificó el proyecto de cierre del módulo 5.

Diagrama: cinco tablas, un solo catálogo

flowchart TB
    subgraph CAT["kiosko_catalog.db -- un solo catalogo"]
        FO["kiosko.fact_orders\n40 filas (L3)"]
        DS["kiosko.dim_store\n3 filas, con country (L3)"]
        DD["kiosko.dim_date\n31 filas (L3)"]
        DP["kiosko.dim_product\n4 filas, 2 snapshots,\nsin historia (L4)"]
        FAS["kiosko.fact_orders_at_scale\n10,000,040 filas,\n2 specs (L5)"]
    end
    L3["Leccion 3"] --> FO
    L3 --> DS
    L3 --> DD
    L4["Leccion 4"] --> DP
    L5["Leccion 5 (esta)"] --> FAS

Profundización: por qué esta tabla no interactúa con las otras cuatro

A diferencia de dim_product en la lección 4 —que se une explícitamente contra fact_orders, dim_store y dim_date para calcular un margen—, kiosko.fact_orders_at_scale no comparte ninguna fila, ninguna clave, ningún cálculo con el resto del lakehouse en esta guía. Eso no es un descuido — es una decisión heredada, explícita, del propio DISEÑO de spark-and-distributed-processing-guide, que generó este dataset sintético (franchise_id de 0 a 249,999) para demostrar volumen, no para modelar un negocio real de franquicias con su propia jerarquía dimensional. La razón por la que esta tabla vive en el mismo catálogo kiosko, entonces, no es "porque se relaciona con las otras cuatro" — es porque es, físicamente, una tabla más del mismo lakehouse de Kiosko, con la misma disciplina de partición oculta y evolución que el resto de esta guía enseñó. Un lakehouse real, en producción, casi siempre tiene tablas que se conectan entre sí (como fact_orders y dim_product) y tablas que conviven en el mismo catálogo sin relacionarse directamente (como esta) — ambos casos son parte normal de operar un catálogo compartido.

Errores comunes

Esperar que kiosko.fact_orders_at_scale comparta algún store_id o cálculo con kiosko.dim_store. Qué pasa: alguien intenta escribir un JOIN entre fact_orders_at_scale y dim_store esperando enriquecer las filas a escala con el country que la lección 3 pobló. Por qué pasa: ambas tablas comparten la columna store_id, y es razonable asumir que cualquier store_id compartido implica una relación de negocio pensada para unirse. Cómo detectarlo: si tu JOIN corre sin error pero el resultado no aparece en ningún assert de esta lección ni de ningún proyecto anterior, no rompiste nada — simplemente estás explorando una combinación que esta guía nunca verificó ni prometió. Cómo corregirlo: store_id en fact_orders_at_scale sí usa los mismos valores (S01, S02, S03) que dim_store —por diseño, ambas tablas describen las mismas tres tiendas de Kiosko—, así que ese JOIN es técnicamente válido y hasta interesante como ejercicio propio; solo ten claro que ningún resultado de esa combinación forma parte de los números canónicos que este capstone verifica.

Correr esta lección sin haber corrido antes las lecciones 3 y 4 en el mismo directorio. Qué pasa: alguien corre kiosko_fact_orders_at_scale.py en un directorio nuevo, sin kiosko_warehouse/ ni kiosko_catalog.db previos. Por qué pasa: a diferencia de los proyectos de cierre de módulos anteriores —que sí eran autocontenidos, pensados para correr solos—, las lecciones 3 a 6 de este módulo comparten, a propósito, el mismo catálogo. Cómo detectarlo: si corres esta lección sola, en realidad no falla —catalog.create_table("kiosko.fact_orders_at_scale", ...) crea el catálogo y el namespace si no existen—, pero al final tendrías un catálogo con solo esta tabla, no el lakehouse completo de cinco tablas que este módulo describe. Cómo corregirlo: para reproducir el lakehouse completo tal como lo describe esta guía, corre las lecciones 3, 4 y 5 en orden, en el mismo directorio de trabajo, sin borrar nada entre ellas.

Ejercicios

Ejercicio 1 — Corre el script tú mismo, en el mismo directorio que las lecciones 3 y 4. Confirma que ves los seis pasos completarse y el mensaje final con las cinco verificaciones. Ten paciencia con el Paso 2 — es el paso más lento de todo este módulo.

Ver solución

Si corriste las lecciones 3 y 4 primero, en el mismo directorio, la salida debería reproducir exactamente la estructura de esta lección: seis pasos numerados, seguidos de la verificación final con 10,000,040 filas y 26,537,606.15 de revenue. El Paso 2 —la generación y carga de diez millones de filas— es, con diferencia, el más lento de las ocho lecciones de este módulo; el resto corre casi instantáneo.

Ejercicio 2 — Calcula, sin correr código, cuántos archivos de datos esperarías ver bajo spec_id=0 después del Paso 2, y compáralo con lo que viste en el módulo 5. Usa files_touched_all del Paso 3 como pista.

Ver solución

files_touched_all en el Paso 3 vale 3 — el mismo número que ya confirmó el proyecto de cierre del módulo 5: un archivo de datos por cada valor distinto de store_id (S01, S02, S03), consecuencia directa de que IdentityTransform sobre store_id agrupa físicamente todas las filas de una misma tienda en el mismo archivo, dentro del mismo append(). Ese número no cambia en este módulo porque el mecanismo de partición —IdentityTransform sobre una columna de baja cardinalidad, tres valores posibles— es exactamente el mismo que ya se demostró en el módulo 5; este capstone lo reproduce, no lo modifica.

Ejercicio 3 — Explica, en tus propias palabras, por qué esta lección no necesita capturar ningún snapshot_id en una variable, a diferencia de la lección 4. En 2-3 frases, justifica la diferencia entre las dos lecciones.

Ver solución

La lección 4 necesita capturar snap_v1 porque su objetivo central es recuperar un estado anterior con time travel —sin la variable capturada, no habría forma de pedirle a Iceberg "muéstrame cómo era dim_product antes del cambio"—. Esta lección, en cambio, nunca necesita volver a un estado anterior de fact_orders_at_scale: la evolución de partición (update_spec()) es una operación hacia adelante, que no borra ni oculta ningún dato existente, y el assert final de esta lección verifica el estado vigente de la tabla (row_count_final, revenue_final), no un snapshot específico del pasado. Por eso table.history() y snapshot_id no aparecen en el flujo principal de esta lección, aunque siguen existiendo —cualquier snapshot intermedio de esta tabla sigue siendo recuperable, simplemente esta lección no necesita pedirlo—.

Resumen y siguiente paso

En esta lección cargaste la quinta y última tabla del lakehouse de Kiosko: kiosko.fact_orders_at_scale, con diez millones de filas particionadas por store_id, evolucionada con DayTransform sobre order_ts, y con una franquicia nueva aterrizando bajo el spec evolucionado — 10,000,040 filas finales, 26,537,606.15 de revenue, dos esquemas de partición conviviendo en la misma tabla. El lakehouse de Kiosko tiene, ahora, sus cinco tablas completas, en el mismo catálogo.

Antes de avanzar deberías poder: explicar por qué fact_orders_at_scale no comparte ningún cálculo de negocio con las otras cuatro tablas del lakehouse, a pesar de vivir en el mismo catálogo; y calcular de memoria por qué el revenue final es 26,537,500.00 + 106.15.

La lección 6 vuelve a kiosko.dim_product — no para cambiarla otra vez, sino para demostrar que el mismo cambio de P002 que la lección 4 aplicó con table.overwrite() produce, con table.upsert(), exactamente el mismo resultado de negocio, por un camino distinto: la vía nativa del formato.

Recursos

  • PyIceberg — documentación oficial (quickstart), el flujo de PartitionSpec, append() y update_spec() que integra esta lección. py.iceberg.apache.org. En inglés.
  • PyIceberg — referencia de API, los transforms (IdentityTransform, DayTransform), table.scan(...).plan_files(), table.inspect.partitions(). py.iceberg.apache.org/api. En inglés.
  • Esta misma guía, módulo 5, lección 8 — fuente del script original que esta lección reutiliza dentro del catálogo único del capstone. ../module-05-hidden-partitioning-and-partition-evolution/es/08-project-kioskos-partitioned-at-scale.md. En español.
  • DISEÑO de spark-and-distributed-processing-guide — fuente de generate_orders_at_scale() y los números exactos del dataset a escala. 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, incluida la lección 6 que sigue. src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.