Módulo 4: Schema Evolution Without Rewriting

Por qué overwrite-partition nunca fue atómico

Descripción

Esta lección retoma, palabra por palabra, la cita de data-engineering-foundations-guide que el módulo 1 de esta guía ya señaló como una promesa pendiente: el patrón overwrite-partitionDELETE de la partición, seguido de INSERT de los datos nuevos— resuelve el problema de los reruns duplicados, pero deja un instante real, entre esas dos operaciones, donde la partición queda vacía. Esta lección no repite esa advertencia en abstracto: construye, con PyIceberg real, un lector concurrente que intenta sorprender a una escritura Iceberg "a medias", y reporta con evidencia ejecutada que nunca lo logra.

Conexión con el módulo. La lección 1 prometió esta comparación como el punto de partida del módulo. Sin esta lección, "una escritura Iceberg es atómica" sería una afirmación de folleto, del mismo tipo que cualquier documentación de producto asegura sobre sí misma. Con ella, es una afirmación que corriste, con tus propios ojos, contra el mismo mecanismo que las lecciones 3 a 8 van a usar para evolucionar el esquema de kiosko.dim_store.

Una analogía: el mostrador que abre o actualiza en un solo trámite

Vuelve al cuaderno de data-engineering-foundations-guide: arrancar la página del martes y escribir una nueva es mejor que pegar una hoja encima, pero sigue siendo dos gestos separados —arrancar, después escribir—, y entre uno y otro, la página no existe. Ahora imagina, en cambio, el mostrador de una oficina de trámites que actualiza un registro de una forma distinta: el empleado prepara el documento nuevo completo, por su cuenta, en su escritorio, sin que nadie más pueda verlo todavía. Solo cuando el documento nuevo está terminado del todo, lo intercambia por el viejo en el archivero — un solo movimiento, un solo instante, sin que el archivero quede vacío ni un microsegundo. Cualquiera que consulte el archivero en cualquier momento, incluso mientras el empleado prepara el documento nuevo en su escritorio, sigue viendo el documento viejo completo, hasta el instante exacto del intercambio — después de ese instante, ve el nuevo, completo. Nunca ve un archivero vacío, ni un documento a medio escribir.

Eso es, con precisión, lo que hace una escritura de Iceberg. El equivalente al "documento nuevo preparado en el escritorio" es el archivo de metadata nuevo, escrito completo en disco antes de que nadie lo necesite. El equivalente al "intercambio en el archivero" es un único movimiento del puntero del catálogo, del archivo de metadata viejo al nuevo. Ningún lector externo —ni siquiera uno que esté consultando la tabla en el instante exacto de la escritura— puede ver un estado intermedio, porque ese estado intermedio nunca se publica: existe, como mucho, en el escritorio del empleado, invisible para cualquiera que consulte el archivero.

Ejemplo trabajado: un lector que intenta sorprender a un overwrite() en el acto

Paso 1 — El experimento: un lector concurrente, polling sin descanso

Esta ilustración usa un catálogo y un warehouse propios, descartables, para no mezclar el experimento con kiosko.dim_store —ese trabajo empieza recién en la lección 3—. El objetivo es puramente medir: ¿puede un lector externo, consultando la tabla lo más rápido posible, capturar algún estado que no sea "el de antes" o "el de después" de un overwrite()?

# atomic_overwrite_demo.py -- ilustracion aislada, NO parte del modelo de Kiosko
import os
import shutil
import threading

import pyarrow as pa
from pyiceberg.catalog import load_catalog
from pyiceberg.schema import Schema
from pyiceberg.types import IntegerType, NestedField

demo_warehouse = os.path.abspath("atomic_demo_warehouse")
demo_db = os.path.abspath("atomic_demo_catalog.db")
os.makedirs(demo_warehouse, exist_ok=True)

CATALOG_KWARGS = dict(type="sql", uri=f"sqlite:///{demo_db}", warehouse=f"file://{demo_warehouse}")
writer_catalog = load_catalog("atomic_demo", **CATALOG_KWARGS)
writer_catalog.create_namespace("atomic_demo")

schema = Schema(NestedField(field_id=1, name="n", field_type=IntegerType(), required=True))
table = writer_catalog.create_table("atomic_demo.race_check", schema=schema)
pa_schema = pa.schema([pa.field("n", pa.int32(), nullable=False)])

# el lector concurrente usa el MISMO nombre de catalogo -- el registro de tablas
# de un SqlCatalog esta separado por catalog_name, no solo por el archivo sqlite
reader_catalog = load_catalog("atomic_demo", **CATALOG_KWARGS)

ROUNDS = 25
ROW_COUNTS = [3_000_000, 500]  # dos tamanos bien distintos: un estado intermedio seria inconfundible

