Módulo 5: Hidden Partitioning And Partition Evolution

Leyendo esquemas de partición viejos y nuevos, juntos

Descripción

La lección 6 evolucionó el PartitionSpec de kiosko.fact_orders_at_scale sin tocar ninguno de los tres archivos existentes. Esta lección hace la prueba definitiva: agrega una franquicia nueva —la número 250,000, la primera que llega después de la evolución— y usa table.inspect.partitions() para mostrar, en una sola consulta, cómo las diez millones de filas viejas y las cuarenta filas nuevas conviven bajo dos esquemas de partición distintos, dentro de la misma tabla, sin ningún conflicto.

Conexión con el módulo. Esta es la lección que cierra el argumento central del módulo: la lección 1 prometió que ambos esquemas iban a convivir; la lección 6 demostró que nada se reescribe; esta lección es la evidencia visual, con datos reales, de esa convivencia. La lección 8 integra las siete lecciones anteriores en un solo proyecto.

Una analogía: el cartero recibe correspondencia bajo la regla nueva

Retomando al cartero de este módulo: la lección 6 lo dejó con una regla nueva —organizar también por día, además de por destinatario—, pero sin haber tocado la correspondencia vieja. Esta lección es el momento en que llega correspondencia genuinamente nueva, y el cartero la archiva usando la regla que tiene vigente ahora: con destinatario y con día. Si en este momento le pides "muéstrame todo lo que tienes archivado", te muestra dos tipos de registro, sin contradicción entre ellos: la correspondencia vieja, archivada con la regla vieja (solo destinatario), y la correspondencia nueva, archivada con la regla nueva (destinatario y día). Ninguna de las dos es "incorrecta" — cada una refleja, con precisión, la regla que estaba vigente en el momento en que se archivó.

Ejemplo trabajado: una franquicia nueva, bajo el spec evolucionado

Paso 1 — Agrega la franquicia 250,000

# l7_read_old_and_new_partitions.py
import os
from datetime import datetime

import pyarrow as pa
import pyarrow.compute as pc
from pyiceberg.catalog import load_catalog

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}",
)
table = catalog.load_table("kiosko.fact_orders_at_scale")


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

# franquicia #250,000 -- la primera que llega DESPUES de evolucionar el spec.
# Mismo generador determinista de siempre, sin random.
NEW_FRANCHISE_ID = 250_000
new_rows = list(generate_orders_at_scale(1))
for r in new_rows:
    r["franchise_id"] = NEW_FRANCHISE_ID
    r["order_id"] = r["order_id"].replace("F000000", f"F{NEW_FRANCHISE_ID:06d}")
    r["order_ts"] = datetime.fromisoformat(r["order_ts"])
new_pa_table = pa.Table.from_pylist(new_rows, schema=PA_SCHEMA)

table.append(new_pa_table)
print(f"franquicia {NEW_FRANCHISE_ID} agregada: {new_pa_table.num_rows} filas nuevas, "
      f"bajo el spec vigente (con order_day)")

Qué esperar (verificado corriendo el script real, continuando sobre la tabla evolucionada en la lección 6):

franquicia 250000 agregada: 40 filas nuevas, bajo el spec vigente (con order_day)

Cuarenta filas nuevas, cargadas con el mismo table.append() de siempre — sin ningún argumento adicional, sin que tuvieras que decirle a Iceberg "usa el spec nuevo". El motor ya sabe cuál es el spec vigente, y lo aplica automáticamente a cualquier escritura nueva.

Paso 2 — Confirma el total, y el revenue de S01 incluyendo la franquicia nueva

row_count_now = table.scan().to_arrow().num_rows
revenue_now = revenue_of(table.scan().to_arrow())
print(f"total ahora: {row_count_now} filas, revenue {round(revenue_now, 2)}")
assert row_count_now == 10_000_000 + 40
assert round(revenue_now, 2) == round(26_537_500.00 + 106.15, 2) == 26_537_606.15

s01_now = table.scan(row_filter="store_id == 'S01'").to_arrow()
s01_revenue_now = revenue_of(s01_now)
print(f"S01 ahora: {s01_now.num_rows} filas, revenue {round(s01_revenue_now, 2)}")
assert round(s01_revenue_now, 2) == round(9_575_000.00 + 38.30, 2) == 9_575_038.30

Qué esperar:

total ahora: 10000040 filas, revenue 26537606.15
S01 ahora: 4000016 filas, revenue 9575038.3

