Módulo 1: From File Format To Table Format
Cargando el fact_orders de Kiosko en Iceberg
Descripción
Esta es la lección central del módulo: la tabla kiosko.fact_orders, vacía desde la lección 5, recibe por fin las cuarenta filas reales de la semana de Kiosko. Vas a reconstruir fact_orders.parquet con pyarrow —el mismo Parquet que spark-and-distributed-processing-guide (módulo 3 de esa guía) ya dejó escrito, reconstruido aquí para que esta guía sea autocontenida— y vas a cargarlo con table.append(), la operación que crea, por primera vez, un snapshot real.
Conexión con el módulo. Esta lección junta todo lo que instalaste en las lecciones 4 y 5 —el catálogo y la tabla vacía— con los datos reales de Kiosko que ya conoces de las seis guías anteriores. La lección 7 verifica, con el mismo número de siempre (106.15), que la migración fue exacta.
Una analogía: pegar la primera colección de fotos en el álbum
El álbum de la lección 5 tiene portada e índice, pero está vacío. Esta lección es el momento de pegar, por primera vez, una colección completa de fotos —las cuarenta órdenes de la semana de Kiosko— en sus páginas. Y fíjate en algo que la analogía predice con precisión: cuando el bibliotecario termina de pegar esta primera colección, no solo actualiza el índice para decir "ahora hay cuarenta fotos" — también archiva, en un registro aparte, el hecho de que esta fue la primera colección jamás pegada en este álbum, con su propia fecha y su propio número de referencia. Ese registro es, con precisión, lo que Iceberg llama un snapshot: no es una copia de las fotos —esas siguen siendo los mismos archivos Parquet de siempre—, es la constancia archivada de "así se veía el álbum completo, justo después de este evento de carga".
Ejemplo trabajado: de datos fijos en Python a una tabla Iceberg real
Paso 1 — La semana fija de Kiosko, idéntica a las seis guías anteriores
# raw_orders.py -- la semana fija de Kiosko, identica a foundations/data-modeling/dbt/spark
RAW_ORDERS = [
# Lunes 2026-08-03 (8 ordenes)
("ORD-1001", "S01", "P001", 3, 0.55, "2026-08-03T08:14:00"),
("ORD-1002", "S01", "P002", 1, 1.20, "2026-08-03T08:20:00"),
("ORD-1003", "S02", "P003", 2, 0.75, "2026-08-03T08:31:00"),
("ORD-1004", "S01", "P004", 1, 4.50, "2026-08-03T09:02:00"),
("ORD-1005", "S03", "P001", 5, 0.55, "2026-08-03T09:15:00"),
("ORD-1006", "S02", "P002", 2, 1.20, "2026-08-03T09:47:00"),
("ORD-1007", "S03", "P003", 1, 0.75, "2026-08-03T10:05:00"),
("ORD-1008", "S01", "P001", 2, 0.55, "2026-08-03T10:22:00"),
# Martes 2026-08-04 (6 ordenes)
("ORD-2001", "S01", "P002", 1, 1.20, "2026-08-04T08:05:00"),
("ORD-2002", "S02", "P001", 4, 0.55, "2026-08-04T08:40:00"),
("ORD-2003", "S03", "P004", 1, 4.50, "2026-08-04T09:12:00"),
("ORD-2004", "S01", "P003", 3, 0.75, "2026-08-04T09:50:00"),
("ORD-2005", "S02", "P002", 2, 1.20, "2026-08-04T10:15:00"),
("ORD-2006", "S03", "P001", 6, 0.55, "2026-08-04T10:33:00"),
# Miercoles 2026-08-05 (2 ordenes)
("ORD-3001", "S02", "P004", 2, 4.50, "2026-08-05T08:10:00"),
("ORD-3002", "S01", "P001", 1, 0.55, "2026-08-05T08:22:00"),
# Jueves 2026-08-06 (5 ordenes)
("ORD-4001", "S01", "P001", 4, 0.55, "2026-08-06T08:10:00"),
("ORD-4002", "S02", "P003", 2, 0.75, "2026-08-06T08:45:00"),
("ORD-4003", "S03", "P002", 1, 1.20, "2026-08-06T09:20:00"),
("ORD-4004", "S01", "P004", 1, 4.50, "2026-08-06T09:55:00"),
("ORD-4005", "S02", "P001", 3, 0.55, "2026-08-06T10:30:00"),
# Viernes 2026-08-07 (7 ordenes)
("ORD-5001", "S01", "P002", 2, 1.20, "2026-08-07T08:05:00"),
("ORD-5002", "S03", "P001", 4, 0.55, "2026-08-07T08:30:00"),
("ORD-5003", "S02", "P004", 1, 4.50, "2026-08-07T08:58:00"),
("ORD-5004", "S01", "P003", 2, 0.75, "2026-08-07T09:22:00"),
("ORD-5005", "S03", "P002", 3, 1.20, "2026-08-07T09:47:00"),
("ORD-5006", "S02", "P001", 5, 0.55, "2026-08-07T10:15:00"),
("ORD-5007", "S01", "P001", 2, 0.55, "2026-08-07T10:40:00"),
# Sabado 2026-08-08 (9 ordenes)
("ORD-6001", "S01", "P001", 6, 0.55, "2026-08-08T08:00:00"),
("ORD-6002", "S02", "P002", 3, 1.20, "2026-08-08T08:18:00"),
("ORD-6003", "S03", "P001", 4, 0.55, "2026-08-08T08:35:00"),
("ORD-6004", "S01", "P004", 2, 4.50, "2026-08-08T08:52:00"),
("ORD-6005", "S02", "P003", 3, 0.75, "2026-08-08T09:10:00"),
("ORD-6006", "S03", "P002", 2, 1.20, "2026-08-08T09:28:00"),
("ORD-6007", "S01", "P003", 1, 0.75, "2026-08-08T09:45:00"),
("ORD-6008", "S02", "P001", 7, 0.55, "2026-08-08T10:02:00"),
("ORD-6009", "S03", "P004", 1, 4.50, "2026-08-08T10:20:00"),
# Domingo 2026-08-09 (3 ordenes)
("ORD-7001", "S01", "P001", 2, 0.55, "2026-08-09T09:15:00"),
("ORD-7002", "S02", "P002", 1, 1.20, "2026-08-09T09:40:00"),
("ORD-7003", "S03", "P001", 3, 0.55, "2026-08-09T10:05:00"),
]
Paso 2 — Reconstruye fact_orders.parquet con pyarrow
# build_fact_orders_parquet.py
from datetime import datetime
import pyarrow as pa
import pyarrow.parquet as pq
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_ids = {s["store_id"] for s in DIM_STORE}
product_ids = {p["product_id"] for p in DIM_PRODUCT}
order_id_col, store_id_col, product_id_col = [], [], []
quantity_col, unit_price_col, revenue_col, order_ts_col = [], [], [], []
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}")
order_id_col.append(order_id)
store_id_col.append(store_id)
product_id_col.append(product_id)
quantity_col.append(quantity)
unit_price_col.append(unit_price)
revenue_col.append(round(quantity * unit_price, 10))
order_ts_col.append(datetime.fromisoformat(ts))
fact_orders_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(order_id_col, type=pa.string()),
pa.array(store_id_col, type=pa.string()),
pa.array(product_id_col, type=pa.string()),
pa.array(quantity_col, type=pa.int32()),
pa.array(unit_price_col, type=pa.float64()),
pa.array(revenue_col, type=pa.float64()),
pa.array(order_ts_col, type=pa.timestamp("us")),
],
schema=fact_orders_schema,
)
pq.write_table(pa_table, "fact_orders.parquet")
print(f"fact_orders.parquet escrito: {pa_table.num_rows} filas, {pa_table.num_columns} columnas")
Qué esperar (verificado corriendo el script real):
fact_orders.parquet escrito: 40 filas, 7 columnas
Paso 3 — Carga el Parquet dentro de la tabla Iceberg con table.append()
# load_into_iceberg.py
import os
import pyarrow.parquet as pq
from pyiceberg.catalog import load_catalog
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.load_table("kiosko.fact_orders")
pa_table = pq.read_table("fact_orders.parquet")
print(f"fact_orders.parquet leido: {pa_table.num_rows} filas")
table.append(pa_table)
# el snapshot-id lo asigna Iceberg en el momento del commit -- se captura
# en variable, nunca se hardcodea (regla dura de esta guia)
snap_id = table.current_snapshot().snapshot_id
print(f"\nPrimer snapshot creado, snapshot_id capturado en variable: {snap_id}")
scanned = table.scan().to_arrow()
print(f"\nlen(table.scan().to_arrow()) = {scanned.num_rows}")
Qué esperar (verificado corriendo el script real; el snapshot_id es un entero grande, asignado por Iceberg en el momento exacto del commit — distinto en cada corrida tuya, nunca el mismo dos veces, así que se muestra aquí como marcador en vez de un número fijo):
fact_orders.parquet leido: 40 filas
Primer snapshot creado, snapshot_id capturado en variable: <snapshot-id asignado en tu corrida, distinto cada vez>
len(table.scan().to_arrow()) = 40
Cuarenta filas leídas del Parquet, cuarenta filas confirmadas dentro de la tabla Iceberg después de table.append() — el mismo número, sin ninguna pérdida ni duplicación. Y por primera vez en esta guía, table.current_snapshot() deja de ser None: existe un snapshot real, con un snapshot_id que tu corrida asignó, distinto del que asignaría cualquier otra corrida —incluida una segunda corrida tuya, si borraras la tabla y la recrearas—. Este es exactamente el motivo por el que la regla de esta guía prohíbe hardcodear un snapshot-id: no es un número reproducible, es un identificador que el propio commit genera, y el código de arriba lo captura en la variable snap_id inmediatamente después de crearse, en vez de asumir cuál va a ser.
Diagrama: qué apareció en disco después del append()
flowchart TB
A["table.append(pa_table)"] --> B["1 archivo de datos nuevo\ndata/00000-0-<uuid>.parquet"]
A --> C["1 manifest file nuevo\nmetadata/<uuid>-m0.avro\n(lista ESE archivo de datos)"]
A --> D["1 manifest list nuevo\nmetadata/snap-<snapshot_id>-0-<uuid>.avro\n(lista ESE manifest file)"]
A --> E["1 metadata.json nuevo\nmetadata/00001-<uuid>.metadata.json\n(apunta al snapshot nuevo)"]
E -->|"el catalogo ahora apunta aqui"| F["kiosko_catalog.db:\nmetadata_location actualizado"]
Después de esta lección, kiosko_warehouse/kiosko/fact_orders/ contiene, verificado en disco:
kiosko_warehouse/kiosko/fact_orders/
├── data/
│ └── 00000-0-<uuid>.parquet (las 40 filas, Parquet normal)
└── metadata/
├── 00000-<uuid>.metadata.json (leccion 5: tabla vacia, snapshot=None)
├── 00001-<uuid>.metadata.json (esta leccion: apunta al primer snapshot)
├── <uuid>-m0.avro (manifest file: lista el archivo de datos)
└── snap-<snapshot_id>-0-<uuid>.avro (manifest list: lista el manifest file)
No necesitas entender todavía, en detalle, qué es un manifest file o un manifest list — el módulo 2 completo de esta guía está dedicado a esa cadena exacta (catálogo → metadata → manifest list → manifest files → archivos de datos), abriendo cada uno de estos archivos con la API de PyIceberg. Lo que vale la pena confirmar aquí es algo más simple y más importante: hay dos archivos .metadata.json, no uno solo. El primero (00000-...) es el que la lección 5 creó, con la tabla vacía. El segundo (00001-...) es el que esta lección acaba de crear — y el primero sigue existiendo, sin haber sido borrado ni sobrescrito. Esa es, en la práctica más concreta posible, la garantía de "cada escritura es un snapshot nuevo, nada se sobrescribe" que vas a explorar a fondo en el módulo 3.
Errores comunes
Intentar table.append() con un pyarrow.Table cuyo esquema no coincide exactamente. Qué pasa: alguien construye su propio pa.Table con, por ejemplo, quantity como int64 en vez de int32, o con las columnas en un orden distinto al del esquema de Iceberg, y table.append() falla con un error de validación de esquema. Por qué pasa: pyarrow es flexible sobre tipos numéricos por defecto (int64 es el tipo que usa si no se especifica nada), y es fácil no fijarse en que el esquema declarado en la lección 5 pidió int32 para quantity. Cómo detectarlo: si table.append() falla mencionando una discrepancia de tipos o de esquema, compara, columna por columna, el pa_table.schema que estás pasando contra el table.schema() de la tabla Iceberg. Cómo corregirlo: el build_fact_orders_parquet.py de esta lección declara el pa.schema() explícitamente, columna por columna, con los mismos tipos exactos del Schema de Iceberg de la lección 5 — replicar esa disciplina evita este error por completo.
Correr load_into_iceberg.py dos veces, y sorprenderse al ver 80 filas en vez de 40. Qué pasa: alguien corre el script del paso 3 una vez, ve 40 filas, lo corre de nuevo por curiosidad (o por error), y ahora table.scan().to_arrow().num_rows da 80. Por qué pasa: table.append() es, exactamente como su nombre lo dice, una operación que agrega filas — no reemplaza el contenido de la tabla. Cada llamada crea un snapshot nuevo con las filas nuevas sumadas a las que ya había, el mismo comportamiento (por diseño) que un INSERT normal en SQL. Cómo detectarlo: si el conteo de filas de tu tabla es un múltiplo de 40, revisa cuántas veces corriste table.append() sobre la misma tabla. Cómo corregirlo: si necesitas volver a empezar desde cero, borra la tabla completa con catalog.drop_table("kiosko.fact_orders") y recréala desde la lección 5 — o, más adelante en esta guía (módulo 3), usa table.overwrite() en vez de table.append() cuando la intención sea reemplazar el contenido, no sumarle.
Ejercicios
Ejercicio 1 — Reproduce la carga completa tú mismo. En tu propia máquina, con la tabla kiosko.fact_orders vacía de la lección 5 todavía disponible, corre los tres pasos de esta lección en orden: reconstruye fact_orders.parquet, cárgalo con table.append(), y confirma len(table.scan().to_arrow()) == 40. Anota el snapshot_id que tu propia corrida asignó.
Ver solución
Si seguiste los tres pasos exactamente, deberías ver fact_orders.parquet escrito: 40 filas, 7 columnas, seguido de fact_orders.parquet leido: 40 filas y len(table.scan().to_arrow()) = 40. Tu snapshot_id va a ser un entero grande, distinto del de cualquier otra persona que haga este mismo ejercicio —y distinto también de cualquier corrida anterior tuya sobre una tabla recreada—, exactamente como advierte esta lección: no hay un valor "correcto" que deba coincidir, la única verificación real es que el número exista (no sea None) y que el conteo de filas sea 40.
Ejercicio 2 — Verifica que los dos archivos de metadata coexisten. Usando ls kiosko_warehouse/kiosko/fact_orders/metadata/ desde tu terminal (no desde Python), confirma que existen exactamente dos archivos .metadata.json después de esta lección, y explica en 1-2 frases por qué el primero no desapareció.
Ver solución
Deberías ver dos archivos, algo como 00000-<uuid>.metadata.json y 00001-<uuid>.metadata.json (los UUIDs exactos van a ser distintos en tu corrida). El primero no desapareció porque Iceberg nunca sobrescribe un archivo de metadata existente — cada cambio de estado de la tabla (crearla vacía, agregarle datos, y cualquier operación futura) escribe un archivo de metadata nuevo, y el catálogo simplemente actualiza su puntero hacia el más reciente. Los archivos viejos siguen en disco, disponibles, hasta que una operación de mantenimiento explícita (expire_snapshots, que vas a ver en el módulo 7) decida limpiarlos.
Ejercicio 3 — Predicción: ¿qué pasaría si raw_orders.py tuviera un store_id inválido? Sin correrlo todavía, predice: si una de las 40 filas de RAW_ORDERS tuviera "S99" como store_id (un valor que no existe en DIM_STORE), ¿en qué paso exacto del ejemplo trabajado de esta lección esperas que aparezca el error?
Ver solución
El error aparecería en el paso 2 —build_fact_orders_parquet.py—, específicamente en la línea if store_id not in store_ids: raise ValueError(...), mucho antes de que el dato llegue siquiera a intentar escribirse como Parquet, y mucho antes de que Iceberg entre en juego. Esta es una validación de integridad referencial hecha a propósito, en Python puro, replicando la misma disciplina de validate_orders() que ya viste en data-engineering-foundations-guide — Iceberg garantiza atomicidad y estructura de la escritura, pero no valida por sí solo que un store_id corresponda a una tienda real de DIM_STORE; esa es, todavía, responsabilidad del código que prepara los datos antes de que lleguen a table.append().
Resumen y siguiente paso
En esta lección reconstruiste fact_orders.parquet con pyarrow —las mismas cuarenta filas de siempre, la misma semana fija del 3 al 9 de agosto de 2026— y lo cargaste dentro de kiosko.fact_orders con table.append(). Confirmaste, con table.current_snapshot() dejando de ser None, que se creó el primer snapshot real de esta tabla, y viste en disco que apareció un archivo de datos, un manifest file, un manifest list, y un segundo archivo de metadata — sin que el primero, de la tabla vacía, desapareciera.
Antes de avanzar deberías poder: explicar la diferencia entre table.append() y una operación que reemplaza contenido; reproducir la carga completa en tu propia máquina; y explicar por qué el snapshot_id nunca se hardcodea en el código de esta guía.
Tienes cuarenta filas cargadas, pero todavía no verificaste el número que de verdad importa: si el revenue total sigue siendo 106.15, exactamente como en las seis guías anteriores. Eso es, con precisión, el trabajo de la lección 7.
Recursos
- PyIceberg — referencia de API, sintaxis exacta de
table.append()ytable.current_snapshot(). py.iceberg.apache.org/api. En inglés. - Apache Iceberg — documentación oficial, "Table Spec", la definición formal de snapshot, manifest list y manifest file que el módulo 2 de esta guía inspecciona a fondo. iceberg.apache.org/spec. En inglés.
- DISEÑO de
spark-and-distributed-processing-guide— fuente defact_orders.parquet, reconstruido en esta lección conpyarrowde forma autocontenida.src/guides/spark-and-distributed-processing-guide/DISENO.md. En español. - DISEÑO de esta guía — la regla dura sobre nunca hardcodear un
snapshot-id, aplicada en esta lección por primera vez.src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.