Módulo 1: From File Format To Table Format

Proyecto: la primera tabla Iceberg de Kiosko

Descripción

Este proyecto cierra el módulo. Tienes PyIceberg instalado (lección 4), sabes crear un catálogo, un namespace y una tabla con esquema explícito (lecciones 4 y 5), sabes cargar datos reales con table.append() (lección 6), y sabes verificar que el resultado es correcto, no solo que "no lanzó ningún error" (lección 7). Falta un solo paso: juntar las cinco piezas en un solo script, corrido de punta a punta, que reproduce —con tus propias manos, en tu propia máquina— el momento exacto en que Kiosko dejó de tener "un archivo Parquet" y empezó a tener "una tabla Iceberg real".

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 pregunta que abrió este módulo en la lección 1: ¿qué le falta a un archivo Parquet suelto para comportarse como una tabla? Este proyecto no cierra esa pregunta todavía por completo —eso toma los ocho módulos completos de esta guía—, pero entrega la primera respuesta ejecutable: una tabla real, con un snapshot real, verificada contra el mismo número que ya confirmaron seis guías anteriores.

Una analogía: el álbum completo, de la portada en blanco a la primera colección archivada

Las lecciones 4 a 7 de este módulo construyeron, una pieza a la vez, el álbum completo de Kiosko: el bibliotecario instalado (lección 4), la portada y el índice impresos (lección 5), la primera colección de fotos pegada (lección 6), y la confirmación de que esas fotos son, de verdad, las correctas (lección 7). Este proyecto es el momento de repetir todo el proceso, de punta a punta, en un solo gesto continuo — como armar el álbum completo frente a un cliente, sin pausas entre pasos, para demostrar que el proceso completo funciona como una sola pieza coherente, no como cinco pasos que solo funcionan por separado.

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

Necesitas, en un directorio de trabajo nuevo:

kiosko_iceberg/
├── raw_orders.py                    (leccion 6: la semana fija de 40 ordenes)
└── kiosko_first_iceberg_table.py    (este proyecto: junta las 5 piezas)

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

La solución de referencia, verificada

# kiosko_first_iceberg_table.py -- proyecto de cierre del modulo 1
# de un archivo Parquet suelto a la primera tabla Iceberg real de Kiosko
import os
from datetime import datetime

import pyarrow as pa
import pyarrow.compute as pc
import pyarrow.parquet as pq
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 = [
    {"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},
]
store_names = {s["store_id"]: s["store_name"] for s in DIM_STORE}


def build_fact_orders_parquet(path: str) -> None:
    store_ids = {s["store_id"] for s in DIM_STORE}
    product_ids = {p["product_id"] for p in DIM_PRODUCT}

    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),
    ])
    pa_table = 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,
    )
    pq.write_table(pa_table, path)


def main() -> None:
    print("=== Kiosko: de fact_orders.parquet a la primera tabla Iceberg ===\n")

    build_fact_orders_parquet("fact_orders.parquet")
    print("Paso 1/5 -- fact_orders.parquet reconstruido con pyarrow")

    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}",
    )
    print(f"Paso 2/5 -- catalogo '{catalog.name}' cargado ({type(catalog).__name__})")

    catalog.create_namespace("kiosko")
    print(f"Paso 3/5 -- namespace creado: {catalog.list_namespaces()}")

    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),
    )
    table = catalog.create_table("kiosko.fact_orders", schema=fact_orders_schema)
    print(f"Paso 4/5 -- tabla creada: {table.name()} (snapshot={table.current_snapshot()})")

    pa_table = pq.read_table("fact_orders.parquet")
    table.append(pa_table)
    snap_id = table.current_snapshot().snapshot_id
    print(f"Paso 5/5 -- table.append() completado, snapshot_id capturado: {snap_id}\n")

    scanned = table.scan().to_arrow()
    total_rows = scanned.num_rows
    total_revenue = round(sum(scanned.column("revenue").to_pylist()), 2)

    print("=== Verificacion final ===\n")
    print(f"len(table.scan().to_arrow()) = {total_rows}")

    by_store = scanned.group_by("store_id").aggregate([("revenue", "sum")])
    print("\nRevenue by store:")
    for row in sorted(by_store.to_pylist(), key=lambda r: r["store_id"]):
        sid = row["store_id"]
        print(f"  {sid} {store_names[sid]:<14}: revenue={round(row['revenue_sum'], 2)}")
    print(f"\nTotal week revenue: {total_revenue}")

    keys = pc.binary_join_element_wise(scanned.column("order_id"), scanned.column("product_id"), "-")
    distinct_keys = pc.count_distinct(keys).as_py()
    print(f"\nGrano: total_rows={total_rows}, distinct_order_product_lines={distinct_keys}")

    assert total_rows == 40, f"esperaba 40 filas, obtuve {total_rows}"
    assert total_revenue == 106.15, f"esperaba 106.15, obtuve {total_revenue}"
    assert distinct_keys == 40, f"esperaba 40 combinaciones distintas, obtuve {distinct_keys}"
    print("\nTodas las verificaciones pasaron: 40 filas, 106.15 de revenue, 0 duplicados.")


