Módulo 7: Catalogs Maintenance And Delta Lake By Contrast

Proyecto: la tabla mantenida de Kiosko

Descripción

Este proyecto cierra el módulo 7. Conociste los cuatro catálogos de producción, sin implementar ninguno (lección 2). Mediste, con evidencia real, cómo cinco noches de un pipeline redundante multiplicaron por más de cuatro el número de snapshots de kiosko.dim_product (lección 3). Viste el problema distinto de los archivos pequeños, sobre kiosko.fact_orders_daily_batches (lección 4, representativo en la parte de compactación). Podaste la tabla de verdad con expire_snapshots(), protegiendo snap_v1 (lección 5). Identificaste los archivos huérfanos que esa poda dejó atrás (lección 6, representativo en la parte de borrado). Y contrastaste todo esto, una sola vez, con Delta Lake (lección 7). Falta un solo paso: juntar las piezas ejecutables de este módulo en un solo script, con assert automáticos que confirmen 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. A diferencia de otros proyectos de cierre de esta guía, este es explícito sobre una frontera doble: reconstruye y verifica de punta a punta todo lo que corrió de verdad en este módulo (expire_snapshots, la acumulación de snapshots, la identificación de huérfanos), y deja constancia, sin fingir ejecutarlas, de las dos operaciones que quedaron representativas (compactación, remove_orphan_files).

Una analogía: el reporte de mantenimiento, con lo hecho y lo pendiente en columnas separadas

Un taller mecánico que entrega un vehículo después de una revisión no mezcla, en un solo párrafo vago, lo que hizo con lo que recomienda para la próxima visita. Entrega una orden de trabajo con dos columnas claras: "hecho hoy" —con la firma del mecánico y las piezas reemplazadas— y "pendiente, recomendado" —con el diagnóstico exacto y qué se necesitaría para resolverlo—. Este proyecto es esa orden de trabajo: una columna con expire_snapshots() ejecutado de verdad, once snapshots podados, verificado con assert; otra columna con la compactación y remove_orphan_files que este entorno no puede correr, documentadas con la misma precisión, sin fingir que ya se hicieron.

El material: un directorio de trabajo nuevo

kiosko_maintained_table_project/
├── raw_orders_by_day.py                     (las 40 ordenes de siempre, agrupadas por dia)
└── kiosko_maintained_table_project.py       (este proyecto: junta las piezas ejecutables)

Con PyIceberg instalado en tu entorno (pip install "pyiceberg[sql-sqlite,pyarrow]", módulo 1, lección 4). Este proyecto es autocontenido: crea kiosko.dim_product y kiosko.fact_orders_daily_batches desde cero, así que corre en un directorio nuevo, separado de las lecciones 3 a 6 de este módulo.

La solución de referencia, verificada

# kiosko_maintained_table_project.py -- proyecto de cierre del modulo 7
import os
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_by_day import RAW_ORDERS_BY_DAY