El total subió en exactamente 40 filas y 106.15 de revenue — la franquicia nueva completa, ni una fila de más ni de menos—, y el filtro oculto por S01 sigue funcionando exactamente igual que en la lección 5, ahora incluyendo las 16 líneas de S01 que trajo la franquicia nueva (4,000,000 + 16 = 4,000,016 filas, 9,575,000.00 + 38.30 = 9,575,038.30 de revenue). La consulta filtrada nunca tuvo que enterarse de que el spec cambió a mitad de camino.

Paso 3 — table.inspect.partitions(): ambos esquemas, en la misma tabla

print("=== table.inspect.partitions() -- ambos esquemas de particion ===")
partitions = table.inspect.partitions().to_pylist()
partitions_sorted = sorted(partitions, key=lambda p: (p["spec_id"], str(p["partition"])))
for p in partitions_sorted:
    print(f"{p['spec_id']:>7}  {str(p['partition']):<45}  {p['record_count']}")

spec_ids_present = sorted({p["spec_id"] for p in partitions})
print(f"\nspec_id presentes: {spec_ids_present}")
assert spec_ids_present == [0, 1]

spec0_rows = sum(p["record_count"] for p in partitions if p["spec_id"] == 0)
spec1_rows = sum(p["record_count"] for p in partitions if p["spec_id"] == 1)
assert spec0_rows == 10_000_000
assert spec1_rows == 40

Qué esperar (salida completa; spec_id=0 es el spec original de la lección 5, spec_id=1 es el evolucionado en la lección 6):

=== table.inspect.partitions() -- ambos esquemas de particion ===
      0  {'store_id': 'S01', 'order_day': None}         4000000
      0  {'store_id': 'S02', 'order_day': None}         3250000
      0  {'store_id': 'S03', 'order_day': None}         2750000
      1  {'store_id': 'S01', 'order_day': datetime.date(2026, 8, 3)}  4
      1  {'store_id': 'S01', 'order_day': datetime.date(2026, 8, 4)}  2
      1  {'store_id': 'S01', 'order_day': datetime.date(2026, 8, 5)}  1
      1  {'store_id': 'S01', 'order_day': datetime.date(2026, 8, 6)}  2
      1  {'store_id': 'S01', 'order_day': datetime.date(2026, 8, 7)}  3
      1  {'store_id': 'S01', 'order_day': datetime.date(2026, 8, 8)}  3
      1  {'store_id': 'S01', 'order_day': datetime.date(2026, 8, 9)}  1
      1  {'store_id': 'S02', 'order_day': datetime.date(2026, 8, 3)}  2
      1  {'store_id': 'S02', 'order_day': datetime.date(2026, 8, 4)}  2
      1  {'store_id': 'S02', 'order_day': datetime.date(2026, 8, 5)}  1
      1  {'store_id': 'S02', 'order_day': datetime.date(2026, 8, 6)}  2
      1  {'store_id': 'S02', 'order_day': datetime.date(2026, 8, 7)}  2
      1  {'store_id': 'S02', 'order_day': datetime.date(2026, 8, 8)}  3
      1  {'store_id': 'S02', 'order_day': datetime.date(2026, 8, 9)}  1
      1  {'store_id': 'S03', 'order_day': datetime.date(2026, 8, 3)}  2
      1  {'store_id': 'S03', 'order_day': datetime.date(2026, 8, 4)}  2
      1  {'store_id': 'S03', 'order_day': datetime.date(2026, 8, 6)}  1
      1  {'store_id': 'S03', 'order_day': datetime.date(2026, 8, 7)}  2
      1  {'store_id': 'S03', 'order_day': datetime.date(2026, 8, 8)}  3
      1  {'store_id': 'S03', 'order_day': datetime.date(2026, 8, 9)}  1

spec_id presentes: [0, 1]

Ahí están, lado a lado, los dos esquemas que este módulo prometió desde la lección 1. Las tres primeras filas (spec_id=0) son las diez millones de filas originales — agrupadas solo por store_id, con order_day=None porque esos archivos nunca guardaron esa dimensión (no existía cuando se escribieron). Las veintitrés filas siguientes (spec_id=1) son la franquicia nueva —solo 40 filas en total, repartidas en combinaciones de store_id y order_day reales, con fechas legibles (datetime.date(2026, 8, 3), no el entero crudo que viste en la lección 4) — porque esos archivos sí se escribieron después de la evolución, con la dimensión completa disponible.