if __name__ == "__main__":
    main()

(raw_orders.py es exactamente el mismo archivo con las cuarenta órdenes fijas de la lección 6 — no se repite aquí por espacio.)

Qué esperar (verificado corriendo python3 kiosko_first_iceberg_table.py real, de punta a punta, en un directorio nuevo; el snapshot_id es el que tu corrida asigna en el momento del commit — distinto cada vez, se muestra aquí como marcador):

=== Kiosko: de fact_orders.parquet a la primera tabla Iceberg ===

Paso 1/5 -- fact_orders.parquet reconstruido con pyarrow
Paso 2/5 -- catalogo 'kiosko' cargado (SqlCatalog)
Paso 3/5 -- namespace creado: [('kiosko',)]
Paso 4/5 -- tabla creada: ('kiosko', 'fact_orders') (snapshot=None)
Paso 5/5 -- table.append() completado, snapshot_id capturado: <snapshot-id asignado en tu corrida, distinto cada vez>

=== Verificacion final ===

len(table.scan().to_arrow()) = 40

Revenue by store:
  S01 Kiosko Centro : revenue=38.3
  S02 Kiosko Norte  : revenue=38.8
  S03 Kiosko Sur    : revenue=29.05

Total week revenue: 106.15

Grano: total_rows=40, distinct_order_product_lines=40

Todas las verificaciones pasaron: 40 filas, 106.15 de revenue, 0 duplicados.

Fíjate en el assert final del script: no es decorativo — si el revenue total no fuera exactamente 106.15, o si el conteo de filas no fuera 40, o si el grano tuviera algún duplicado, el script terminaría con un AssertionError en vez de con el mensaje de éxito. Esta es la misma disciplina de "verificar, no confiar" que la lección 7 explicó en su Profundización, ahora convertida en una comprobación automática que corre cada vez que ejecutas este proyecto.

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

flowchart LR
    A["Leccion 1-3:\nel problema y el vocabulario\n(archivo vs tabla)"] --> B["Leccion 4:\nPyIceberg instalado,\ncatalogo cargado"]
    B --> C["Leccion 5:\nnamespace + tabla\ncon esquema, vacia"]
    C --> D["Leccion 6:\ntable.append()\nprimer snapshot real"]
    D --> E["Leccion 7:\n106.15 verificado,\ngrano sin duplicados"]
    E --> F["Este proyecto:\nlas 5 piezas, un solo script,\nassert automatico"]
    F --> G["Modulo 2:\nanatomia de la tabla\n(que hay, de verdad, en disco)"]

En disco, después de este proyecto

kiosko_warehouse/kiosko/fact_orders/
├── data/
│   └── 00000-0-<uuid>.parquet
└── metadata/
    ├── 00000-<uuid>.metadata.json   (tabla vacia, paso 4)
    ├── 00001-<uuid>.metadata.json   (primer snapshot, paso 5)
    ├── <uuid>-m0.avro               (manifest file)
    └── snap-<snapshot_id>-0-<uuid>.avro   (manifest list)

Exactamente la misma estructura que ya viste en la lección 6 — este proyecto no agrega ningún archivo nuevo al warehouse, solo reproduce el mismo flujo con las cinco piezas juntas en un único script, en vez de tres scripts separados.

Cerrando la promesa de la lección 1, punto por punto