def main() -> None:
    print("=== Kiosko: kiosko.dim_product, mantenida de punta a punta ===\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")

    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),
    )
    table = catalog.create_table("kiosko.dim_product", schema=dim_product_schema)

    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},
    ]

    # Paso 1/8 -- reconstruye el estado del modulo 3: V1 + el cambio real de P002
    table.append(pa.Table.from_pylist(DIM_PRODUCT_V1, schema=pa_schema))
    snap_v1 = table.current_snapshot().snapshot_id
    table.overwrite(pa.Table.from_pylist(DIM_PRODUCT_V2, schema=pa_schema))
    snapshots_after_m3 = len(table.history())
    print(f"Paso 1/8 -- estado del modulo 3 reconstruido: {snapshots_after_m3} snapshots, snap_v1 capturado")

    # Paso 2/8 -- 5 noches redundantes (leccion 3 de este modulo)
    nights = ["2026-08-16", "2026-08-17", "2026-08-18", "2026-08-19", "2026-08-20"]
    for _ in nights:
        table.overwrite(pa.Table.from_pylist(DIM_PRODUCT_V2, schema=pa_schema))
    snapshots_before_expire = len(table.history())
    all_files_before = table.inspect.all_data_files().num_rows
    live_files_before = table.inspect.files().num_rows
    print(f"Paso 2/8 -- 5 noches redundantes aplicadas: {snapshots_before_expire} snapshots, "
          f"{all_files_before} archivos rastreados, {live_files_before} vivo(s)")

    # Paso 3/8 -- construye la lista de poda: todo menos snap_v1 y el vigente
    history = table.history()
    current_snapshot_id = table.current_snapshot().snapshot_id
    to_expire = [e.snapshot_id for e in history if e.snapshot_id not in (snap_v1, current_snapshot_id)]
    print(f"Paso 3/8 -- {len(to_expire)} snapshots seleccionados para expirar "
          f"(protegidos: snap_v1 y el vigente)")

    # Paso 4/8 -- expire_snapshots() REAL
    table.maintenance.expire_snapshots().by_ids(to_expire).commit()
    table.refresh()
    snapshots_after_expire = len(table.history())
    all_files_after = table.inspect.all_data_files().num_rows
    print(f"Paso 4/8 -- expire_snapshots().by_ids() corrido: {snapshots_after_expire} snapshots, "
          f"{all_files_after} archivos rastreados")

    # Paso 5/8 -- verifica que el time travel y el estado vigente siguen intactos
    v1_rows = table.scan(snapshot_id=snap_v1).to_arrow().to_pylist()
    p002_v1 = next(r for r in v1_rows if r["product_id"] == "P002")
    current_rows = table.scan().to_arrow().to_pylist()
    p002_current = next(r for r in current_rows if r["product_id"] == "P002")
    print(f"Paso 5/8 -- AS OF snap_v1: P002={p002_v1['category']}/{p002_v1['unit_cost']} | "
          f"vigente: P002={p002_current['category']}/{p002_current['unit_cost']}")

    # Paso 6/8 -- identifica archivos huerfanos que quedaron en disco (diagnostico, solo lectura)
    tracked = {row["file_path"] for row in table.inspect.all_data_files().select(["file_path"]).to_pylist()}
    data_dir = os.path.join(warehouse_path, "kiosko", "dim_product", "data")
    on_disk = {"file://" + os.path.join(data_dir, fn) for fn in os.listdir(data_dir) if fn.endswith(".parquet")}
    orphans = on_disk - tracked
    orphan_bytes = sum(os.path.getsize(f.replace("file://", "")) for f in orphans)
    print(f"Paso 6/8 -- {len(on_disk)} archivos fisicos en disco, {len(orphans)} huerfanos "
          f"({orphan_bytes} bytes) -- remove_orphan_files es representativo, no se ejecuta aqui")

    # Paso 7/8 -- kiosko.fact_orders_daily_batches -- el problema de archivos pequenos (leccion 4)
    fact_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),
    )
    fact_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),
    ])
    fact_table = catalog.create_table("kiosko.fact_orders_daily_batches", schema=fact_schema)
    for day, orders in RAW_ORDERS_BY_DAY:
        rows = []
        for order_id, store_id, product_id, quantity, unit_price, ts in orders:
            rows.append({
                "order_id": order_id, "store_id": store_id, "product_id": product_id,
                "quantity": quantity, "unit_price": unit_price,
                "revenue": round(quantity * unit_price, 10),
                "order_ts": datetime.fromisoformat(ts),
            })
        fact_table.append(pa.Table.from_pylist(rows, schema=fact_pa_schema))
    fact_rows = fact_table.scan().to_arrow().to_pylist()
    fact_total_rows = len(fact_rows)
    fact_total_revenue = round(sum(r["revenue"] for r in fact_rows), 2)
    fact_live_files = fact_table.inspect.files().num_rows
    print(f"Paso 7/8 -- kiosko.fact_orders_daily_batches: {fact_total_rows} filas, "
          f"revenue={fact_total_revenue}, {fact_live_files} archivos vivos "
          f"(compactacion es representativa, no se ejecuta aqui)")

    print("\n=== Verificacion final ===\n")

    assert snapshots_after_m3 == 3, "el modulo 3 deberia dejar 3 snapshots (append, delete, append)"
    assert snapshots_before_expire == 13, "5 noches redundantes (2 snapshots c/u) + 3 = 13"
    assert len(to_expire) == 11
    assert snapshots_after_expire == 2, "solo snap_v1 y el vigente deberian sobrevivir"
    assert all_files_before == 7
    assert all_files_after == 2
    assert p002_v1["category"] == "snacks" and p002_v1["unit_cost"] == 0.60
    assert p002_current["category"] == "health-snacks" and p002_current["unit_cost"] == 0.68
    assert len(orphans) == 5, "los 5 archivos que expire_snapshots() dejo de rastrear siguen en disco"
    assert fact_total_rows == 40
    assert fact_total_revenue == 106.15
    assert fact_live_files == 7, "7 append() diarios, ninguno redundante, sin compactar"

    print("Todas las verificaciones pasaron:")
    print(f"  - kiosko.dim_product: {snapshots_after_m3} -> {snapshots_before_expire} snapshots "
          f"(5 noches redundantes) -> {snapshots_after_expire} (post expire_snapshots, REAL)")
    print("  - snap_v1 protegido: P002 sigue siendo snacks/0.60 via time travel")
    print("  - vigente intacto: P002 sigue siendo health-snacks/0.68")
    print(f"  - {len(orphans)} archivos huerfanos identificados ({orphan_bytes} bytes), "
          "remove_orphan_files queda representativo")
    print(f"  - kiosko.fact_orders_daily_batches: {fact_total_rows} filas, revenue={fact_total_revenue}, "
          f"{fact_live_files} archivos vivos, compactacion queda representativa")