Fíjate también en algo que confirma, con datos, la disciplina de este módulo: S03 no tiene ninguna fila para 2026-08-05 —ni en el spec viejo (agregado, no lo distinguirías) ni en el nuevo—, porque la semana real de Kiosko nunca tuvo una orden de S03 ese día (revisa RAW_ORDERS del módulo 1: el miércoles solo tuvo dos órdenes, de S01 y S02). El desglose por día que ves aquí es exactamente la misma semana real que conoces desde la primera guía del ecosistema, ahora visible a través de una partición nueva.

Diagrama: dos specs, una tabla

flowchart TB
    T["kiosko.fact_orders_at_scale"]
    T --> S0["spec_id=0: solo store_id\n10,000,000 filas\norder_day=None"]
    T --> S1["spec_id=1: store_id + order_day\n40 filas (franquicia 250,000)\norder_day poblado"]
    S0 -.->|"archivos de la leccion 5,\nnunca reescritos"| F0["3 archivos de datos"]
    S1 -.->|"archivos de la leccion 7,\nescritos bajo el spec nuevo"| F1["varios archivos chicos\n(uno por combinacion\nstore_id x order_day)"]
    Q["table.scan(row_filter=\"store_id == 'S01'\")"] -->|"lee AMBOS specs\nsin distincion para quien consulta"| S0
    Q --> S1

Errores comunes

Pensar que spec_id es una columna real de la tabla, consultable con row_filter. Qué pasa: alguien, al ver spec_id en el resultado de table.inspect.partitions(), intenta escribir table.scan(row_filter="spec_id == 0") para leer solo las filas viejas. Por qué pasa: spec_id aparece junto a columnas de negocio reales (store_id, order_day) en la misma tabla de resultados, así que es fácil asumir que tiene el mismo estatus. Cómo detectarlo: si tu row_filter menciona spec_id y falla con un error de columna desconocida, revisa esta lección — spec_id es metadata sobre la partición, generada por table.inspect.partitions(), nunca una columna del esquema de negocio de la tabla. Cómo corregirlo: spec_id es información de introspección, útil para entender cómo está organizada la tabla por dentro (exactamente el propósito de esta lección) — no es parte del row_filter de una consulta normal, que solo conoce columnas reales como store_id, product_id u order_ts.

Esperar que order_day=None en las filas viejas signifique que esas filas tienen una fecha faltante o corrupta. Qué pasa: alguien, al ver order_day: None para las diez millones de filas originales, se preocupa pensando que order_ts se perdió o quedó mal cargado para esos datos. Por qué pasa: None suele asociarse con "dato faltante" en el sentido de calidad de datos, no con "esta fila nunca calculó este valor derivado". Cómo detectarlo: si tu preocupación es sobre la calidad de order_ts en las filas viejas, verifica directamente con table.scan(row_filter="franchise_id < 250000").to_arrow().column("order_ts") — vas a ver que order_ts está completo y correcto en cada fila. Cómo corregirlo: order_day es un valor de partición, no una columna del esquema de negocio — no existe como columna consultable en absoluto, es una propiedad de cómo table.inspect.partitions() agrupa archivos. Los archivos viejos simplemente nunca calcularon ese agrupamiento, porque el campo de partición no existía cuando se escribieron; el dato de negocio real, order_ts, nunca estuvo en riesgo.

Ejercicios

Ejercicio 1 — Reproduce la convivencia de ambos specs tú mismo. Con la tabla evolucionada de la lección 6 disponible, corre los tres pasos de esta lección. Confirma 10,000,040 filas totales, 9,575,038.30 para S01, y los dos spec_id (0 y 1) presentes en table.inspect.partitions().

Ver solución

Si tu tabla llegó a esta lección con exactamente el estado de las lecciones 5 y 6, tu salida debería coincidir número por número. Presta atención especial a spec_id presentes: [0, 1] — si solo ves [1], probablemente tu tabla no tenía datos previos a la evolución (revisa que la lección 5 haya cargado las diez millones de filas antes de evolucionar el spec en la lección 6); si solo ves [0], probablemente el append() de la franquicia nueva de esta lección no se ejecutó.