Lo que la lección 1 prometióEvidencia de que este módulo lo entregó
Nombrar las cuatro veces que Parquet solo no alcanzóLección 2: código citado literal de foundations, data-modeling, dbt y spark
Definir con precisión formato de archivo vs formato de tablaLección 3: la distinción y las seis capacidades que Iceberg agrega
Instalar PyIceberg de verdad, local, sin JVMpyiceberg==0.11.1 instalado y verificado (lección 4)
Crear un catálogo, un namespace y una tabla con esquema realSqlCatalog + kiosko + kiosko.fact_orders con 7 columnas tipadas (lecciones 4-5)
Cargar el fact_orders.parquet de Kiosko dentro de Icebergtable.append() ejecutado, primer snapshot real (lección 6)
Verificar que el revenue sigue siendo 106.15S01=38.3/S02=38.8/S03=29.05, total 106.15, sin duplicados (lección 7, y este proyecto)

Ninguna fila de esta tabla resuelve todavía uno de los cuatro problemas de la lección 2 a fondo —eso empieza recién en el módulo 2 (anatomía) y se resuelve, uno por uno, en los módulos 3 a 6—. Lo que este módulo entrega es la base: una tabla Iceberg real, cargada, verificada, sobre la que el resto de esta guía construye cada garantía específica.

Errores comunes

Correr el proyecto sobre una tabla que ya existe de una lección anterior, y toparse con TableAlreadyExistsError. Qué pasa: alguien corre este proyecto en el mismo directorio donde ya completó las lecciones 4 a 7, y catalog.create_table("kiosko.fact_orders", ...) falla porque la tabla ya está registrada en kiosko_catalog.db. Por qué pasa: este proyecto repite, a propósito, los mismos pasos de creación que ya corriste antes, para que sea autocontenido y reproducible desde cero. Cómo detectarlo: si ves TableAlreadyExistsError (o un mensaje equivalente) al correr kiosko_first_iceberg_table.py, ya tienes un catálogo con esa tabla registrada en el mismo directorio. Cómo corregirlo: corre este proyecto en un directorio de trabajo nuevo, separado de donde hiciste las lecciones 4 a 7 —tal como sugiere "El material" de esta lección—, o borra kiosko_catalog.db y kiosko_warehouse/ del directorio anterior antes de repetir el proyecto ahí.

Ejecutar el script sin raw_orders.py en el mismo directorio. Qué pasa: alguien copia solo kiosko_first_iceberg_table.py, sin raw_orders.py, y el script falla de inmediato con ModuleNotFoundError: No module named 'raw_orders'. Por qué pasa: el import de la primera línea del script espera encontrar ese archivo en el mismo directorio de trabajo (o en el PYTHONPATH). Cómo detectarlo: el error aparece antes de imprimir siquiera la primera línea de main() — si tu salida no muestra ni el encabezado === Kiosko: ... ===, revisa que raw_orders.py exista junto al script. Cómo corregirlo: copia raw_orders.py de la lección 6 al mismo directorio que kiosko_first_iceberg_table.py, exactamente como muestra "El material" de esta lección.

Interpretar que el assert final es opcional, y quitarlo "para que corra más rápido". Qué pasa: alguien, adaptando este script para su propio uso, borra las tres líneas de assert al final, pensando que son solo una formalidad. Por qué pasa: los assert no cambian el resultado visible en pantalla si todo sale bien —el mensaje de éxito aparece de todas formas—, así que pueden sentirse redundantes. Cómo detectarlo: sin los assert, un bug futuro (por ejemplo, si alguien modifica RAW_ORDERS sin darse cuenta) produciría una salida silenciosamente incorrecta, sin ninguna señal de alarma. Cómo corregirlo: mantén los assert — son, con precisión, la diferencia entre "el script corrió sin errores" y "el script confirmó que el resultado es el correcto", la misma distinción que motivó toda la lección 7.

Ejercicios

Ejercicio 1 — Corre el proyecto completo tú mismo, desde cero. En un directorio nuevo, con solo raw_orders.py y kiosko_first_iceberg_table.py, corre python3 kiosko_first_iceberg_table.py. Confirma que ves los cinco 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 (lección 4), la salida debería reproducir exactamente la estructura de esta lección: cinco pasos numerados, seguidos de la verificación final con 106.15 de total y 40 == 40 en el grano. Tu snapshot_id en el Paso 5/5 va a ser un entero distinto al de cualquier otra corrida — eso es esperado, no un error.

