Módulo 5: Hidden Partitioning And Partition Evolution
Proyecto: el `fact_orders_at_scale` de Kiosko, particionado
Descripción
Este proyecto cierra el módulo 5. Viste el archivero visible de Spark, con sus carpetas rotuladas a mano (lección 2). Hiciste la misma pregunta de negocio sin conocer ningún layout físico (lección 3). Conociste los tres transforms de partición y cuándo usar cada uno (lección 4). Creaste kiosko.fact_orders_at_scale con un PartitionSpec real, cargaste diez millones de filas, y confirmaste la consulta oculta con evidencia de poda (lección 5). Evolucionaste ese spec sin reescribir nada (lección 6). Y viste, con table.inspect.partitions(), cómo ambos esquemas conviven en la misma tabla (lección 7). Falta un solo paso: juntar las siete piezas en un solo script, corrido de punta a punta, con assert automáticos que confirman cada afirmación.
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 promesa que abrió este módulo en la lección 1: una tabla particionada a escala, consultada sin conocer su layout físico, evolucionada hacia adelante sin reescribir un solo archivo de los diez millones de filas ya existentes.
Una analogía: el cartero, de punta a punta, en un solo turno
Las lecciones 2 a 7 de este módulo construyeron, una pieza a la vez, la evidencia completa de la partición oculta y su evolución: el contraste con el archivero visible de Spark (lección 2), la consulta que no necesita conocerlo (lección 3), el vocabulario de los tres transforms (lección 4), la tabla real a escala con su poda medible (lección 5), la evolución sin reescritura (lección 6), y la convivencia de ambos esquemas (lección 7). Este proyecto es el turno completo del cartero, de principio a fin, sin cortes: recibe diez millones de piezas, las organiza, responde una consulta, cambia de criterio a mitad de camino, recibe más correspondencia, y te muestra, al final del turno, el archivo completo — viejo y nuevo, sin contradicción.
El material: todo lo que este módulo construyó, en un solo lugar
Necesitas, en un directorio de trabajo nuevo:
kiosko_fact_orders_at_scale_project/
├── kiosko_scale.py (el generador determinista, identico al del modulo)
└── kiosko_partitioned_at_scale.py (este proyecto: junta las 7 piezas)
Con PyIceberg instalado en tu entorno (pip install "pyiceberg[sql-sqlite,pyarrow]", módulo 1, lección 4). Este proyecto es autocontenido: crea kiosko.fact_orders_at_scale desde cero, así que no depende de ningún archivo de las lecciones anteriores de este módulo — solo de que el catálogo kiosko exista en el directorio donde lo corras (si ya tienes otras tablas de los módulos 1 a 4 en ese mismo catálogo, este proyecto las deja intactas).
Nota de escala. Este proyecto carga las 10,000,000 filas completas —las mismas 250,000 franquicias de la lección 5, no una muestra reducida—; en una laptop moderna, la generación y carga completas toman aproximadamente medio minuto. Si tu máquina es más limitada, NUM_FRANCHISES es la única constante que necesitas reducir — la fórmula (NUM_FRANCHISES × 40 filas, 106.15 × NUM_FRANCHISES de revenue) se sostiene para cualquier valor, exactamente como demostró el ejercicio 1 de la lección de spark-and-distributed-processing-guide que originó este dataset.
La solución de referencia, verificada
# kiosko_scale.py -- el generador determinista, identico al del modulo 5
from typing import Iterator, Dict, Any
KIOSKO_WEEK = [
{"order_id": "ORD-1001", "store_id": "S01", "product_id": "P001", "quantity": 3, "unit_price": 0.55, "order_ts": "2026-08-03T08:14:00"},
# ... las 40 filas completas de la semana real de Kiosko
]
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"],
}
# kiosko_partitioned_at_scale.py -- proyecto de cierre del modulo 5
# fact_orders_at_scale particionada, consultada de forma oculta, evolucionada sin reescribir
import os
from datetime import datetime
import pyarrow as pa
import pyarrow.compute as pc
from pyiceberg.catalog import load_catalog
from pyiceberg.schema import Schema
from pyiceberg.types import DoubleType, IntegerType, NestedField, StringType, TimestampType
from pyiceberg.partitioning import PartitionField, PartitionSpec
from pyiceberg.transforms import DayTransform, IdentityTransform
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)
def data_file_count(table) -> int:
return len(table.inspect.files().to_pylist())
def main() -> None:
print("=== Kiosko: fact_orders_at_scale particionada, oculta, evolucionada ===\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")
table = catalog.create_table(
"kiosko.fact_orders_at_scale", schema=FACT_ORDERS_AT_SCALE_SCHEMA, partition_spec=INITIAL_SPEC,
)
print(f"Paso 1/8 -- catalogo listo, kiosko.fact_orders_at_scale creada con spec inicial:")
print(f" {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)
files_before_evolution = data_file_count(table)
print(f"Paso 2/8 -- {row_count} filas cargadas ({NUM_FRANCHISES} franquicias x 40), "
f"revenue total {round(total_revenue, 2)}, {files_before_evolution} archivo(s) de datos")
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/8 -- 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")
files_after_evolution = data_file_count(table)
print(f"Paso 4/8 -- spec evolucionado con add_field('order_ts', DayTransform(), 'order_day'). "
f"Archivos: {files_before_evolution} -> {files_after_evolution} (sin cambios: "
f"{files_before_evolution == files_after_evolution})")
print(f" spec vigente: {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/8 -- 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
}
print(f"Paso 6/8 -- table.inspect.partitions() tiene {len(partitions)} filas de particion, "
f"spec_id presentes: {spec_ids_present}, filas por spec: {rows_by_spec}")
row_count_final = table.scan().to_arrow().num_rows
revenue_final = revenue_of(table.scan().to_arrow())
print(f"Paso 7/8 -- estado final: {row_count_final} filas, revenue total {round(revenue_final, 2)}")
total_snapshots = len(table.history())
print(f"Paso 8/8 -- table.history() tiene {total_snapshots} snapshots en total\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 == 1 and files_touched_all == 3, \
"la consulta filtrada debe tocar menos archivos que el scan completo"
assert files_before_evolution == files_after_evolution, \
"evolucionar el spec no debe reescribir ningun archivo de datos"
assert spec_ids_present == [0, 1], "ambos esquemas de particion deben convivir en la misma tabla"
assert rows_by_spec[0] == 10_000_000, "las filas del bulk load deben seguir bajo el spec original"
assert rows_by_spec[1] == 40, "solo la franquicia nueva debe estar bajo el spec evolucionado"
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 cargadas bajo un PartitionSpec inicial (IdentityTransform sobre store_id)")
print(" - store_id == 'S01' devuelve 9,575,000.00 sin que el codigo mencionara ninguna carpeta")
print(" - el scan filtrado toco 1 de 3 archivos -- la poda real de la particion oculta")
print(" - update_spec().add_field(DayTransform) no reescribio un solo archivo de datos existente")
print(" - inspect.partitions() muestra los dos esquemas de particion conviviendo en la misma tabla")
if __name__ == "__main__":
main()
Qué esperar (verificado corriendo python3 kiosko_partitioned_at_scale.py real, de punta a punta, en un directorio nuevo; ningún snapshot-id se imprime como literal, siguiendo la regla dura de esta guía):
=== Kiosko: fact_orders_at_scale particionada, oculta, evolucionada ===
Paso 1/8 -- catalogo listo, kiosko.fact_orders_at_scale creada con spec inicial:
[
1000: store_id: identity(3)
]
Paso 2/8 -- 10000000 filas cargadas (250000 franquicias x 40), revenue total 26537500.0, 3 archivo(s) de datos
Paso 3/8 -- consulta oculta store_id == 'S01': 4000000 filas, revenue 9575000.0, toco 1/3 archivos
Paso 4/8 -- spec evolucionado con add_field('order_ts', DayTransform(), 'order_day'). Archivos: 3 -> 3 (sin cambios: True)
spec vigente: [
1000: store_id: identity(3)
1001: order_day: day(7)
]
Paso 5/8 -- franquicia 250000 (40 filas) aterriza bajo el spec nuevo (con order_day)
Paso 6/8 -- table.inspect.partitions() tiene 23 filas de particion, spec_id presentes: [0, 1], filas por spec: {0: 10000000, 1: 40}
Paso 7/8 -- estado final: 10000040 filas, revenue total 26537606.15
Paso 8/8 -- table.history() tiene 2 snapshots en total
=== Verificacion final ===
Todas las verificaciones pasaron:
- 10,000,000 filas cargadas bajo un PartitionSpec inicial (IdentityTransform sobre store_id)
- store_id == 'S01' devuelve 9,575,000.00 sin que el codigo mencionara ninguna carpeta
- el scan filtrado toco 1 de 3 archivos -- la poda real de la particion oculta
- update_spec().add_field(DayTransform) no reescribio un solo archivo de datos existente
- inspect.partitions() muestra los dos esquemas de particion conviviendo en la misma tabla
Fíjate en el paso 8: table.history() reporta dos entradas, no ocho ni cinco, a pesar de que este script pasó por ocho etapas narrativas distintas. Las dos entradas son, en orden: el append() de las diez millones de filas (paso 2), y el append() de la franquicia nueva (paso 5). Ni la evolución del spec (paso 4) ni las dos consultas de lectura (paso 3) agregaron ninguna entrada al historial de snapshots — exactamente el mismo patrón que ya viste en el proyecto del módulo 4: las operaciones de metadata (esquema o partición) nunca crean snapshots; solo las escrituras de datos lo hacen.
Diagrama: de dónde venías, a dónde llegaste
flowchart LR
A["Modulo 1-4:\nfact_orders, dim_product,\ntime travel, dim_store evolucionada"] --> B["Leccion 2:\narchivero visible de Spark,\nen disco, con codigo real"]
B --> C["Leccion 3:\nmisma consulta,\nsin conocer el layout"]
C --> D["Leccion 4:\nIdentityTransform,\nBucketTransform, DayTransform"]
D --> E["Leccion 5:\n10M filas, S01 = 9,575,000.00,\n1 de 3 archivos tocados"]
E --> F["Leccion 6:\nupdate_spec() agrega order_day,\n0 archivos reescritos"]
F --> G["Leccion 7:\ninspect.partitions():\nspec_id 0 y 1 conviviendo"]
G --> H["Este proyecto:\nlas 7 piezas, un solo script,\nassert automatico"]
H --> I["Modulo 6:\nMERGE INTO y\nupserts nativos"]
Cerrando la promesa del módulo, punto por punto
| Lo que la lección 1 prometió | Evidencia de que este módulo lo entregó |
|---|---|
| El costo del particionado por carpetas de Spark, con evidencia real | Lección 2: layout Hive real en disco, store_id desaparece de una lectura que no lo declara |
| La misma consulta, sin conocer el layout | Lección 3: row_filter="store_id == 'S01'" idéntico contra una tabla sin partición |
| Vocabulario preciso de los tres transforms | Lección 4: IdentityTransform, BucketTransform, DayTransform, ejecutados y contrastados |
kiosko.fact_orders_at_scale particionada, con S01 = 9,575,000.00 | Lección 5 y este proyecto: 10,000,000 filas, poda real de 1 de 3 archivos |
| Evolución de spec sin reescribir | Lección 6 y este proyecto: files_before_evolution == files_after_evolution, verificado con assert |
| Ambos esquemas de partición conviviendo | Lección 7 y este proyecto: spec_ids_present == [0, 1], rows_by_spec == {0: 10_000_000, 1: 40} |
Este proyecto no tocó kiosko.fact_orders, kiosko.dim_product ni kiosko.dim_store — esas tablas siguen exactamente como las dejaron los módulos 1 a 4. Lo que este proyecto entrega es exactamente lo que prometió: una tabla particionada a escala completa, consultada sin conocer su layout físico, evolucionada hacia adelante sin reescribir ni una de las diez millones de filas ya escritas, con la garantía verificada con código real, no solo citada de la documentación.
Errores comunes
Correr este proyecto sobre un catálogo que ya tiene kiosko.fact_orders_at_scale de una lección anterior de este módulo. Qué pasa: alguien corre este proyecto en el mismo directorio donde ya completó las lecciones 5 a 7, y catalog.create_table(...) falla porque la tabla ya está registrada. 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_partitioned_at_scale.py, ya tienes un catálogo con kiosko.fact_orders_at_scale registrada en el mismo directorio. Cómo corregirlo: corre este proyecto en un directorio de trabajo nuevo, separado de donde hiciste las lecciones 5 a 7 — tal como sugiere "El material" de esta lección.
Esperar que el paso 6 imprima record_count sumando exactamente 10,000,040 en una sola fila. Qué pasa: alguien, al ver rows_by_spec: {0: 10000000, 1: 40}, espera encontrar también una fila consolidada con el total combinado (10,000,040) en algún lugar de la salida de table.inspect.partitions(). Por qué pasa: es natural buscar un "total general" en cualquier reporte tabular, como el TOTAL al pie de una factura. Cómo detectarlo: si tu código busca una fila con spec_id vacío o None que represente el agregado, revisa el paso 6 de este proyecto — rows_by_spec se calcula sumando en Python, después de leer el resultado de inspect.partitions(), no porque el método lo entregue así. Cómo corregirlo: table.inspect.partitions() siempre devuelve una fila por cada combinación única de spec_id y valores de partición — nunca un total consolidado; cualquier agregación adicional que necesites (como rows_by_spec en este proyecto) es responsabilidad de tu propio código, sobre el resultado ya obtenido.
Ejercicios
Ejercicio 1 — Corre el proyecto completo tú mismo, desde cero. En un directorio nuevo, corre python3 kiosko_partitioned_at_scale.py. Confirma que ves los ocho pasos completarse y el mensaje final con las cinco verificaciones.
Ver solución
Si PyIceberg está instalado en tu entorno, la salida debería reproducir exactamente la estructura de esta lección: ocho pasos numerados, seguidos de la verificación final con los cinco mensajes de éxito. La generación de las diez millones de filas es la parte más lenta del script —espera algo cercano a medio minuto en una laptop moderna—; el resto de los pasos son casi instantáneos.
Ejercicio 2 — Rompe un assert a propósito, y observa el fallo. Cambia temporalmente NUM_FRANCHISES de 250_000 a 100_000, corre el script de nuevo, y observa qué assert falla primero. Después revierte el cambio.
Ver solución
El primer assert que falla es assert row_count == NUM_FRANCHISES * 40 == 10_000_000 — porque con 100_000 franquicias, row_count es 4_000_000, que ya no coincide con el literal 10_000_000 de la segunda mitad de esa comparación encadenada. Este ejercicio demuestra que los assert de este proyecto no solo verifican consistencia interna (row_count == NUM_FRANCHISES * 40, que se sostendría para cualquier valor) — también verifican, con un literal explícito, que el dataset completo de esta guía es específicamente el de 250,000 franquicias, ni más ni menos. Si necesitaras correr este proyecto con un NUM_FRANCHISES distinto por limitaciones de tu máquina, tendrías que ajustar también los literales de los assert, exactamente como sugiere la nota de escala al principio de esta lección.
Ejercicio 3 — Explica, en tus propias palabras, por qué este proyecto verifica files_touched_s01 == 1 en vez de solo verificar s01_scan.num_rows == 4_000_000. En 3-4 frases, justifica por qué el assert sobre archivos tocados es tan importante como el assert sobre el conteo de filas.
Ver solución
Verificar solo el conteo de filas —que la consulta filtrada devuelva las 4,000,000 filas correctas de S01— confirmaría que el resultado de negocio es correcto, pero no confirmaría por qué es eficiente obtenerlo. Una implementación distinta, que escaneara los tres archivos completos y descartara las filas que no son de S01 después de leerlas todas (el comportamiento que este módulo completo argumenta que la partición evita), podría llegar exactamente al mismo resultado final, sin ninguno de los beneficios de poda que este módulo se propuso demostrar. Verificar files_touched_s01 == 1 (contra files_touched_all == 3) confirma la afirmación central del módulo —que particionar por store_id le permite al motor ignorar archivos completos sin abrirlos—, con evidencia directa sobre cuántos archivos tocó el plan de ejecución, no solo sobre el resultado final de la consulta. Es la misma disciplina de "verificar el camino, no solo el destino" que ya viste en los proyectos de los módulos 3 y 4.
Resumen y siguiente paso: el cierre de este módulo
Con este proyecto cierras el módulo 5. Integraste las siete lecciones anteriores —el contraste con el archivero visible de Spark, la consulta oculta, los tres transforms de partición, la tabla a escala con poda real, la evolución sin reescritura, y la convivencia de ambos esquemas— en un solo script, corrido de punta a punta, con assert automáticos que confirman cada afirmación con evidencia, no con una promesa.
Kiosko tiene, por primera vez en este ecosistema, una tabla particionada a escala real cuyo esquema de partición evolucionó después de que los datos ya existían, sin que esa evolución tocara un solo archivo Parquet ya escrito — y una consulta de negocio que nunca tuvo que mencionar, en ningún momento, cómo esos archivos están organizados por dentro.
Hacia dónde sigues. El módulo 6 —MERGE INTO y upserts nativos— pone lado a lado las tres formas en que Kiosko ya resolvió "actualizar P002 sin perder su historia" —MERGE INTO a mano en DuckDB (data-modeling), dbt snapshot automatizado (dbt)— y agrega la cuarta: MERGE INTO nativo de Iceberg vía Spark SQL, y table.upsert() de PyIceberg como alternativa cien por ciento Python. Es, además, el único módulo de esta guía que necesita la JVM — reutilizando el PySpark que spark-and-distributed-processing-guide ya dejó instalado, declarado así de explícito desde su primera lección.
Recursos
- PyIceberg — documentación oficial (quickstart), el flujo completo de catálogo, tabla,
PartitionSpec,append()yupdate_spec()que integra este proyecto. py.iceberg.apache.org. En inglés. - PyIceberg — referencia de API,
PartitionSpec/PartitionField, los transforms,table.update_spec(),table.scan(...).plan_files(),table.inspect.partitions(). py.iceberg.apache.org/api. En inglés. - Apache Iceberg — documentación oficial, "Partitioning" (partición oculta, transforms, evolución de partición), la base formal de todo lo que este proyecto verifica con
assert. iceberg.apache.org/docs/latest/partitioning. En inglés. - DISEÑO de
spark-and-distributed-processing-guide— fuente degenerate_orders_at_scale()y los números exactos del dataset a escala que este proyecto reconstruye y carga 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, incluido el módulo 6 que sigue.
src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.