if __name__ == "__main__":
    main()

(raw_orders_by_day.py agrupa las mismas cuarenta órdenes del módulo 1, lección 6, por el día real en que ocurrió cada una — no se repite aquí por espacio; es exactamente el diccionario que ya viste completo en la lección 4 de este módulo.)

Qué esperar (verificado corriendo python3 kiosko_maintained_table_project.py real, de punta a punta, en un directorio nuevo; ningún snapshot_id se imprime como literal):

=== Kiosko: kiosko.dim_product, mantenida de punta a punta ===

Paso 1/8 -- estado del modulo 3 reconstruido: 3 snapshots, snap_v1 capturado
Paso 2/8 -- 5 noches redundantes aplicadas: 13 snapshots, 7 archivos rastreados, 1 vivo(s)
Paso 3/8 -- 11 snapshots seleccionados para expirar (protegidos: snap_v1 y el vigente)
Paso 4/8 -- expire_snapshots().by_ids() corrido: 2 snapshots, 2 archivos rastreados
Paso 5/8 -- AS OF snap_v1: P002=snacks/0.6 | vigente: P002=health-snacks/0.68
Paso 6/8 -- 7 archivos fisicos en disco, 5 huerfanos (9185 bytes) -- remove_orphan_files es representativo, no se ejecuta aqui
Paso 7/8 -- kiosko.fact_orders_daily_batches: 40 filas, revenue=106.15, 7 archivos vivos (compactacion es representativa, no se ejecuta aqui)

=== Verificacion final ===

Todas las verificaciones pasaron:
  - kiosko.dim_product: 3 -> 13 snapshots (5 noches redundantes) -> 2 (post expire_snapshots, REAL)
  - snap_v1 protegido: P002 sigue siendo snacks/0.60 via time travel
  - vigente intacto: P002 sigue siendo health-snacks/0.68
  - 5 archivos huerfanos identificados (9185 bytes), remove_orphan_files queda representativo
  - kiosko.fact_orders_daily_batches: 40 filas, revenue=106.15, 7 archivos vivos, compactacion queda representativa

Once assert, ninguno decorativo: confirman el número exacto de snapshots en cada etapa (3 → 13 → 2), que la lista de poda fue exactamente la esperada (11), que el time travel y el estado vigente sobrevivieron intactos a la poda, que los archivos huérfanos identificados coinciden con lo que la lección 6 predijo (5, 9185 bytes), y que el problema de archivos pequeños de la lección 4 sigue siendo reproducible (7 archivos vivos, 40 filas, 106.15 de revenue).

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