Ejercicio 2 — Rompe el assert a propósito, y observa el fallo. Cambia temporalmente la línea assert total_revenue == 106.15, ... a assert total_revenue == 999.99, ..., corre el script de nuevo, y observa qué pasa. Después revierte el cambio.

Ver solución

El script debería fallar con un AssertionError: esperaba 106.15, obtuve 106.15 (el mensaje incluye el valor real, 106.15, para que quede claro qué esperaba el assert y qué encontró de verdad) — la ejecución se detiene ahí, sin llegar a imprimir el mensaje final de éxito. Este ejercicio demuestra, con evidencia directa, que los assert de este script cumplen una función real: si el número de negocio no coincide con lo esperado, el script te avisa de forma ruidosa e inmediata, en vez de terminar en silencio con un resultado incorrecto sin que nadie lo note.

Ejercicio 3 — Explica, en tus propias palabras, por qué este proyecto no "aprende nada nuevo" y aun así vale la pena hacerlo. En 3-4 frases, justifica por qué integrar código que ya escribiste en las lecciones 4 a 7, sin agregar ningún concepto nuevo, es un ejercicio que vale la pena — y no solo una repetición.

Ver solución

Escribir cada pieza por separado (catálogo, namespace, tabla, carga, verificación) en lecciones distintas es la forma correcta de aprender cada concepto sin sobrecarga cognitiva — pero un pipeline real nunca corre como cinco scripts separados que alguien ejecuta a mano, uno por uno, en el orden correcto. Integrar las cinco piezas en un solo script con main() y assert finales es exactamente el tipo de trabajo que un data engineer real hace después de prototipar cada pieza por separado: convertir "sé cómo hacer cada paso" en "tengo un artefacto que reproduce el resultado completo, de forma confiable, cada vez que lo corro". Es la misma transición que ya viste en python-for-data-engineering-guide, cuando kiosko_pipeline empaquetó funciones ya probadas en un solo comando ejecutable.

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

Con este proyecto cierras el módulo 1. Integraste las cinco piezas de las lecciones 4 a 7 —instalación, catálogo, namespace y tabla, carga, verificación— en un solo script, corrido de punta a punta, con assert automáticos que confirman 40 filas, 106.15 de revenue total, y cero duplicados en el grano. Kiosko tiene, por primera vez en este ecosistema, una tabla Iceberg real: kiosko.fact_orders, con un snapshot, un esquema declarado, y el mismo dato exacto que ya confirmaron seis motores distintos antes.

Diste, en ocho lecciones, el primer paso de esta guía: de un archivo Parquet suelto, sin memoria de sí mismo, a una tabla con catálogo, esquema y snapshot. Ninguno de los cuatro problemas de la lección 2 está resuelto todavía a fondo —eso es exactamente correcto en este punto de la guía—. Lo que tienes ahora es la base real sobre la que cada módulo siguiente construye una garantía específica.

Hacia dónde sigues. El módulo 2 —Anatomía de una tabla Iceberg— abre el directorio kiosko_warehouse/ que este módulo creó y explica, archivo por archivo, la cadena completa: catálogo → metadata → manifest list → manifest files → archivos de datos. Vas a inspeccionar, tanto en disco como con table.inspect.snapshots()/table.inspect.manifests()/table.inspect.files() de PyIceberg, exactamente qué apunta a qué, y por qué nada de lo que viste en este módulo se sobrescribe nunca.

Recursos

  • PyIceberg — documentación oficial (quickstart), el flujo completo de instalación, catálogo, namespace, tabla y carga que integra este proyecto. py.iceberg.apache.org. En inglés.
  • PyIceberg — referencia de API, load_catalog(), create_namespace(), create_table(), table.append(), table.scan(). py.iceberg.apache.org/api. En inglés.
  • Apache Iceberg — documentación oficial, versión de referencia 1.11.0, conceptos de tabla como temas documentados de primer nivel. iceberg.apache.org/docs/latest. En inglés.
  • DISEÑO de esta guía — el mapa completo de los ocho módulos, incluido el módulo 2 que sigue. src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.