Módulo 6: Merge Into And Native Upserts

Proyecto: el upsert nativo de Kiosko

Descripción

Este proyecto cierra el módulo 6. Recordaste las tres formas que Kiosko ya usó para resolver el cambio de P002 (lección 2). Instalaste el runtime de Iceberg para Spark y documentaste, con evidencia real, hasta dónde llega en este entorno (lección 3). Aprendiste la sintaxis general de MERGE INTO (lección 4) y la aplicaste al caso concreto de Kiosko (lección 5). Corriste table.upsert() de verdad, sin ninguna reserva (lección 6). Y armaste el criterio para elegir entre las tres técnicas de Iceberg (lección 7). Falta un solo paso: juntar el resultado ejecutable de este módulo en un solo script, con assert automáticos — y, junto a él, la referencia completa del MERGE INTO de Spark, documentada como lo que es: representativa, no ejecutada en este entorno.

Conexión con el módulo. Este proyecto no introduce ningún concepto nuevo — es la integración final de las siete lecciones anteriores. A diferencia de los proyectos de cierre de otros módulos de esta guía, este no puede integrar las cinco técnicas en un solo script ejecutado: la incompatibilidad real documentada en la lección 3 sigue vigente. Lo que este proyecto entrega, con honestidad, es el resultado que sí corre de punta a punta —table.upsert()— y la referencia completa de lo que correría con Spark, en un entorno donde las versiones sí coincidan.

Una analogía: el resultado que sí quedó archivado, y el plano del que no

Un arquitecto que presenta un proyecto de dos edificios, uno ya construido y uno todavía en planos por un problema de permisos, no oculta el segundo — lo presenta como lo que es: un plano completo, verificado, listo para construirse en cuanto el permiso se resuelva. Este proyecto hace exactamente eso: el edificio de table.upsert() está construido, con assert que lo confirman ladrillo por ladrillo; el de MERGE INTO vía Spark es un plano completo y verificado —la misma sintaxis exacta de la lección 5—, listo para construirse el día que iceberg-spark-runtime publique una variante compatible con pyspark==4.2.0.

El material: un directorio de trabajo nuevo

kiosko_native_upsert_project/
└── kiosko_native_upsert_project.py       (este proyecto: junta las piezas ejecutables)

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_product_upsert_demo desde cero, así que no depende de ningún archivo de las lecciones anteriores de este módulo — y no toca kiosko.dim_product, kiosko.fact_orders, kiosko.dim_store ni kiosko.fact_orders_at_scale, las tablas reales que los módulos 1 a 5 ya construyeron, si corres este proyecto en el mismo directorio donde las tienes.

La solución de referencia, verificada — la parte que sí corre

# kiosko_native_upsert_project.py -- proyecto de cierre del modulo 6
# table.upsert() de punta a punta, con assert automaticos
import os

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

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