Ejercicio 2 — Agrega una segunda franquicia nueva, y confirma que también aterriza bajo spec_id=1. Repite el paso 1 de esta lección con NEW_FRANCHISE_ID = 250_001, y confirma con table.inspect.partitions() que las filas de ambas franquicias nuevas (250,000 y 250,001) aparecen bajo el mismo spec_id=1 que ya viste — el spec no vuelve a evolucionar solo porque llegaron más datos.

Ver solución
NEW_FRANCHISE_ID_2 = 250_001
more_rows = list(generate_orders_at_scale(1))
for r in more_rows:
    r["franchise_id"] = NEW_FRANCHISE_ID_2
    r["order_id"] = r["order_id"].replace("F000000", f"F{NEW_FRANCHISE_ID_2:06d}")
    r["order_ts"] = datetime.fromisoformat(r["order_ts"])
more_pa_table = pa.Table.from_pylist(more_rows, schema=PA_SCHEMA)
table.append(more_pa_table)

partitions_v2 = table.inspect.partitions().to_pylist()
spec_ids_v2 = sorted({p["spec_id"] for p in partitions_v2})
print("spec_id presentes tras la segunda franquicia:", spec_ids_v2)
assert spec_ids_v2 == [0, 1]

El resultado sigue siendo [0, 1] — ningún spec nuevo se creó solo por escribir más datos. Un spec_id nuevo solo aparece cuando alguien evoluciona el PartitionSpec explícitamente, con update_spec(), exactamente como hizo la lección 6. Escribir datos, por más que sean muchas franquicias nuevas, siempre usa el spec vigente en ese momento — nunca crea uno nuevo por sí solo.

Ejercicio 3 — Explica, en 2-3 frases, por qué esta tabla nunca necesitó una operación de "migración" o "backfill" de order_day para las filas viejas. Contrasta esto con lo que probablemente esperarías en un almacén de datos relacional tradicional, donde agregar una columna calculada a una tabla existente sí suele requerir un paso explícito de backfill.

Ver solución

En un almacén relacional tradicional, agregar una columna con datos derivados (como order_day a partir de order_ts) típicamente implica un UPDATE masivo o un job de backfill que recorre cada fila existente y calcula el valor nuevo — una operación costosa, que puede tomar horas en una tabla grande, y que bloquea o compite por recursos con el tráfico normal. Aquí no hizo falta ningún backfill porque order_day nunca es una columna de la tabla en el sentido tradicional: es un valor de partición, calculado únicamente para decidir en qué archivo escribir una fila en el momento de escribirla. Las filas viejas nunca necesitaron ese cálculo retroactivo porque, para propósitos de lectura de negocio (order_ts sigue ahí, completo), nunca les faltó nada — lo único que "falta" es un agrupamiento de archivos que, para esas filas, simplemente nunca se hizo, y no hace ninguna falta que se haga.

Resumen y siguiente paso

En esta lección agregaste una franquicia nueva a kiosko.fact_orders_at_scale, después de haber evolucionado su PartitionSpec en la lección 6, y confirmaste con table.inspect.partitions() que las diez millones de filas originales (spec_id=0, solo store_id) y las cuarenta filas nuevas (spec_id=1, store_id + order_day) conviven, sin conflicto, dentro de la misma tabla. Verificaste, también, que la consulta oculta de la lección 3 sigue funcionando exactamente igual —row_filter="store_id == 'S01'"— sin que el código tuviera que enterarse de que el spec cambió a mitad de camino.

Antes de avanzar deberías poder: leer una salida de table.inspect.partitions() con múltiples spec_id, y explicar qué significa cada uno; y explicar, sin ambigüedad, por qué las filas viejas nunca necesitaron un backfill.

Tienes las siete piezas completas: el archivero visible de Spark, la consulta oculta, los tres transforms, la tabla particionada a escala, la evolución sin reescritura, y la convivencia de ambos esquemas. La lección 8 las integra en un solo proyecto, de punta a punta, con assert automáticos sobre cada una.

Recursos

  • PyIceberg — referencia de API, table.inspect.partitions(), la estructura exacta de su resultado (spec_id, partition, record_count). py.iceberg.apache.org/api. En inglés.
  • Apache Iceberg — documentación oficial, "Partitioning", sección "Partition Evolution" (la garantía formal de que datos escritos bajo specs distintos coexisten sin conflicto en la misma tabla). iceberg.apache.org/docs/latest/partitioning. En inglés.
  • DISEÑO de esta guía — la sección del módulo 5, con la verificación exacta esperada (table.inspect.partitions() mostrando ambos esquemas de partición conviviendo). src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.