table.append(pa.Table.from_pylist([{"n": i} for i in range(ROW_COUNTS[0])], schema=pa_schema))
reader_catalog.load_table("atomic_demo.race_check")  # precalienta la conexion del lector

all_observed = set()
total_polls = 0

for round_i in range(ROUNDS):
    target_count = ROW_COUNTS[round_i % 2]
    replacement = pa.Table.from_pylist([{"n": i} for i in range(target_count)], schema=pa_schema)

    observed_this_round = []
    stop_polling = threading.Event()

    def poll_reader():
        while not stop_polling.is_set():
            reader_table = reader_catalog.load_table("atomic_demo.race_check")
            observed_this_round.append(reader_table.scan().to_arrow().num_rows)

    poller = threading.Thread(target=poll_reader)
    poller.start()
    table.overwrite(replacement)  # el lector sigue haciendo polling MIENTRAS esto corre
    stop_polling.set()
    poller.join(timeout=5)

    total_polls += len(observed_this_round)
    all_observed.update(observed_this_round)

print(f"{ROUNDS} rondas de table.overwrite() alternando entre {ROW_COUNTS[0]:,} y {ROW_COUNTS[1]} filas")
print(f"Total de lecturas concurrentes del lector externo durante las escrituras: {total_polls}")
print(f"Conjunto COMPLETO de valores de num_rows observados en todas las rondas: {sorted(all_observed)}")
print(f"Estado final de la tabla: {table.scan().to_arrow().num_rows} filas")

Qué esperar (verificado corriendo el script real; el número exacto de lecturas por ronda varía con la velocidad de tu propia máquina, pero el conjunto de valores observados es la parte que importa):

25 rondas de table.overwrite() alternando entre 3,000,000 y 500 filas
Total de lecturas concurrentes del lector externo durante las escrituras: 490
Conjunto COMPLETO de valores de num_rows observados en todas las rondas: [500, 3000000]
Estado final de la tabla: 3000000 filas

Cuatrocientas noventa lecturas concurrentes, mientras 25 escrituras reemplazaban la tabla completa, alternando entre tres millones de filas y quinientas. Ninguna de esas 490 lecturas vio 0. Ninguna vio ningún número que no fuera exactamente 3000000 o 500 — nunca un valor a medio camino, nunca el estado vacío que el delete interno de overwrite() (visto en el módulo 3) sí produce como snapshot archivado. El lector externo, corriendo en otro hilo, consultando el catálogo lo más rápido que puede, nunca alcanza a ver nada que no sea "el estado completo de antes" o "el estado completo de después".

Paso 2 — Por qué esto no es casualidad: una sola confirmación al catálogo, nunca dos

La razón no es que el experimento tuviera suerte 490 veces seguidas — es que, por diseño, no hay ningún punto intermedio que un lector externo pueda observar. Table.overwrite(), por dentro, abre una Transaction:

# el codigo real de pyiceberg 0.11.1 -- Table.overwrite()
with self.transaction() as tx:
    tx.overwrite(df=df, overwrite_filter=overwrite_filter, ...)

El propio docstring de transaction() en PyIceberg 0.11.1 lo dice sin rodeos: "Create a new transaction object to first stage the changes, and then commit them to the catalog." Fíjate en las dos palabras clave: stage, primero; commit, después, como un paso separado y único. Todo lo que ocurre dentro del bloque with —incluido el snapshot delete y el snapshot append que el módulo 3 encontró dentro de un solo overwrite()— se acumula en memoria, sin tocar el catálogo para nada. Recién cuando el bloque with termina, commit_transaction() empaqueta todos los cambios acumulados en una sola llamada a catalog.commit_table(...), que hace, en el catálogo SQL de esta guía, exactamente esto:

-- el patron real que usa SqlCatalog.commit_table() en pyiceberg 0.11.1
UPDATE iceberg_tables
SET metadata_location = :nueva_ruta
WHERE catalog_name = :catalogo
  AND table_namespace = :namespace
  AND table_name = :tabla
  AND metadata_location = :ruta_vieja_esperada

Un único UPDATE, con una condición WHERE que exige que la ruta vieja siga siendo exactamente la que este proceso vio al empezar. Antes de correr ese UPDATE, el archivo de metadata nuevo ya está completo en disco —ya tiene los dos snapshots, delete y append, ya escritos—. El UPDATE no construye nada: solo mueve un puntero, de una fila ya completa a otra fila ya completa, en una sola sentencia SQL que la propia base de datos garantiza atómica. No hay ningún momento en que el catálogo apunte a "la mitad" de un cambio, porque el catálogo nunca conoció ningún estado intermedio — solo conoció el de antes, y después, en un solo paso, el de después.

