Módulo 6: Merge Into And Native Upserts
El upsert de PyIceberg: la alternativa Python-nativa
Descripción
Esta lección corre de verdad, sin ninguna reserva, lo que las lecciones 4 y 5 no pudieron ejecutar en este entorno: la misma idea de MERGE INTO —actualizar si existe, insertar si no—, pero como un único método de Python, table.upsert(), sin una sola línea de SQL y sin necesitar la JVM. Vas a crear una tabla nueva, cargarla con V1, y confirmar, con salida literal, que table.upsert() detecta el cambio real de P002 —y solo ese— entre cuatro filas que le entregas completas.
Conexión con el módulo. La lección 3 documentó una incompatibilidad real entre iceberg-spark-runtime-4.0 y Spark 4.2.0, que dejó las lecciones 4 y 5 marcadas como representativas. table.upsert() no depende de Spark en absoluto —vive enteramente en PyIceberg, el mismo paquete que ya usaste sin interrupciones en los módulos 1 a 5—, así que esta lección vuelve a la disciplina de esta guía: código ejecutado de verdad, salida copiada literal.
Por qué esta lección no toca kiosko.dim_product
Antes del código, una aclaración necesaria. kiosko.dim_product —la tabla real, la que construiste en el módulo 3— ya vive en su estado V2 desde hace tres módulos: table.overwrite() ya aplicó el cambio, y snap_v1 sigue disponible para recuperar el estado anterior por time travel. Repetir el mismo cambio sobre esa misma tabla no demostraría nada nuevo —upsert() compararía V2 contra V2, y no encontraría ninguna diferencia que actualizar—.
Para comparar upsert() con las otras cuatro técnicas de este módulo, en igualdad de condiciones —partiendo de V1, igual que hicieron la lección 5 (con local.kiosko.dim_product, en un catálogo aparte) y las técnicas de DuckDB y dbt (lección 2)—, esta lección reproduce el experimento en una tabla nueva y dedicada: kiosko.dim_product_upsert_demo. Mismo catálogo kiosko de siempre, tabla nueva, para no tocar la historia real que ya construiste. Es la misma disciplina que ya viste en los proyectos de cierre de cada módulo: un directorio de trabajo nuevo, un estado reproducible desde cero.
Una analogía: el mismo mostrador, ahora en modo autoservicio
El mostrador del banco de la lección 1 seguía necesitando que alguien lo atendiera con SQL. table.upsert() es la misma ventanilla, convertida en un kiosco de autoservicio: le entregas el estado completo que crees correcto —las cuatro filas, tal como las conoces hoy—, y el sistema decide, por su cuenta, cuáles de esas filas representan un cambio real y cuáles ya estaban al día, sin que tú tengas que calcular el delta de antemano ni escribir un ON/WHEN MATCHED.
Ejemplo trabajado: upsert(), ejecutado de verdad
Paso 1 — Crea la tabla de demostración, cargada con V1
# upsert_demo.py
import os
import pyarrow as pa
from pyiceberg.catalog import load_catalog
from pyiceberg.schema import Schema
from pyiceberg.types import DoubleType, NestedField, StringType
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}",
)
catalog.create_namespace("kiosko")
dim_product_schema = Schema(
NestedField(field_id=1, name="product_id", field_type=StringType(), required=True),
NestedField(field_id=2, name="product_name", field_type=StringType(), required=True),
NestedField(field_id=3, name="category", field_type=StringType(), required=True),
NestedField(field_id=4, name="unit_cost", field_type=DoubleType(), required=True),
)
table = catalog.create_table("kiosko.dim_product_upsert_demo", schema=dim_product_schema)
dim_product_pa_schema = pa.schema([
pa.field("product_id", pa.string(), nullable=False),
pa.field("product_name", pa.string(), nullable=False),
pa.field("category", pa.string(), nullable=False),
pa.field("unit_cost", pa.float64(), nullable=False),
])
DIM_PRODUCT_V1 = [
{"product_id": "P001", "product_name": "Bottled Water 600ml", "category": "beverages", "unit_cost": 0.40},
{"product_id": "P002", "product_name": "Energy Bar", "category": "snacks", "unit_cost": 0.60},
{"product_id": "P003", "product_name": "Instant Coffee Sachet", "category": "beverages", "unit_cost": 0.35},
{"product_id": "P004", "product_name": "Phone Charger Cable", "category": "electronics", "unit_cost": 2.10},
]
pa_table_v1 = pa.Table.from_pylist(DIM_PRODUCT_V1, schema=dim_product_pa_schema)
table.append(pa_table_v1)
print("V1 cargado. Filas:", table.scan().to_arrow().num_rows)
Qué esperar (verificado corriendo el script real):
V1 cargado. Filas: 4
Nada nuevo hasta acá — el mismo create_table() + append() que ya usaste decenas de veces desde el módulo 1.
Paso 2 — table.upsert(), con el estado completo V2
DIM_PRODUCT_V2 = [
{"product_id": "P001", "product_name": "Bottled Water 600ml", "category": "beverages", "unit_cost": 0.40},
{"product_id": "P002", "product_name": "Energy Bar", "category": "health-snacks", "unit_cost": 0.68},
{"product_id": "P003", "product_name": "Instant Coffee Sachet", "category": "beverages", "unit_cost": 0.35},
{"product_id": "P004", "product_name": "Phone Charger Cable", "category": "electronics", "unit_cost": 2.10},
]
pa_table_v2 = pa.Table.from_pylist(DIM_PRODUCT_V2, schema=dim_product_pa_schema)
result = table.upsert(pa_table_v2, join_cols=["product_id"])
print("UpsertResult:", result)
print("rows_updated:", result.rows_updated)
print("rows_inserted:", result.rows_inserted)
Fíjate en algo distinto de la lección 5: acá le entregas a upsert() las cuatro filas completas —no solo el delta—, exactamente como hicieron data-modeling y dbt en la lección 2. La diferencia es que upsert(), a diferencia de MERGE INTO tal como lo escribiste en la lección 5, sí compara internamente cada valor no-clave contra lo que ya existe, y solo cuenta como "actualizada" la fila donde algo cambió de verdad.
Qué esperar (verificado corriendo el script real):
UpsertResult: UpsertResult(rows_updated=1, rows_inserted=0)
rows_updated: 1
rows_inserted: 0
Exactamente lo que predijiste en el Ejercicio 3 de la lección 2: una fila actualizada —P002, la única que cambió—, cero filas insertadas —ningún product_id nuevo—. P001, P003 y P004 estaban presentes en las cuatro filas que le entregaste, pero upsert() los reconoció como idénticos a lo ya archivado, y no los tocó.
Paso 3 — Confirma el estado final
print("\ntable.scan().to_arrow() tras el upsert (ordenado por product_id):")
rows = sorted(table.scan().to_arrow().to_pylist(), key=lambda r: r["product_id"])
for row in rows:
print(f" {row['product_id']} {row['product_name']:<22} {row['category']:<14} unit_cost={row['unit_cost']}")
Qué esperar (verificado corriendo el script real):
table.scan().to_arrow() tras el upsert (ordenado por product_id):
P001 Bottled Water 600ml beverages unit_cost=0.4
P002 Energy Bar health-snacks unit_cost=0.68
P003 Instant Coffee Sachet beverages unit_cost=0.35
P004 Phone Charger Cable electronics unit_cost=2.1
Nota sobre el orden. Si corres table.scan().to_arrow().to_pylist() sin ordenar explícitamente, la fila de P002 aparece primero, no en el orden P001..P004 que viste en el módulo 3 — porque upsert() escribió un archivo de datos nuevo, separado, solo con la fila actualizada, y table.scan() no garantiza ningún orden entre archivos. Esta guía nunca asumió un orden implícito en ninguna lección anterior por casualidad: table.scan() refleja el orden físico en que Iceberg lee sus archivos de datos, nunca un ORDER BY — si tu código depende de un orden específico, ordénalo tú mismo, exactamente como hace sorted(...) en este paso.
Paso 4 — Lo que reveló table.history(): dos snapshots, no uno
print("\ntable.history() tiene", len(table.history()), "entradas")
snaps = table.inspect.snapshots().select(["operation", "summary"])
for row in snaps.to_pylist():
s = dict(row["summary"])
print(f" operation={row['operation']:<9} added-records={s.get('added-records','-')} "
f"deleted-records={s.get('deleted-records','-')} total-records={s.get('total-records','-')}")
Qué esperar (verificado corriendo el script real; snapshot_id/committed_at son de tu propia corrida, distintos cada vez):
table.history() tiene 3 entradas
operation=append added-records=4 deleted-records=- total-records=4
operation=overwrite added-records=3 deleted-records=4 total-records=-
operation=append added-records=1 deleted-records=- total-records=4
Tres entradas: el append() del Paso 1 (las cuatro filas de V1), y dos más, no una, producidas por la única llamada a upsert() del Paso 2. Esto no debería sorprenderte del todo — es la misma lección que ya aprendiste en el módulo 3, lección 3, ahora aplicada a un mecanismo distinto: una operación de alto nivel de PyIceberg puede resolverse internamente como varios snapshots.
Profundización: qué hace upsert() por dentro, con evidencia del propio código fuente
Vale la pena entender por qué aparecen exactamente esas dos entradas, no una ni tres. La implementación de table.upsert() en PyIceberg 0.11.1 —código abierto, verificable— hace, en esencia, dos pasos:
- Identifica qué filas del
sourcecoinciden con una fila deltargety tienen, al menos, un valor no-clave distinto —la misma comparación fila por fila que ya viste en el Paso 2, la razón por la querows_updated=1y no4—. Con esas filas —soloP002, en este caso—, llama internamente atable.overwrite(rows_to_update, overwrite_filter=...), con un filtro que retira únicamente las filas que coinciden porproduct_id. Como esas filas conviven, en el mismo archivo Parquet, con las tres que no cambiaron, Iceberg tiene que reescribir ese archivo — retira el archivo viejo completo (deleted-records=4) y escribe uno nuevo, solo con las tres filas que no iban a actualizarse (added-records=3). Esta es la entradaoperation=overwritede la salida de arriba. - Agrega las filas actualizadas como un
append()separado — la fila nueva deP002, con sus valoreshealth-snacks/0.68(added-records=1). Esta es la segunda entradaoperation=append.
Dos snapshots, uno para "lo que sigue igual, reacomodado" y otro para "lo que cambió, agregado" — un mecanismo distinto del overwrite() completo del módulo 3 (que hacía delete de todo seguido de append de todo), pero con la misma lección de fondo: cuenta siempre los snapshots con table.history(), nunca asumas que "una llamada, un snapshot".
Diagrama: dos filas quietas, una que se mueve
flowchart TD
A["table.upsert(V2_completo, join_cols=['product_id'])"] --> B["Paso interno 1:\ncompara V2 contra target\nsolo P002 tiene un valor distinto"]
B --> C["overwrite(solo P002 nuevo,\nfiltro=product_id=='P002')\nreescribe el archivo: quedan P001,P003,P004"]
B --> D["append(P002 health-snacks/0.68)\nfila nueva, archivo nuevo"]
C --> E["UpsertResult\nrows_updated=1, rows_inserted=0"]
D --> E
Errores comunes
Llamar a table.upsert() sin join_cols, y esperar que PyIceberg adivine la llave. Qué pasa: alguien corre table.upsert(pa_table_v2), sin el argumento join_cols, confiando en que PyIceberg va a usar product_id automáticamente porque es "obviamente" la llave. Verifícalo tú mismo:
table.upsert(pa_table_v2) # sin join_cols
ValueError: Join columns could not be found, please set identifier-field-ids or pass in explicitly.
Por qué pasa: upsert() sí puede inferir la columna de unión automáticamente, pero solo si el esquema de la tabla declara explícitamente identifier-field-ids —una marca formal de "esta columna identifica de forma única cada fila"—, algo que el esquema de dim_product de esta guía nunca declaró (ni en el módulo 3, ni acá). Cómo detectarlo: el mensaje de error es explícito y dice exactamente qué falta — join_cols o identifier-field-ids. Cómo corregirlo: pasa join_cols=["product_id"] de forma explícita, como hace esta lección — es más seguro que depender de una configuración de esquema que tendrías que recordar declarar en el momento de crear la tabla.
Asumir que upsert() con el estado completo (V2, cuatro filas) es equivalente a table.overwrite() del módulo 3. Qué pasa: alguien ve que esta lección le pasa las cuatro filas a upsert(), igual que overwrite() recibía las cuatro filas completas en el módulo 3, y concluye que son la misma operación con otro nombre. Por qué pasa: ambas reciben el mismo tipo de entrada —el estado completo—. Cómo detectarlo: cuenta los snapshots — overwrite() del módulo 3 produjo delete + append (retira todo, agrega todo); upsert() de esta lección produjo overwrite (parcial, solo reescribe el archivo que contenía la fila que cambió) + append (solo la fila nueva). Cómo corregirlo: upsert() siempre hace, primero, el trabajo de comparar valor por valor —el mismo trabajo que el WHEN MATCHED AND (...) explícito tuvo que hacer a mano en el MERGE de DuckDB (lección 2)—, y solo toca lo que cambió de verdad; overwrite() sin filtro nunca compara nada, simplemente reemplaza el 100% de la tabla, haya cambiado algo o no.
Ejercicios
Ejercicio 1 — Reproduce el experimento completo tú mismo, y confirma los tres números. Con PyIceberg instalado, corre los cuatro pasos de esta lección en un directorio nuevo. Confirma rows_updated=1, rows_inserted=0, y len(table.history()) == 3.
Ver solución
Si tu entorno tiene PyIceberg 0.11.1 instalado (pip install "pyiceberg[sql-sqlite,pyarrow]"), tu salida debería coincidir exactamente con la de esta lección en los tres números —rows_updated=1, rows_inserted=0, tres entradas en table.history()—. Los snapshot_id van a ser distintos de los mostrados acá — eso es exactamente lo esperado.
Ejercicio 2 — Corre upsert() una segunda vez, con el mismo V2, y predice el resultado antes de correrlo. Sin cambiar nada, llama table.upsert(pa_table_v2, join_cols=["product_id"]) de nuevo, sobre la tabla que dejó esta lección.
Ver solución
result2 = table.upsert(pa_table_v2, join_cols=["product_id"])
print("UpsertResult (segunda corrida):", result2)
Salida esperada: UpsertResult(rows_updated=0, rows_inserted=0). Ninguna fila cambia, porque las cuatro filas de V2 ya son idénticas a lo que la tabla tiene archivado — la misma propiedad de idempotencia que ya viste en el MERGE de DuckDB (lección 2) y en dbt snapshot, ahora confirmada también para upsert(). table.history() seguiría en tres entradas —una corrida sin cambios reales no agrega ningún snapshot nuevo, porque get_rows_to_update() no encuentra ninguna fila que reescribir y rows_to_insert queda vacío.
Ejercicio 3 — Explica, en tus propias palabras, por qué upsert() necesitó reescribir un archivo completo (deleted-records=4) para actualizar una sola fila. En 2-3 frases, usando lo que sabes sobre archivos Parquet desde el módulo 2, explica por qué "actualizar una fila" en Iceberg casi nunca significa "tocar solo esa fila" a nivel de archivo físico.
Ver solución
Parquet es un formato de archivo inmutable — una vez escrito, ningún proceso puede modificar una fila específica dentro de un archivo .parquet existente sin reescribirlo completo. Como las cuatro filas de V1 viven juntas en un único archivo de datos (las cuatro se cargaron con un solo append() en el Paso 1), actualizar el valor de P002 obliga a Iceberg a retirar la referencia a ese archivo completo (deleted-records=4) y escribir uno nuevo con las filas que sí sobreviven sin cambios (added-records=3), mientras la fila actualizada se archiva por separado. Es el mismo principio que ya viste en el módulo 2 —los archivos de datos nunca se modifican en el lugar, solo se reemplazan por completo— aplicado acá a una actualización de una sola fila.
Resumen y siguiente paso
En esta lección corriste table.upsert() de verdad, sin ninguna reserva: creaste kiosko.dim_product_upsert_demo desde V1, y confirmaste que upsert(V2_completo, join_cols=["product_id"]) detecta, por sí solo, que solo P002 cambió —UpsertResult(rows_updated=1, rows_inserted=0)—, sin que tú tuvieras que calcular ningún delta de antemano. Viste, con table.history(), que esa única llamada produjo dos snapshots internos —un overwrite parcial y un append—, y por qué, leyendo el propio código fuente de PyIceberg.
Antes de avanzar deberías poder: explicar por qué esta lección usa una tabla nueva en vez de kiosko.dim_product; reproducir el experimento completo con los tres números exactos; y explicar la diferencia entre upsert() (compara y solo toca lo que cambió) y MERGE INTO tal como lo escribiste en la lección 5 (actualiza cualquier coincidencia, sin comparar valores, a menos que tú agregues esa condición a mano).
La lección 7 da el criterio final: cuándo elegir MERGE INTO en SQL, cuándo elegir upsert() en Python — y cuándo ninguno de los dos, y sigues necesitando table.overwrite() del módulo 3.
Recursos
- PyIceberg — referencia de API,
table.upsert(),UpsertResult, y el argumentojoin_cols. py.iceberg.apache.org/api. En inglés. - PyIceberg — PyPI, versión vigente 0.11.1, la misma versión que instalaste en el módulo 1 y que corre esta lección sin cambios. pypi.org/project/pyiceberg. En inglés.
- Esta misma guía, módulo 3, lección 3 — fuente del hallazgo original de que una operación puede producir más de un snapshot, retomado aquí para
upsert().03-changing-p002-with-a-plain-overwrite.md. En español. - DISEÑO de esta guía — el mapa completo de los ocho módulos, incluida la sección del módulo 6, y la regla dura de nunca hardcodear un
snapshot_id.src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.