Módulo 4: Schema Evolution Without Rewriting
Proyecto: el dim_store de Kiosko evolucionado
Descripción
Este proyecto cierra el módulo 4. Sabes por qué una escritura Iceberg es atómica, con evidencia de 490 lecturas concurrentes que nunca capturaron un estado a medias (lección 2). Sabes agregar una columna sin tocar ningún archivo de datos (lección 3), y renombrar o borrar una sin ese mismo costo (lección 4). Poblaste country de forma determinística desde city (lección 5), confirmaste que un snapshot anterior a la evolución sigue leyendo su propio esquema (lección 6), y viste, con un conflicto real entre dos escritores, exactamente qué garantiza "ACID" en el contexto de una tabla Iceberg (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: evolucionar el esquema de kiosko.dim_store —agregar, poblar, renombrar, borrar— sin reescribir un solo archivo de datos existente, y sin que ningún lector, en ningún momento, vea la tabla a medio camino entre un estado y el siguiente.
Una analogía: el formulario completo, revisado de punta a punta
Las lecciones 2 a 7 de este módulo construyeron, una pieza a la vez, la evidencia completa de que la evolución de esquema de Iceberg es segura: la prueba de que ningún lector ve un estado intermedio (lección 2), el mecanismo de agregar una columna sin tocar datos (lección 3), la misma garantía para renombrar y borrar (lección 4), el pago real —country poblada— (lección 5), la confirmación de que el pasado se sigue leyendo con su propio esquema (lección 6), y el límite preciso de la palabra "ACID" (lección 7). Este proyecto es el momento de repetir todo el proceso, de punta a punta, en un solo gesto continuo — el mismo tipo de integración que ya hiciste al cerrar los módulos 1 y 3.
El material: todo lo que este módulo construyó, en un solo lugar
Necesitas, en un directorio de trabajo nuevo:
kiosko_dim_store_evolution/
└── kiosko_evolved_dim_store.py (este proyecto: junta las 6 piezas de datos)
Con PyIceberg instalado en tu entorno (pip install "pyiceberg[sql-sqlite,pyarrow]", módulo 1, lección 4). Este proyecto es autocontenido: crea kiosko.dim_store 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 kiosko.fact_orders o kiosko.dim_product de los módulos 1 a 3 en ese mismo catálogo, este proyecto las deja intactas; si no las tienes, kiosko.dim_store se crea igual, sola).
La solución de referencia, verificada
# kiosko_evolved_dim_store.py -- proyecto de cierre del modulo 4
# dim_store evoluciona (add country, rename+drop temp_notes) sin reescribir un solo Parquet existente
import os
import pyarrow as pa
from pyiceberg.catalog import load_catalog
from pyiceberg.schema import Schema
from pyiceberg.types import NestedField, StringType
DIM_STORE_V1 = [
{"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"},
]
CITY_TO_COUNTRY = {"Bogota": "Colombia", "Lima": "Peru", "Santiago": "Chile"}
DIM_STORE_SCHEMA_V1 = Schema(
NestedField(field_id=1, name="store_id", field_type=StringType(), required=True),
NestedField(field_id=2, name="store_name", field_type=StringType(), required=True),
NestedField(field_id=3, name="city", field_type=StringType(), required=True),
)
def dim_store_v1_pa_table() -> pa.Table:
schema = pa.schema([
pa.field("store_id", pa.string(), nullable=False),
pa.field("store_name", pa.string(), nullable=False),
pa.field("city", pa.string(), nullable=False),
])
return pa.Table.from_pylist(DIM_STORE_V1, schema=schema)
def dim_store_with_country_pa_table() -> pa.Table:
schema = pa.schema([
pa.field("store_id", pa.string(), nullable=False),
pa.field("store_name", pa.string(), nullable=False),
pa.field("city", pa.string(), nullable=False),
pa.field("country", pa.string(), nullable=True),
])
rows = [{**s, "country": CITY_TO_COUNTRY[s["city"]]} for s in DIM_STORE_V1]
return pa.Table.from_pylist(rows, schema=schema)
def data_file_names(table) -> list:
return sorted(f["file_path"].split("/")[-1] for f in table.inspect.files().to_pylist())
def main() -> None:
print("=== Kiosko: dim_store evolucionada sin reescribir datos ===\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")
print(f"Paso 1/8 -- catalogo '{catalog.name}' y namespace 'kiosko' listos")
dim_store = catalog.create_table("kiosko.dim_store", schema=DIM_STORE_SCHEMA_V1)
dim_store.append(dim_store_v1_pa_table())
snap_before_evolution = dim_store.current_snapshot().snapshot_id
files_before_evolution = data_file_names(dim_store)
print(f"Paso 2/8 -- dim_store cargada: {dim_store.scan().to_arrow().num_rows} filas, "
f"snap_before_evolution capturado, {len(files_before_evolution)} archivo(s) de datos")
with dim_store.update_schema() as update:
update.add_column("country", StringType())
files_after_add_column = data_file_names(dim_store)
rows_after_add_column = dim_store.scan().to_arrow().to_pylist()
print(f"Paso 3/8 -- add_column('country') aplicado. Archivos de datos sin cambios: "
f"{files_before_evolution == files_after_add_column}. "
f"country en filas existentes: {[r['country'] for r in rows_after_add_column]}")
with dim_store.update_schema() as update:
update.rename_column("store_name", "outlet_name")
with dim_store.update_schema() as update:
update.rename_column("outlet_name", "store_name")
with dim_store.update_schema() as update:
update.add_column("temp_notes", StringType())
with dim_store.update_schema() as update:
update.delete_column("temp_notes")
files_after_schema_ops = data_file_names(dim_store)
schema_names_after_ops = [f.name for f in dim_store.schema().fields]
print(f"Paso 4/8 -- rename store_name<->outlet_name + add/drop temp_notes. "
f"Archivos de datos sin cambios: {files_before_evolution == files_after_schema_ops}. "
f"Esquema final: {schema_names_after_ops}")
dim_store.overwrite(dim_store_with_country_pa_table())
print("Paso 5/8 -- table.overwrite() puebla 'country' desde 'city' para las 3 filas existentes")
current_rows = dim_store.scan().to_arrow().to_pylist()
print("Paso 6/8 -- estado vigente de dim_store:")
for row in sorted(current_rows, key=lambda r: r["store_id"]):
print(f" {row['store_id']} {row['store_name']:<15} {row['city']:<10} country={row['country']}")
old_scan = dim_store.scan(snapshot_id=snap_before_evolution).to_arrow()
old_rows = old_scan.to_pylist()
print(f"Paso 7/8 -- table.scan(snapshot_id=snap_before_evolution) -- "
f"columnas: {old_scan.schema.names}, filas: {len(old_rows)}")
total_snapshots = len(dim_store.history())
print(f"Paso 8/8 -- table.history() tiene {total_snapshots} entradas en total\n")
print("=== Verificacion final ===\n")
assert files_before_evolution == files_after_add_column, "add_column no debe tocar los archivos de datos"
assert files_before_evolution == files_after_schema_ops, "rename/add/drop no deben tocar los archivos de datos"
assert all(r["country"] is None for r in rows_after_add_column), \
"tras add_column, las filas existentes deben tener country=None (sin poblar todavia)"
assert schema_names_after_ops == ["store_id", "store_name", "city", "country"], \
"el esquema final no debe conservar temp_notes ni el rename"
country_by_store = {r["store_id"]: r["country"] for r in current_rows}
assert country_by_store == {"S01": "Colombia", "S02": "Peru", "S03": "Chile"}
assert old_scan.schema.names == ["store_id", "store_name", "city"], \
"el snapshot anterior a la evolucion debe leerse con el esquema de 3 columnas, sin country"
assert len(old_rows) == 3
assert all("country" not in r for r in old_rows)
print("Todas las verificaciones pasaron:")
print(" - add_column('country') no toco ningun archivo de datos existente")
print(" - rename_column + add/delete_column('temp_notes') tampoco tocaron archivos de datos")
print(" - country quedo poblado deterministicamente desde city (S01=Colombia, S02=Peru, S03=Chile)")
print(" - table.scan(snapshot_id=snap_before_evolution) sigue leyendo el esquema de 3 columnas, sin country")
if __name__ == "__main__":
main()
Qué esperar (verificado corriendo python3 kiosko_evolved_dim_store.py real, de punta a punta, en un directorio nuevo; ningún snapshot_id se imprime como literal — la lección 3 de este módulo, heredando la regla del módulo 3, explica por qué):
=== Kiosko: dim_store evolucionada sin reescribir datos ===
Paso 1/8 -- catalogo 'kiosko' y namespace 'kiosko' listos
Paso 2/8 -- dim_store cargada: 3 filas, snap_before_evolution capturado, 1 archivo(s) de datos
Paso 3/8 -- add_column('country') aplicado. Archivos de datos sin cambios: True. country en filas existentes: [None, None, None]
Paso 4/8 -- rename store_name<->outlet_name + add/drop temp_notes. Archivos de datos sin cambios: True. Esquema final: ['store_id', 'store_name', 'city', 'country']
Paso 5/8 -- table.overwrite() puebla 'country' desde 'city' para las 3 filas existentes
Paso 6/8 -- estado vigente de dim_store:
S01 Kiosko Centro Bogota country=Colombia
S02 Kiosko Norte Lima country=Peru
S03 Kiosko Sur Santiago country=Chile
Paso 7/8 -- table.scan(snapshot_id=snap_before_evolution) -- columnas: ['store_id', 'store_name', 'city'], filas: 3
Paso 8/8 -- table.history() tiene 3 entradas en total
=== Verificacion final ===
Todas las verificaciones pasaron:
- add_column('country') no toco ningun archivo de datos existente
- rename_column + add/delete_column('temp_notes') tampoco tocaron archivos de datos
- country quedo poblado deterministicamente desde city (S01=Colombia, S02=Peru, S03=Chile)
- table.scan(snapshot_id=snap_before_evolution) sigue leyendo el esquema de 3 columnas, sin country
Fíjate en el paso 8: table.history() reporta tres entradas, no cuatro ni cinco, a pesar de que este script corrió add_column, dos rename_column, add_column de nuevo, delete_column, y finalmente overwrite(). Las tres entradas son, en orden: el append original (paso 2), y el delete+append internos del overwrite() que pobló country (paso 5) — exactamente el mismo patrón que el módulo 3 encontró para el cambio de P002. Ninguna de las cinco operaciones de esquema del paso 3 y el paso 4 agregó una sola entrada al historial de snapshots, porque ninguna de ellas es una escritura de datos — es, con precisión, la evidencia consolidada de todo este módulo, en un solo número.
Diagrama: de dónde venías, a dónde llegaste
flowchart LR
A["Modulo 1-3:\nfact_orders, dim_product,\ntime travel"] --> B["Leccion 2:\n490 lecturas, 0 estados\nintermedios -- atomicidad"]
B --> C["Leccion 3:\nadd_column('country')\n0 archivos tocados"]
C --> D["Leccion 4:\nrename + add/drop\ntemp_notes, misma garantia"]
D --> E["Leccion 5:\ncountry poblado:\nColombia/Peru/Chile"]
E --> F["Leccion 6:\nsnap_before_evolution\nlee su propio esquema"]
F --> G["Leccion 7:\nCommitFailedException real,\nel limite preciso de ACID"]
G --> H["Este proyecto:\nlas 6 piezas de datos,\nun solo script, assert automatico"]
H --> I["Modulo 5:\nparticion oculta y\nevolucion de particion"]
Cerrando la promesa del módulo, punto por punto
| Lo que la lección 1 prometió | Evidencia de que este módulo lo entregó |
|---|---|
Por qué overwrite-partition nunca fue atómico, con la cita exacta | Lección 2: la cita literal de data-engineering-foundations-guide, y 490 lecturas concurrentes que nunca vieron un estado intermedio en Iceberg |
| Agregar una columna sin reescribir datos | Lección 3: add_column('country'), misma lista de archivos antes y después, table.history() sin cambios |
| Renombrar y borrar columnas con la misma garantía | Lección 4: rename_column de ida y vuelta, temp_notes agregada y borrada, field_id nunca reutilizado |
country poblada desde city, de forma determinística | Lección 5 y este proyecto: country_by_store == {"S01": "Colombia", "S02": "Peru", "S03": "Chile"}, verificado con assert |
| Un snapshot anterior a la evolución sigue leyendo su propio esquema | Lección 6 y este proyecto: old_scan.schema.names == ["store_id", "store_name", "city"], sin country |
| Qué garantiza "ACID" aquí, con precisión | Lección 7: CommitFailedException real, con el mensaje exacto de un conflicto de escritores, y el límite explícito frente a transacciones multi-tabla |
Este módulo no evolucionó el esquema de kiosko.fact_orders ni de kiosko.dim_product —esas tablas siguen exactamente como las dejó el módulo 3—, ni tocó partición alguna —eso llega en el módulo 5—. Lo que este módulo entrega es exactamente lo que prometió: una tabla nueva, kiosko.dim_store, evolucionada de tres columnas a cuatro, con datos poblados de forma determinística, sin reescribir un solo archivo de datos existente en ningún paso, y con la garantía de atomicidad y aislamiento 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.dim_store de una lección anterior de este módulo. Qué pasa: alguien corre este proyecto en el mismo directorio donde ya completó las lecciones 3 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_evolved_dim_store.py, ya tienes un catálogo con kiosko.dim_store registrada en el mismo directorio. Cómo corregirlo: corre este proyecto en un directorio de trabajo nuevo, separado de donde hiciste las lecciones 3 a 7 — tal como sugiere "El material" de esta lección.
Esperar que table.history() reporte 8 entradas, una por cada paso numerado del script. Qué pasa: alguien, viendo los "Paso 1/8" a "Paso 8/8" impresos por el script, asume que cada uno corresponde a un snapshot nuevo, y se sorprende al ver solo 3 en table.history(). Por qué pasa: los números de "Paso N/8" numeran las etapas narrativas del script —creación del catálogo, carga, evolución de esquema, población, lectura histórica, verificación—, no cada llamada individual a update_schema() ni cada operación de escritura. Cómo detectarlo: si tu conteo esperado de snapshots coincide con el número de "pasos" impresos en pantalla, revisa la distinción entre operaciones de esquema (que no crean snapshots) y operaciones de datos (que sí) que las lecciones 3 y 5 de este módulo ya establecieron. Cómo corregirlo: cuenta snapshots con table.history() o table.inspect.snapshots() directamente — la respuesta correcta es 3: el append original, y el delete+append internos del único overwrite() que este proyecto ejecuta.
Ejercicios
Ejercicio 1 — Corre el proyecto completo tú mismo, desde cero. En un directorio nuevo, corre python3 kiosko_evolved_dim_store.py. Confirma que ves los ocho pasos completarse y el mensaje final con las cuatro 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 tres países correctos, el esquema de tres columnas confirmado para el snapshot histórico, y el mensaje de éxito con las cuatro condiciones confirmadas.
Ejercicio 2 — Rompe un assert a propósito, y observa el fallo. Cambia temporalmente CITY_TO_COUNTRY["Lima"] de "Peru" a "Perú" (con tilde), corre el script de nuevo, y observa qué assert falla primero. Después revierte el cambio.
Ver solución
El assert que falla es assert country_by_store == {"S01": "Colombia", "S02": "Peru", "S03": "Chile"} — porque el diccionario ahora produce "Perú" con tilde para S02, que ya no coincide, carácter por carácter, con el valor esperado "Peru" sin tilde. Este ejercicio demuestra dos cosas a la vez: que los assert de este proyecto están encadenados a los valores canónicos exactos del mapeo CITY_TO_COUNTRY, y que esta guía, siguiendo la convención dura de todo el ecosistema de que los identificadores y datos de código van en inglés (sin acentos ni caracteres especiales), usa "Peru" sin tilde como el valor de negocio correcto, aunque la prosa en español de esta misma lección sí escriba "Perú" con tilde al referirse al país.
Ejercicio 3 — Explica, en tus propias palabras, por qué este proyecto verifica files_before_evolution == files_after_schema_ops en vez de solo verificar el resultado final de country. En 3-4 frases, justifica por qué el assert de archivos de datos es tan importante como el assert del valor final de country.
Ver solución
Verificar solo el resultado final —que country tenga los tres países correctos— confirmaría que los datos son correctos, pero no confirmaría cómo se llegó ahí: una implementación distinta, que reescribiera todos los archivos de datos en cada operación de esquema (el comportamiento que este módulo completo argumenta que Iceberg evita), podría llegar exactamente al mismo resultado final, sin ninguna de las ventajas que este módulo se propuso demostrar. Verificar files_before_evolution == files_after_schema_ops confirma la afirmación central del módulo —que agregar, renombrar y borrar columnas son operaciones de metadata, no de datos—, con evidencia directa sobre el sistema de archivos, no solo con el resultado correcto de una consulta. Es la misma disciplina de "verificar el camino, no solo el destino" que ya viste en el proyecto del módulo 3, donde el assert del margen "roto" era tan importante como el del margen "correcto".
Resumen y siguiente paso: el cierre de este módulo
Con este proyecto cierras el módulo 4. Integraste las siete lecciones anteriores —la cita exacta de data-engineering-foundations-guide y su contraste con la atomicidad de Iceberg, el mecanismo de agregar/renombrar/borrar columnas sin reescribir datos, la población real de country, la lectura correcta del pasado, y el límite preciso de "ACID"— 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 de dimensión cuyo esquema evolucionó después de que los datos ya existían, sin que ninguna de esas operaciones tocara un solo archivo Parquet ya escrito. kiosko.dim_store pasó de tres columnas a cuatro, con country derivada de forma determinística de city — y el mecanismo que lo hizo posible es el mismo que garantiza, en cualquier tabla Iceberg, que ningún lector vea nunca un esquema a medio cambiar.
Hacia dónde sigues. El módulo 5 —Partición oculta y evolución de partición— toma kiosko.fact_orders_at_scale, la tabla de diez millones de filas heredada de spark-and-distributed-processing-guide, y contrasta el particionado por carpetas de Spark (partitionBy("store_id"), visible, hay que conocer la estructura) con la partición oculta de Iceberg —la consulta filtra por la columna de negocio, el motor decide el layout— y su evolución hacia adelante, sin reescribir los diez millones de filas ya existentes.
Recursos
- PyIceberg — documentación oficial (quickstart), el flujo completo de catálogo, tabla,
append(),update_schema()yoverwrite()que integra este proyecto. py.iceberg.apache.org. En inglés. - PyIceberg — referencia de API,
table.update_schema(),table.scan(snapshot_id=...),table.history(),table.inspect.files(). py.iceberg.apache.org/api. En inglés. - Apache Iceberg — documentación oficial, "Evolution" y "Reliability", fuente formal de las garantías que este proyecto verifica con
assert. iceberg.apache.org/docs/latest/evolution · iceberg.apache.org/docs/latest/reliability. En inglés. - DISEÑO de
data-engineering-foundations-guide— fuente del esquema original destoresy de la cita exacta sobre el riesgo deoverwrite-partitionque este módulo completo contrasta.src/guides/data-engineering-foundations-guide/DISENO.md. En español. - DISEÑO de esta guía — el mapa completo de los ocho módulos, incluido el módulo 5 que sigue.
src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.