Módulo 5: Hidden Partitioning And Partition Evolution

Particionando `fact_orders_at_scale`

Descripción

Esta es la lección donde todo lo conceptual de las lecciones 2 a 4 se vuelve una tabla real, a escala completa. Vas a crear kiosko.fact_orders_at_scale con un PartitionSpec inicial —IdentityTransform sobre store_id, el mismo criterio que Spark ya usó para sus carpetas—, vas a cargar las diez millones de filas reales de Kiosko a escala, y vas a confirmar dos cosas con evidencia ejecutada: que la consulta oculta de la lección 3 sigue dando el resultado correcto (9,575,000.00 para S01), y que esta vez sí hay algo que podar — con un número exacto de archivos tocados, no solo una promesa.

Conexión con el módulo. Esta lección junta el archivero de la lección 2, la consulta sin carpetas de la lección 3, y el vocabulario de transforms de la lección 4, en una sola tabla ejecutada de punta a punta. Las lecciones 6 y 7 van a construir sobre exactamente esta misma tabla — no la vuelven a crear desde cero.

Declaración de dataset, heredada sin regenerar el argumento. Las diez millones de filas de esta lección son sintéticas, generadas por generate_orders_at_scale(250_000) —el mismo generador determinista, sin random, que spark-and-distributed-processing-guide ya justificó en su módulo 4—. Kiosko real —tres tiendas, cuarenta órdenes en una semana— nunca produce ese volumen; esta guía reconstruye los datos, con el mismo generador, porque necesita cargarlos de verdad dentro de una tabla Iceberg para que la partición y su evolución tengan algo real que mover.

Una analogía: el archivero, reconstruido por el propio cartero

Vuelve a la analogía de la lección 1: esta lección es el momento en que el cartero recibe, por primera vez, un volumen real de correspondencia —diez millones de piezas—, y organiza sus propios sacos según su propio criterio, sin que nadie tenga que indicarle nada más allá de "organiza por destinatario". A partir de aquí, cualquiera que le pida "la correspondencia de S01" va a recibir una respuesta rápida y correcta, sin haber tenido que aprender nunca cómo el cartero organiza sus sacos por dentro.

Ejemplo trabajado: crear, cargar, consultar, medir

Paso 1 — Declara el esquema y el PartitionSpec inicial

# l5_partitioning_fact_orders_at_scale.py
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 PartitionSpec, PartitionField
from pyiceberg.transforms import IdentityTransform

from kiosko_scale import generate_orders_at_scale

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}",
)

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),
)

# el mismo store_id que Spark ya uso para particionar por carpetas (leccion 2),
# ahora como PartitionSpec de Iceberg -- los field_id de particion empiezan en
# 1000 por convencion del formato (leccion 4)
INITIAL_SPEC = PartitionSpec(
    PartitionField(source_id=3, field_id=1000, transform=IdentityTransform(), name="store_id"),
)

table = catalog.create_table(
    "kiosko.fact_orders_at_scale",
    schema=FACT_ORDERS_AT_SCALE_SCHEMA,
    partition_spec=INITIAL_SPEC,
)
print("kiosko.fact_orders_at_scale creada. PartitionSpec inicial:")
print(table.spec())

Qué esperar (verificado corriendo el script real):

kiosko.fact_orders_at_scale creada. PartitionSpec inicial:
[
  1000: store_id: identity(3)
]

source_id=3 porque store_id es el tercer campo del Schema que acabas de declarar arriba (order_id=1, franchise_id=2, store_id=3); field_id=1000 porque este es el primer campo de partición de esta tabla (lección 4). El texto identity(3) que imprime PyIceberg confirma ambos números de un vistazo: el transform (identity) y la columna de origen (3, es decir, store_id).

Paso 2 — Genera y carga las diez millones de filas reales

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())


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

NUM_FRANCHISES = 250_000
rows = list(generate_orders_at_scale(NUM_FRANCHISES))
for r in rows:
    r["order_ts"] = datetime.fromisoformat(r["order_ts"])