def main() -> None:
    print("=== Kiosko: el upsert nativo, de punta a punta ===\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("Paso 1/6 -- catalogo 'kiosko' listo")

    table = catalog.create_table("kiosko.dim_product_upsert_demo", schema=DIM_PRODUCT_SCHEMA)
    pa_table_v1 = pa.Table.from_pylist(DIM_PRODUCT_V1, schema=PA_SCHEMA)
    table.append(pa_table_v1)
    rows_after_v1 = table.scan().to_arrow().num_rows
    print(f"Paso 2/6 -- kiosko.dim_product_upsert_demo creada y cargada con V1: {rows_after_v1} filas "
          f"(P002 = snacks/0.60)")

    pa_table_v2 = pa.Table.from_pylist(DIM_PRODUCT_V2, schema=PA_SCHEMA)
    result = table.upsert(pa_table_v2, join_cols=["product_id"])
    print(f"Paso 3/6 -- table.upsert(V2_completo) corrido: rows_updated={result.rows_updated}, "
          f"rows_inserted={result.rows_inserted}")

    final_rows = sorted(table.scan().to_arrow().to_pylist(), key=lambda r: r["product_id"])
    p002 = next(r for r in final_rows if r["product_id"] == "P002")
    print(f"Paso 4/6 -- estado final: {len(final_rows)} filas, "
          f"P002 = {p002['category']}/{p002['unit_cost']}")

    history_len = len(table.history())
    ops = [row["operation"] for row in table.inspect.snapshots().select(["operation"]).to_pylist()]
    print(f"Paso 5/6 -- table.history() tiene {history_len} entradas, operaciones: {ops}")

    others_unchanged = all(
        r["category"] == next(o for o in DIM_PRODUCT_V1 if o["product_id"] == r["product_id"])["category"]
        and r["unit_cost"] == next(o for o in DIM_PRODUCT_V1 if o["product_id"] == r["product_id"])["unit_cost"]
        for r in final_rows if r["product_id"] != "P002"
    )
    print(f"Paso 6/6 -- P001, P003, P004 sin cambios: {others_unchanged}\n")

    print("=== Verificacion final ===\n")

    assert rows_after_v1 == 4
    assert result.rows_updated == 1, "el upsert debe actualizar exactamente 1 fila (P002)"
    assert result.rows_inserted == 0, "el upsert no debe insertar ninguna fila nueva"
    assert len(final_rows) == 4, "el grano debe seguir siendo una fila por producto"
    assert p002["category"] == "health-snacks" and p002["unit_cost"] == 0.68
    assert others_unchanged, "P001, P003 y P004 no deben haber cambiado"
    assert history_len == 3, "append(V1) + upsert (overwrite parcial + append) = 3 entradas"
    assert ops == ["append", "overwrite", "append"], "el upsert debe resolverse como overwrite+append"

    print("Todas las verificaciones pasaron:")
    print("  - kiosko.dim_product_upsert_demo: 4 filas, una por producto, sin columnas de historia")
    print("  - table.upsert(V2_completo, join_cols=['product_id']) detecto solo el cambio real de P002")
    print("  - rows_updated=1, rows_inserted=0 -- sin que nadie calculara el delta a mano")
    print("  - P001, P003, P004 confirmados sin cambios")
    print("  - table.history() confirma 3 snapshots: append + overwrite (parcial) + append")


if __name__ == "__main__":
    main()

Qué esperar (verificado corriendo python3 kiosko_native_upsert_project.py real, de punta a punta, en un directorio nuevo):

=== Kiosko: el upsert nativo, de punta a punta ===

Paso 1/6 -- catalogo 'kiosko' listo
Paso 2/6 -- kiosko.dim_product_upsert_demo creada y cargada con V1: 4 filas (P002 = snacks/0.60)
Paso 3/6 -- table.upsert(V2_completo) corrido: rows_updated=1, rows_inserted=0
Paso 4/6 -- estado final: 4 filas, P002 = health-snacks/0.68
Paso 5/6 -- table.history() tiene 3 entradas, operaciones: ['append', 'overwrite', 'append']
Paso 6/6 -- P001, P003, P004 sin cambios: True

=== Verificacion final ===

Todas las verificaciones pasaron:
  - kiosko.dim_product_upsert_demo: 4 filas, una por producto, sin columnas de historia
  - table.upsert(V2_completo, join_cols=['product_id']) detecto solo el cambio real de P002
  - rows_updated=1, rows_inserted=0 -- sin que nadie calculara el delta a mano
  - P001, P003, P004 confirmados sin cambios
  - table.history() confirma 3 snapshots: append + overwrite (parcial) + append

Ocho assert, ninguno decorativo: confirman que el grano se mantuvo (4 filas, nunca más), que upsert() detectó exactamente el cambio real (rows_updated=1, rows_inserted=0), que las tres filas sin cambios efectivamente no cambiaron, y que el mecanismo interno se resolvió exactamente como explicó la lección 6 (append + overwrite + append).

La referencia completa — la parte representativa: MERGE INTO vía Spark

Para que este proyecto documente el resultado completo del módulo, y no solo la mitad que corrió en este entorno, acá está el MERGE INTO de la lección 5, íntegro, marcado explícitamente como lo que es:

-- Qué esperar (representativo) -- sintaxis identica a la leccion 5, no ejecutada
-- en este entorno por la incompatibilidad real documentada en la leccion 3.

CREATE NAMESPACE IF NOT EXISTS local.kiosko;

CREATE TABLE local.kiosko.dim_product (
    product_id STRING, product_name STRING, category STRING, unit_cost DOUBLE
) USING iceberg;

INSERT INTO local.kiosko.dim_product VALUES
    ('P001', 'Bottled Water 600ml',   'beverages',   0.40),
    ('P002', 'Energy Bar',            'snacks',      0.60),
    ('P003', 'Instant Coffee Sachet', 'beverages',   0.35),
    ('P004', 'Phone Charger Cable',   'electronics', 2.10);

CREATE TABLE local.kiosko.dim_product_staging (
    product_id STRING, product_name STRING, category STRING, unit_cost DOUBLE
) USING iceberg;

INSERT INTO local.kiosko.dim_product_staging VALUES
    ('P002', 'Energy Bar', 'health-snacks', 0.68);

MERGE INTO local.kiosko.dim_product AS target
USING local.kiosko.dim_product_staging AS source
ON target.product_id = source.product_id
WHEN MATCHED THEN UPDATE SET
    target.product_name = source.product_name,
    target.category     = source.category,
    target.unit_cost    = source.unit_cost
WHEN NOT MATCHED THEN INSERT (product_id, product_name, category, unit_cost)
VALUES (source.product_id, source.product_name, source.category, source.unit_cost);

Qué esperar (representativo): exactamente el mismo resultado de negocio que confirmó el script de upsert() de arriba — cuatro filas, P002 con category='health-snacks', unit_cost=0.68, P001/P003/P004 intactos. La diferencia, ya documentada a fondo en las lecciones 4, 5 y 7 de este módulo: este SQL corre distribuido, sobre Spark, y necesita la JVM; el script de arriba corre en un solo proceso Python, sin ninguna de las dos.

Diagrama: de dónde venías, a dónde llegaste

flowchart LR
    A["Modulo 1-5:\nkiosko.dim_product en V2\nvia overwrite() + time travel"] --> B["Leccion 2:\n3 tecnicas recordadas,\nDuckDB, dbt, esta guia"]
    B --> C["Leccion 3:\nSpark + runtime instalado,\nincompatibilidad real documentada"]
    C --> D["Leccion 4-5:\nMERGE INTO -- sintaxis real,\nsalida representativa"]
    D --> E["Leccion 6:\ntable.upsert() -- ejecutado\nde verdad, rows_updated=1"]
    E --> F["Leccion 7:\ncriterio de eleccion,\n5 factores"]
    F --> G["Este proyecto:\nupsert ejecutado + assert,\nMERGE documentado junto a el"]
    G --> H["Modulo 7:\ncatalogos, mantenimiento,\nDelta Lake por contraste"]

Cerrando la promesa del módulo, punto por punto

Lo que la lección 1 prometióEvidencia de que este módulo lo entregó
Recordar las tres formas que Kiosko ya resolvió el cambio de P002Lección 2: código y salida citados literal de data-modeling, dbt, y el módulo 3
MERGE INTO nativo de Iceberg vía Spark SQLLecciones 3-5: entorno instalado de verdad, sintaxis verificada contra la documentación oficial, incompatibilidad real documentada con evidencia
table.upsert() como alternativa 100% PythonLección 6 y este proyecto: UpsertResult(rows_updated=1, rows_inserted=0), ejecutado de verdad, sin reserva
Criterio para elegir entre las técnicasLección 7: cinco factores, tabla comparativa, árbol de decisión
Los dos catálogos (kiosko/local) declarados como distintos, sin afirmar interoperabilidadDeclarado explícito desde la lección 1, respetado en cada lección técnica de este módulo

Este proyecto no tocó kiosko.dim_product, kiosko.fact_orders, kiosko.dim_store ni kiosko.fact_orders_at_scale — esas tablas siguen exactamente como las dejaron los módulos 1 a 5. Lo que este proyecto entrega es exactamente lo que este módulo prometió: dos formas nuevas de aplicar el cambio de P002 sin reconstruir la tabla completa, una ejecutada de punta a punta con verificación automática, la otra documentada con la misma honestidad que el resto de esta guía aplica a cualquier resultado que no pudo ejecutar de verdad.

Errores comunes

Correr este proyecto sobre un catálogo que ya tiene kiosko.dim_product_upsert_demo de la lección 6. Qué pasa: alguien corre este proyecto en el mismo directorio donde ya completó la lección 6, 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 dejó la lección 6. Cómo detectarlo: si ves TableAlreadyExistsError al correr kiosko_native_upsert_project.py, ya tienes un catálogo con esa tabla registrada en el mismo directorio. Cómo corregirlo: corre este proyecto en un directorio de trabajo nuevo, separado de donde hiciste la lección 6 — tal como sugiere "El material" de esta lección.

Intentar correr el bloque SQL de MERGE INTO de esta lección, esperando que funcione. Qué pasa: alguien copia el SQL de la sección "La referencia completa" de esta lección a un spark-sql real, esperando verlo correr sin problemas. Por qué pasa: el SQL está escrito de forma completa y correcta —no es pseudocódigo—, así que parece razonable que corra. Cómo detectarlo: si tu entorno tiene exactamente la misma combinación de versiones que esta guía (pyspark==4.2.0 con iceberg-spark-runtime-4.0_2.13:1.11.0), vas a encontrar el mismo IncompatibleClassChangeError documentado en la lección 3. Cómo corregirlo: si tienes acceso a un entorno con una combinación de versiones compatible —por ejemplo, pyspark==4.0.x con este mismo runtime, ver el Ejercicio 3 de la lección 3—, este SQL debería correr sin cambios. Si no, trátalo como lo que esta lección declara que es: una referencia verificada contra la documentación oficial, no una demostración ejecutada en este entorno.

Ejercicios

Ejercicio 1 — Corre el proyecto completo tú mismo, desde cero. En un directorio nuevo, corre python3 kiosko_native_upsert_project.py. Confirma que ves los seis pasos completarse y el mensaje final con las cinco verificaciones.

Ver solución

Si PyIceberg está instalado en tu entorno, la salida debería reproducir exactamente la estructura de esta lección: seis pasos numerados, seguidos de la verificación final con los cinco mensajes de éxito. Todo el script corre en menos de un segundo — no hay ningún dataset a escala en este proyecto, a diferencia del proyecto del módulo 5.

Ejercicio 2 — Rompe un assert a propósito, y observa el fallo. Cambia temporalmente DIM_PRODUCT_V2 para que P001 también tenga un unit_cost distinto (por ejemplo, 0.45 en vez de 0.40), corre el script de nuevo, y observa qué assert falla primero. Después revierte el cambio.

Ver solución

El primer assert en fallar debería ser assert result.rows_updated == 1, "el upsert debe actualizar exactamente 1 fila (P002)" — porque ahora hay dos filas con un valor distinto (P001 y P002), así que result.rows_updated sería 2, no 1. Este ejercicio confirma, con evidencia directa, que table.upsert() realmente compara cada fila de forma independiente —no es una casualidad que solo P002 se contara como actualizada en la corrida original, es el resultado directo de que solo P002 tenía un valor distinto entre V1 y V2—.

Ejercicio 3 — Explica, en tus propias palabras, por qué este proyecto documenta el MERGE INTO de Spark en vez de simplemente omitirlo. En 3-4 frases, justifica por qué incluir código que no pudo ejecutarse de verdad en este entorno es coherente con la disciplina del resto de esta guía.

Ver solución

Omitir el MERGE INTO de Spark dejaría al proyecto incompleto frente a lo que el módulo 1 prometió: dos técnicas nuevas, no una. La disciplina de esta guía —declarada desde el DISEÑO— nunca fue "solo mostrar lo que funciona sin fricción", sino "mostrar exactamente qué se verificó de verdad y qué no, sin ambigüedad" — la misma regla que ya aplicaron los bloques "Qué esperar" con snapshot_id no hardcodeado, o la advertencia explícita sobre los dos catálogos distintos. Documentar el MERGE INTO como referencia verificada contra la documentación oficial, con la incompatibilidad real explicada en la lección 3, es más honesto y más útil que fingir una ejecución que no ocurrió, o que borrar la mitad del módulo porque una pieza de infraestructura externa no cooperó.

Resumen y siguiente paso: el cierre de este módulo

Con este proyecto cierras el módulo 6. Integraste el resultado ejecutable —table.upsert(), con assert que confirman cada afirmación— junto a la referencia completa y honesta de lo que MERGE INTO vía Spark habría producido, en un entorno con versiones compatibles. Kiosko tiene, ahora, dos formas más de aplicar un cambio parcial a una tabla Iceberg sin reconstruir el estado completo — sumadas a las tres que ya conocía desde data-modeling, dbt, y el módulo 3 de esta misma guía.

Hacia dónde sigues. El módulo 7 —Catálogos, mantenimiento y Delta Lake por contraste— nombra los catálogos de producción que esta guía nunca implementó (REST, AWS Glue Catalog, Unity Catalog, Polaris), y vuelve a kiosko.dim_product —la tabla real, con sus snapshots acumulados desde el módulo 3— para podarla con seguridad: expirar snapshots viejos, compactar archivos pequeños, sin perder el time travel que sí necesitas. Y cierra con Delta Lake, nombrado una sola vez, por contraste.

Recursos

  • PyIceberg — documentación oficial (quickstart), el flujo completo de catálogo, tabla, append() y upsert() que integra este proyecto. py.iceberg.apache.org. En inglés.
  • PyIceberg — referencia de API, table.upsert(), UpsertResult, table.inspect.snapshots(). py.iceberg.apache.org/api. En inglés.
  • Apache Iceberg — documentación oficial, "Spark Writes", sección MERGE INTO, fuente del bloque SQL de referencia de este proyecto. iceberg.apache.org/docs/latest/spark-writes/#merge-into. En inglés.
  • GitHub — apache/iceberg#15238, la incompatibilidad real documentada que explica por qué la sección SQL de este proyecto es representativa. github.com/apache/iceberg/issues/15238. En inglés.
  • DISEÑO de esta guía — el mapa completo de los ocho módulos, incluido el módulo 7 que sigue. src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.