flowchart LR
    A["Modulo 1-6:\nkiosko.dim_product real,\nsnap_v1 protegido"] --> B["Leccion 2:\n4 catalogos de produccion\nnombrados, no implementados"]
    B --> C["Leccion 3:\n5 noches redundantes,\n3 -> 13 snapshots"]
    C --> D["Leccion 4:\narchivos pequenos,\nfact_orders_daily_batches"]
    D --> E["Leccion 5:\nexpire_snapshots() REAL,\n13 -> 2 snapshots"]
    E --> F["Leccion 6:\n5 huerfanos identificados,\nremove_orphan_files representativo"]
    F --> G["Leccion 7:\nDelta Lake por contraste,\nconvergencia 2026"]
    G --> H["Este proyecto:\ntodo integrado,\n11 assert automaticos"]
    H --> I["Modulo 8:\ncapstone del lakehouse\nde Kiosko"]

Cerrando la promesa del módulo, punto por punto

Lo que la lección 1 prometióEvidencia de que este módulo lo entregó
Nombrar los catálogos de producción, sin implementarlosLección 2: REST, Glue, Unity Catalog, Polaris, cada uno con su garantía citada contra documentación oficial
Medir por qué los snapshots acumulan costoLección 3 y este proyecto: 3 → 13 snapshots, 7 archivos rastreados, 1 vivo, verificado con assert
Compactar archivos pequeñosLección 4: kiosko.fact_orders_daily_batches, 7 archivos vivos reales, compactación documentada como representativa (Spark, rewrite_data_files)
Expirar snapshots viejos con seguridadLección 5 y este proyecto: expire_snapshots().by_ids() ejecutado de verdad, snap_v1 protegido, time travel verificado intacto
Eliminar archivos huérfanos sin perder el time travel que sí se necesitaLección 6 y este proyecto: 5 huérfanos identificados con código real (9185 bytes), remove_orphan_files documentado como representativo
Delta Lake nombrado una sola vez, por contrasteLección 7: mecanismo de metadata, sintaxis de time travel, mantenimiento, y la convergencia de 2026 — sin construir ninguna tabla Delta

Este proyecto no dejó ninguna promesa de la lección 1 sin evidencia — incluidas las dos que este entorno no pudo ejecutar de verdad, documentadas con la misma honestidad que el resto de esta guía aplicó al MERGE INTO de Spark en el módulo 6.

Errores comunes

Correr este proyecto sobre un catálogo que ya tiene kiosko.dim_product o kiosko.fact_orders_daily_batches 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 reconstrucció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_maintained_table_project.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.

Asumir que all_files_after == 2 en el Paso 4 significa que el espacio en disco ya se liberó. Qué pasa: alguien, viendo que el assert del Paso 4 pasa con all_files_after == 2, concluye que el proyecto ya completó toda la limpieza posible. Por qué pasa: es fácil no distinguir, otra vez, entre "archivos rastreados por la metadata" (lo que expire_snapshots() sí cambia) y "archivos físicos en disco" (lo que solo remove_orphan_files cambiaría). Cómo detectarlo: revisa el assert len(orphans) == 5 del Paso 6 — ese mismo script, corriendo inmediatamente después de la poda, confirma que siguen existiendo cinco archivos físicos sin ninguna referencia. Cómo corregirlo: lee los dos assert en conjunto, no por separado — juntos cuentan la historia completa: la metadata ya está podada (2 archivos rastreados), pero el disco todavía tiene trabajo pendiente (5 huérfanos), exactamente la distinción que las lecciones 5 y 6 de este módulo desarrollaron a fondo.

Ejercicios

Ejercicio 1 — Corre el proyecto completo tú mismo, desde cero. En un directorio nuevo, con raw_orders_by_day.py en el mismo lugar, corre python3 kiosko_maintained_table_project.py. Confirma que ves los siete pasos completarse y el mensaje final con las cinco verificaciones.

Ver solución

Si raw_orders_by_day.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 cinco mensajes de éxito. Los snapshot_id de tu corrida van a ser distintos de cualquier ejemplo anterior de esta guía — eso es exactamente lo esperado.