pa_table = pa.Table.from_pylist(rows, schema=PA_SCHEMA)

table.append(pa_table)

# el snapshot-id lo asigna Iceberg en el commit -- se captura en variable,
# nunca se hardcodea (regla dura de esta guia, aplicada desde el modulo 1)
snap_after_bulk_load = table.current_snapshot().snapshot_id
row_count = table.scan().to_arrow().num_rows
print(f"table.scan().to_arrow().num_rows = {row_count}")
assert row_count == NUM_FRANCHISES * 40 == 10_000_000

total_revenue = revenue_of(pa_table)
print(f"revenue total = {round(total_revenue, 2)}")
assert round(total_revenue, 2) == 26_537_500.00

Qué esperar (verificado corriendo el script real; la generación y conversión a Arrow de las diez millones de filas toma unos 30 segundos en una laptop moderna, el append() en sí menos de un segundo — números de referencia, no una medición de rendimiento a reproducir byte a byte):

table.scan().to_arrow().num_rows = 10000000
revenue total = 26537500.0

Diez millones de filas exactas, 26,537,500.00 de revenue exacto — los mismos dos números que spark-and-distributed-processing-guide ya verificó con su propio motor. snap_after_bulk_load quedó capturado en variable, listo para cualquier uso posterior que lo necesite (esta lección no lo vuelve a usar, pero la lección 6 lo hace de forma equivalente con su propio snapshot).

Paso 3 — La consulta oculta, a escala completa

print("=== La consulta oculta: filtra por store_id, nunca por carpeta ===")
s01_scan = table.scan(row_filter="store_id == 'S01'").to_arrow()
s01_revenue = revenue_of(s01_scan)
print(f'table.scan(row_filter="store_id == \'S01\'").to_arrow()')
print(f"  filas: {s01_scan.num_rows}, revenue: {round(s01_revenue, 2)}")
assert s01_scan.num_rows == 4_000_000
assert round(s01_revenue, 2) == 9_575_000.00

Qué esperar:

=== La consulta oculta: filtra por store_id, nunca por carpeta ===
table.scan(row_filter="store_id == 'S01'").to_arrow()
  filas: 4000000, revenue: 9575000.0

9,575,000.00 — el mismo número exacto que aparece en el DISEÑO de esta guía, idéntico al que spark-and-distributed-processing-guide obtiene filtrando el mismo store_id sobre su propio Parquet particionado. La línea de código es, literalmente, la misma que ya usaste en la lección 3 contra kiosko.fact_orders sin partición — ningún cambio de sintaxis, solo un resultado que ahora sí tiene partición real por debajo.

Paso 4 — La evidencia de poda: cuántos archivos tocó cada scan

print("=== Evidencia de poda: cuantos archivos de datos toco cada scan ===")
all_files = list(table.scan().plan_files())
s01_files = list(table.scan(row_filter="store_id == 'S01'").plan_files())
print(f"  scan completo (sin filtro): {len(all_files)} archivo(s) de datos")
print(f"  scan filtrado (store_id == 'S01'): {len(s01_files)} archivo(s) de datos")
assert len(all_files) == 3
assert len(s01_files) == 1

Qué esperar:

=== Evidencia de poda: cuantos archivos de datos toco cada scan ===
  scan completo (sin filtro): 3 archivo(s) de datos
  scan filtrado (store_id == 'S01'): 1 archivo(s) de datos

Esta es la diferencia real frente a la lección 3, donde ambos números —con filtro y sin filtro— eran 1, porque no había nada que podar. Aquí, un append() de diez millones de filas sobre una tabla particionada por store_id produjo exactamente tres archivos de datos —uno por cada valor distinto de la columna de partición—, y el scan filtrado por S01 tocó exactamente uno de esos tres, ignorando por completo los otros dos sin que el código lo pidiera explícitamente. Eso es partición oculta funcionando de verdad: la misma sintaxis de la lección 3, ahora con un beneficio medible.