Diagrama: dos caminos hacia "reemplazar los datos de una fecha"

flowchart TB
    subgraph fnd["overwrite-partition -- data-engineering-foundations-guide M6"]
        F1["DELETE FROM staging_demo\nWHERE dt = fecha"] --> F2["instante real:\nla particion esta vacia,\nvisible para cualquier lector"]
        F2 --> F3["INSERT INTO staging_demo\n... datos nuevos"]
    end
    subgraph ice["table.overwrite() -- Apache Iceberg, esta leccion"]
        I1["Transaction: stage\nsnapshot delete + snapshot append\n(en memoria, invisible desde afuera)"] --> I2["metadata nuevo escrito COMPLETO\nen disco, todavia no referenciado"]
        I2 --> I3["UN SOLO UPDATE atomico\ndel puntero del catalogo"]
        I3 --> I4["visible desde afuera:\nsolo 'antes completo'\no 'despues completo'"]
    end

Profundización: lo que la documentación oficial de Iceberg confirma

Esta lección no inventa el término "atómico" — lo toma directo de la documentación oficial de Apache Iceberg, sección "Reliability": "Commits replace the path of the current table metadata file using an atomic operation. This ensures that all updates to table data and metadata are atomic, and is the basis for serializable isolation." Y sobre qué pasa cuando dos escrituras compiten por el mismo instante: "Iceberg supports multiple concurrent writes using optimistic concurrency. Each writer assumes that no other writers are operating and writes out new table metadata for an operation. Then, the writer attempts to commit by atomically swapping the new table metadata file for the existing metadata file. If the atomic swap fails because another writer has committed, the failed writer retries [...]". La lección 7 de este módulo va a construir, con código real, exactamente ese escenario de dos escritores compitiendo — por ahora, retén la idea central: la atomicidad de Iceberg no depende de que nadie más esté escribiendo al mismo tiempo; depende de que el intercambio del puntero sea, siempre, una sola operación indivisible, gane quien gane la carrera.

Vale la pena decirlo con la misma honestidad que data-engineering-foundations-guide mostró sobre su propio patrón: esto no hace que overwrite-partition esté "mal" — sigue siendo, hoy, un patrón válido para pipelines batch simples sin ningún formato de tabla transaccional debajo. Lo que cambia es que, con Iceberg, la garantía que esa guía tuvo que nombrar como un costo a asumir —"hay un costo real en esta decisión, y vale la pena nombrarlo con honestidad"— deja de ser un costo: el formato de tabla la resuelve de fábrica, sin que nadie tenga que envolver un DELETE y un INSERT en una transacción manual.

Errores comunes

Pensar que "atómico" significa que table.overwrite() no puede fallar nunca. Qué pasa: alguien concluye que, como el overwrite() es atómico, no hace falta manejar ningún error al llamarlo. Por qué pasa: "atómico" suena, por asociación, a "infalible" o "seguro en todos los sentidos". Cómo detectarlo: si tu código no contempla la posibilidad de que table.overwrite() lance una excepción, no interiorizaste todavía la diferencia. Cómo corregirlo: atómico significa que la operación o se aplica completa, o no se aplica en absoluto — puede fallar del todo (por ejemplo, si otro escritor ganó la carrera del UPDATE ... WHERE metadata_location = ..., la lección 7 muestra ese caso exacto con una excepción real), pero nunca deja la tabla en un estado parcialmente aplicado. Fallar-completo es exactamente lo que la atomicidad garantiza; fallar-a-medias es lo que previene.

Asumir que esta garantía aplica a varias tablas a la vez, como en una transacción de base de datos relacional. Qué pasa: alguien, después de ver esta lección, espera que actualizar kiosko.dim_store y kiosko.fact_orders en el mismo script se comporte como una transacción SQL que puede hacer ROLLBACK de ambas tablas si una de las dos falla. Por qué pasa: "ACID" es un término que la mayoría aprendió primero en el contexto de una base de datos relacional, donde sí cubre transacciones multi-tabla. Cómo detectarlo: si tu código depende de que un error al escribir la segunda de dos tablas Iceberg deshaga automáticamente lo que ya se escribió en la primera, estás asumiendo una garantía que Iceberg no da. Cómo corregirlo: la atomicidad de Iceberg es por tabla, no por transacción multi-tabla — cada commit_table() es su propio movimiento de puntero, independiente del de cualquier otra tabla. La lección 7 de este módulo precisa exactamente hasta dónde llega esta garantía, y dónde termina.

Ejercicios