Ejercicio 2 — Rompe un assert a propósito, y observa el fallo. Cambia temporalmente el número de noches redundantes de cinco a tres (nights = ["2026-08-16", "2026-08-17", "2026-08-18"]), 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 snapshots_before_expire == 13 — con tres noches en vez de cinco, cada una produciendo dos snapshots, el total después del Paso 2 sería 3 + 3*2 = 9, no 13. Si corrigieras ese número a 9, el siguiente en fallar sería assert len(to_expire) == 11 (ahora serían 9 - 2 = 7 candidatos), y en cascada, assert all_files_before == 7 también cambiaría. Este ejercicio confirma que los números de este proyecto están encadenados de forma precisa a la cantidad exacta de escrituras redundantes — cualquier cambio en el Paso 2 se propaga, de forma predecible, a cada verificación posterior.

Ejercicio 3 — Explica, en tus propias palabras, por qué este proyecto incluye tanto kiosko.dim_product (el problema de snapshots viejos) como kiosko.fact_orders_daily_batches (el problema de archivos pequeños) en el mismo script, en vez de limitarse a uno solo. En 3-4 frases, justifica esta decisión con lo que aprendiste en las lecciones 3 y 4.

Ver solución

Limitarse a un solo problema dejaría al proyecto incompleto frente a lo que la lección 1 de este módulo prometió: dos ejes de costo distintos, no uno. La lección 3 y la lección 4 de este módulo fueron explícitas en que "snapshots viejos que ya nadie necesita" y "archivos pequeños dentro del snapshot vigente" son problemas ortogonales —una tabla puede tener uno, el otro, ambos, o ninguno—, y confundirlos es exactamente uno de los errores comunes que la lección 4 documentó. Incluir ambos en el mismo proyecto de cierre, con assert separados para cada uno, confirma que quien complete este módulo puede reconocer y diagnosticar los dos problemas de forma independiente, en vez de asumir que resolver uno resuelve automáticamente el otro.

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

Con este proyecto cierras el módulo 7. Integraste, en un solo script con once assert automáticos, todo lo que este módulo ejecutó de verdad: la acumulación real de snapshots (3 → 13), la poda real con snap_v1 protegido (13 → 2), la identificación real de archivos huérfanos (5, 9185 bytes), y el problema real de archivos pequeños sobre una tabla de hechos cargada día por día (7 archivos vivos). Y dejaste constancia explícita, sin fingir ejecutarlas, de las dos operaciones que este entorno no puede correr en Python puro: compactación y remove_orphan_files, ambas documentadas con su sintaxis exacta de Spark.

Kiosko tiene, ahora, una tabla que no solo sabe recuperar su historia con time travel (módulo 3) — sabe, también, podar la parte de esa historia que ya no aporta nada, sin arriesgar la parte que sí.

Hacia dónde sigues. El módulo 8 —el capstone de esta guía— ensambla el lakehouse completo de Kiosko: fact_orders, dim_store (con country), dim_product (historizado por time travel, sin columnas), dim_date, y fact_orders_at_scale (particionada y evolucionada), todos juntos, con el mismo revenue total (106.15) y el mismo margen correcto de P002 (10.8) que ya confirmaron data-modeling-for-analytics-guide y dbt-analytics-engineering-guide. Y cierra con el mapa hacia las guías hermanas de este ecosistema — incluida esta, con sus catálogos de producción todavía sin implementar y su compactación todavía representativa, las dos fronteras exactas que aws-core-services-guide retoma.

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.maintenance, table.inspect.all_data_files(), table.inspect.files(), la base de los pasos 3 a 6 de este proyecto. py.iceberg.apache.org/api. En inglés.
  • Apache Iceberg — documentación oficial, "Maintenance", la fuente de las tres operaciones que este módulo completo organizó. iceberg.apache.org/docs/latest/maintenance. En inglés.
  • Esta misma guía, módulo 3, lección 8 — fuente del patrón de proyecto de cierre con assert encadenados que este proyecto reutiliza. ../module-03-snapshots-and-time-travel/es/08-project-kioskos-time-traveled-dim-product.md. En español.
  • DISEÑO de esta guía — el mapa completo de los ocho módulos, incluido el módulo 8 que sigue. src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.