Diagrama: de la carga a la poda

flowchart TB
    A["generate_orders_at_scale(250_000)\n10,000,000 filas"] --> B["table.append(pa_table)"]
    B --> C["3 archivos de datos\n(uno por store_id, por el PartitionSpec)"]
    C --> D1["data/store_id=S01/...\n4,000,000 filas"]
    C --> D2["data/store_id=S02/...\n3,250,000 filas"]
    C --> D3["data/store_id=S03/...\n2,750,000 filas"]
    E["table.scan(row_filter=\"store_id == 'S01'\")"] -->|"consulta store_id,\nnunca menciona una carpeta"| C
    C -.->|"plan_files() decide,\nsolo abre D1"| D1

Errores comunes

Esperar que table.append() produzca un archivo por fila, o un solo archivo gigante, en vez de uno por valor de partición. Qué pasa: alguien, al ver que la carga produjo exactamente 3 archivos para 10,000,000 filas, se sorprende — esperaba muchos más (uno por lote de escritura interno) o exactamente uno (todo junto, sin partición real). Por qué pasa: sin haber visto antes el mecanismo completo, es fácil no anticipar que el número de archivos de un append() está gobernado, precisamente, por cuántos valores distintos de partición aparecen en los datos que estás cargando. Cómo detectarlo: si tu conteo de archivos de datos no coincide con el número de valores distintos de tu columna particionada (más alguna fragmentación adicional en cargas muy grandes o en múltiples llamadas a append()), revisa qué PartitionSpec tiene tu tabla. Cómo corregirlo: con IdentityTransform sobre una columna de tres valores, y una sola llamada a append() con todos los datos en memoria de una vez, 3 archivos es exactamente lo esperado — uno por store_id. Cargas repetidas (varios append() sucesivos) sí producen más archivos, uno por cada combinación de llamada y valor de partición — ese es, precisamente, el problema de "muchos archivos chicos" que la compactación del módulo 7 va a resolver.

Medir "más rápido" comparando el tiempo de esta consulta contra la de la lección 3, en vez de comparar archivos tocados. Qué pasa: alguien cronometra ambas consultas —contra kiosko.fact_orders sin partición, y contra kiosko.fact_orders_at_scale particionada— y trata la diferencia de tiempo como la prueba central de que la partición "funciona". Por qué pasa: un cronómetro es intuitivo y fácil de usar, y el resultado —más rápido— confirma la expectativa. Cómo detectarlo: si tu evidencia de "la partición ayuda" es un número de segundos, en vez de un número de archivos, revisa por qué esta lección mide con plan_files() en su lugar. Cómo corregirlo: un tiempo de ejecución depende de la máquina, de la carga del sistema en ese momento, del caché del sistema operativo — no es reproducible entre corridas ni entre personas. El número de archivos que plan_files() reporta sí lo es: 1 de 3 es un hecho estructural de la tabla, verificable con un assert, sin depender de ningún reloj.

Ejercicios

Ejercicio 1 — Reproduce la carga completa tú mismo, y verifica los cuatro números centrales. En tu propia máquina, con kiosko.fact_orders_at_scale recién creada, corre los cuatro pasos de esta lección. Confirma 10,000,000 filas, 26,537,500.00 de revenue total, 9,575,000.00 para S01, y 1 de 3 archivos tocados por el scan filtrado.

Ver solución

Si seguiste los cuatro pasos exactamente, tu salida debería coincidir, número por número, con la de esta lección — los cuatro assert del ejemplo trabajado están ahí precisamente para que no tengas que confiar en tu propia lectura visual del resultado. Si algún número no coincide, revisa primero si generate_orders_at_scale(250_000) corrió completo, sin interrupciones — un dataset parcial es la causa más común de un conteo distinto.