Ejercicio 1 — Reproduce el experimento tú mismo, con tus propios números. Corre el script completo de esta lección en tu máquina, y después modifícalo para usar ROW_COUNTS = [1_000_000, 50] en vez de [3_000_000, 500]. Confirma que el conjunto de valores observados sigue siendo exactamente esos dos números, nunca ningún otro.

Ver solución

El resultado debería reproducirse con cualquier par de tamaños que uses: el conjunto all_observed siempre queda limitado a los dos valores de ROW_COUNTS, sin importar cuántas rondas corras ni qué tan rápido sea tu hardware. El número total de lecturas concurrentes sí va a variar (depende de la velocidad de tu disco y tu CPU), pero la conclusión —cero estados intermedios observados— es la parte determinista del experimento.

Ejercicio 2 — Explica por qué este experimento usa dos tamaños de fila "bien distintos" en vez de, por ejemplo, 500 y 501. En 1-2 frases, justifica la elección de diseño de usar 3_000_000 y 500 en vez de dos números casi iguales.

Ver solución

Usar dos números muy distintos hace que cualquier estado intermedio sea inconfundible: si el lector alguna vez capturara, por ejemplo, 1_500_000 filas (la mitad del proceso de reemplazo), sería obvio que no es ni el estado "antes" ni el "después". Con 500 y 501, un estado intermedio real podría, por pura casualidad, coincidir con uno de los dos valores esperados y pasar inadvertido, debilitando la fuerza de la evidencia.

Ejercicio 3 — Relee la cita de data-engineering-foundations-guide y reescribe, en tus propias palabras, la diferencia central con esta lección. Sin copiar ninguna frase de esta lección, explica en 2-3 frases qué tienen en común overwrite-partition y table.overwrite() de Iceberg (el objetivo: reemplazar datos existentes por datos nuevos), y qué los separa (dónde vive la ventana de riesgo, si es que existe).

Ver solución

Ambos patrones persiguen el mismo objetivo: reemplazar el contenido existente de una porción de una tabla por contenido nuevo, de forma segura de re-correr. La diferencia es dónde ocurre el "intercambio": overwrite-partition ejecuta el DELETE y el INSERT como dos sentencias SQL separadas, así que cualquier lector que consulte la tabla justo entre ambas ve la partición vacía — una ventana de riesgo real, que la propia guía de foundations reconoció con honestidad. table.overwrite() de Iceberg prepara el resultado completo por fuera del catálogo (metadata nuevo, con todos los snapshots que haga falta) y solo después intercambia el puntero del catálogo en una única operación — nunca hay una ventana donde el catálogo apunte a un estado incompleto, sin importar cuántos pasos internos tuvo la escritura.

Resumen y siguiente paso

En esta lección construiste, con PyIceberg real, un lector concurrente que intentó sorprender a un table.overwrite() en pleno proceso —490 lecturas, durante 25 escrituras que reemplazaron millones de filas por cientos— y nunca capturó nada que no fuera el estado completo de antes o el de después. Viste por qué: Transaction acumula todos los cambios en memoria y escribe el metadata completo antes de tocar el catálogo, y el catálogo mismo confirma el cambio con un único UPDATE atómico, condicionado a que nadie más se haya adelantado. Y contrastaste esto, con la cita exacta, contra la ventana de riesgo real que data-engineering-foundations-guide reconoció en su propio overwrite-partition.

Antes de avanzar deberías poder: explicar, con tus propias palabras, por qué table.overwrite() es atómico aunque produzca más de un snapshot por dentro; y distinguir la garantía de atomicidad-por-tabla de Iceberg de una transacción multi-tabla de una base de datos relacional tradicional.

Con la atomicidad de una escritura ya demostrada, la lección 3 construye sobre el mismo mecanismo —el intercambio único del puntero del catálogo— para una operación distinta: table.update_schema(), que va a agregar la primera columna nueva de este módulo sin tocar ni un solo archivo Parquet existente.

Recursos

  • Apache Iceberg — documentación oficial, "Reliability", la fuente de las dos citas de esta lección sobre atomicidad de commits y concurrencia optimista. iceberg.apache.org/docs/latest/reliability. En inglés.
  • PyIceberg — referencia de API y código fuente, Table.transaction(), Transaction.commit_transaction(), Table._do_commit(), el mecanismo de "stage, luego commit" que esta lección cita literal. py.iceberg.apache.org/api. En inglés.
  • DISEÑO de data-engineering-foundations-guide — fuente del patrón overwrite-partition y la cita exacta sobre la ventana de riesgo, retomada palabra por palabra en esta lección. src/guides/data-engineering-foundations-guide/DISENO.md. En español.
  • DISEÑO de esta guía — la sección "Evolución de esquema" (M4), fuente exacta de la comparación entre overwrite-partition y una escritura Iceberg atómica. src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.