Ejercicio 2 — Repite la consulta filtrada para S02 y S03, y confirma el desglose completo. Usando la misma sintaxis del paso 3, filtra por store_id == 'S02' y store_id == 'S03', y confirma con assert los revenues 9,700,000.00 y 7,262,500.00 respectivamente — el mismo desglose por tienda que ya conoces desde el módulo 1.

Ver solución
s02_scan = table.scan(row_filter="store_id == 'S02'").to_arrow()
s03_scan = table.scan(row_filter="store_id == 'S03'").to_arrow()
assert round(revenue_of(s02_scan), 2) == 9_700_000.00
assert round(revenue_of(s03_scan), 2) == 7_262_500.00
print("Verificacion: desglose completo por tienda confirmado con assert")

Los tres números —9,575,000.00, 9,700,000.00, 7,262,500.00— suman exactamente 26,537,500.00, el revenue total del paso 2. Esta es la misma disciplina de verificación cruzada que ya usaste en guías anteriores del ecosistema: sumar las partes y confirmar que coinciden con el todo.

Ejercicio 3 — Predicción: ¿cuántos archivos tocaría un scan filtrado por product_id, en vez de store_id? Sin correrlo, predice: si escribieras table.scan(row_filter="product_id == 'P002'").plan_files() contra esta misma tabla, ¿cuántos de los tres archivos de datos esperas que toque?

Ver solución

Los tres. kiosko.fact_orders_at_scale está particionada únicamente por store_idproduct_id no forma parte de ningún PartitionField del spec actual—, así que filtrar por product_id no le da al motor ninguna pista sobre qué archivos puede ignorar: cada uno de los tres archivos (uno por tienda) contiene los cuatro productos mezclados, así que los tres tienen que abrirse para encontrar las filas de P002. Esto no es un error ni una limitación de PyIceberg — es, con precisión, la definición de qué significa "particionar por una columna": la poda solo funciona para las columnas que de verdad forman parte del PartitionSpec, nunca para cualquier columna del esquema.

Resumen y siguiente paso

En esta lección creaste kiosko.fact_orders_at_scale con un PartitionSpec real —IdentityTransform sobre store_id—, cargaste las diez millones de filas completas de Kiosko a escala, y confirmaste dos hechos con assert: la consulta oculta sigue devolviendo el resultado correcto (9,575,000.00 para S01, sin mencionar ninguna carpeta), y esta vez sí hay poda real —un archivo de tres, no uno de uno—.

Antes de avanzar deberías poder: explicar por qué la carga produjo exactamente tres archivos de datos; reproducir la consulta oculta y su evidencia de poda; y anticipar qué pasaría si filtraras por una columna que no forma parte del PartitionSpec (ejercicio 3).

Tienes una tabla particionada y funcionando. Pero el particionado que elegiste hoy —solo por store_id— no tiene por qué ser el definitivo para siempre. La lección 6 evoluciona este mismo PartitionSpec, agregando una segunda dimensión, sin reescribir ni una sola de las diez millones de filas que ya cargaste.

Recursos

  • PyIceberg — referencia de API, PartitionSpec, PartitionField, catalog.create_table(..., partition_spec=...), table.scan(...).plan_files(). py.iceberg.apache.org/api. En inglés.
  • Apache Iceberg — documentación oficial, "Partitioning" (partición oculta y poda de archivos como consecuencia directa del PartitionSpec). iceberg.apache.org/docs/latest/partitioning. En inglés.
  • DISEÑO de spark-and-distributed-processing-guide — fuente de los números exactos del dataset a escala (10,000,000 filas, 26,537,500.00 de revenue, 9,575,000.00/9,700,000.00/7,262,500.00 por tienda) que esta lección verifica con assert. src/guides/spark-and-distributed-processing-guide/DISENO.md. En español.
  • DISEÑO de esta guía — la sección del módulo 5, con la especificación exacta de kiosko.fact_orders_at_scale y su PartitionSpec inicial. src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.