Módulo 6: Merge Into And Native Upserts

Presentación del módulo: `MERGE INTO` y upserts nativos

Por qué existe este módulo

El módulo 3 de esta guía resolvió el cambio de P002 con una técnica que ningún módulo anterior del ecosistema tenía disponible: table.overwrite() sobre kiosko.dim_product —sin valid_from, sin valid_to, sin is_current— y table.scan(snapshot_id=snap_v1) para recuperar el estado anterior. Funcionó, con evidencia ejecutada: el margen correcto de P002 (10.8) salió exactamente igual que en data-modeling-for-analytics-guide y dbt-analytics-engineering-guide, sin que nadie diseñara una sola columna de historia.

Pero fíjate en algo que el módulo 3 dejó, a propósito, sin resolver: table.overwrite() reemplaza toda la tabla. Funciona perfecto cuando tienes, en Python, el estado completo de los cuatro productos —DIM_PRODUCT_V2, escrito a mano, con los cuatro dict completos—. Pero en un pipeline real, lo que normalmente llega no es "el estado completo de la dimensión", sino un cambio: una fila de un feed CDC, una línea de un archivo exportado hoy por el sistema de catálogo de Kiosko, con el precio nuevo de un solo producto. Si tu única herramienta es overwrite(), tienes que reconstruir las cuatro filas completas cada vez —incluidas las tres que no cambiaron— solo para poder escribir la que sí. Ese es, con precisión, el problema que este módulo resuelve: aplicar un cambio parcial, sin reconstruir la tabla entera, y sin perder el grano de una fila por producto.

No es un problema nuevo para el ecosistema. Ya lo resolviste, dos veces, con dos herramientas completamente distintas: data-modeling-for-analytics-guide (módulo 4) escribió un MERGE INTO a mano sobre DuckDB, comparando una tabla de staging contra dim_product_scd en un solo statement. dbt-analytics-engineering-guide (módulo 5) automatizó exactamente esa misma idea con dbt snapshot, sin que nadie escribiera un UPDATE ni un INSERT. Y esta misma guía, en su módulo 3, resolvió el caso favorable de Kiosko con una tercera técnica todavía: time travel puro, cero columnas. Este módulo pone las tres, lado a lado, con código citado literal —sin volver a correr ninguna—, y agrega una cuarta: MERGE INTO, ahora nativo del propio formato de tabla, corrido con SQL real sobre Iceberg vía Spark. Y una quinta, hermana de la cuarta: table.upsert() de PyIceberg, la misma idea sin una sola línea de SQL.

El caso que nos acompaña: el mismo cambio de P002, una cuarta y una quinta vez

Este módulo no inventa un caso nuevo. Es, otra vez, P002 Energy Bar: category='snacks', unit_cost=0.60 pasa a category='health-snacks', unit_cost=0.68, vigente desde el 2026-08-15 — el mismo cambio, con los mismos dos números (margin=10.8 correcto, margin=9.36 roto) que ya viste tres veces en este ecosistema. Lo que cambia aquí es el vehículo: en vez de reescribir la tabla completa (módulo 3) o de un UPDATE/INSERT manual sobre DuckDB (data-modeling) o de un comando declarativo de dbt (dbt), este módulo aplica el cambio como lo que realmente es en cualquier pipeline de producción — un delta, una sola fila con el valor nuevo de P002, fusionada contra el estado existente en un solo comando transaccional.

Dos piezas nuevas de vocabulario, ambas en inglés, se suman al léxico de Iceberg que ya conoces: MERGE INTO, el statement de SQL que Spark hereda casi textual de cualquier motor relacional moderno, adaptado para escribir sobre una tabla Iceberg; y upsert, el nombre genérico —de la industria completa, no solo de Iceberg— para la operación "actualiza si existe, inserta si no existe", que PyIceberg expone como un único método de Python: table.upsert().

Una analogía: el mostrador del banco que hace un solo trámite

Imagina que tienes que actualizar un dato en tu cuenta bancaria —digamos, tu número de teléfono— y no sabes si el banco ya tiene un registro tuyo o si eres cliente nuevo. La forma vieja de resolver esto, en un banco mal diseñado, son dos ventanillas separadas: una para "abrir cuenta nueva" y otra para "actualizar cuenta existente", y tú tienes que saber, de antemano, cuál te corresponde —y si te equivocas de ventanilla, el trámite falla o, peor, crea una cuenta duplicada—.

MERGE INTO y upsert son el mostrador rediseñado: una sola ventanilla. Le entregas tu documento —tu product_id, en el caso de Kiosko—, y el sistema decide por sí mismo, en el mismo trámite, si tiene que actualizar un registro que ya existe o crear uno nuevo. Tú nunca declaras cuál de los dos casos es el tuyo; el mostrador lo resuelve comparando lo que le entregaste contra lo que ya tiene archivado. WHEN MATCHED es "ya eres cliente, actualizo tu registro"; WHEN NOT MATCHED es "eres nuevo, te doy de alta" — las dos ramas de la misma transacción, nunca dos trámites separados que alguien tenga que secuenciar a mano.

Diagrama: cinco vehículos, un mismo destino

flowchart TB
    P["El cambio de P002:\nsnacks/0.60 -> health-snacks/0.68"]

    P --> A["1. data-modeling (M4):\nMERGE INTO a mano sobre DuckDB\nvalid_from/valid_to/is_current"]
    P --> B["2. dbt (M5):\ndbt snapshot\ndbt_valid_from/dbt_valid_to/dbt_scd_id"]
    P --> C["3. esta guia (M3):\ntable.overwrite() + time travel\nCERO columnas de historia"]
    P --> D["4. esta guia (M6, Spark):\nMERGE INTO nativo de Iceberg\nSQL real via iceberg-spark-runtime"]
    P --> E["5. esta guia (M6, PyIceberg):\ntable.upsert()\nPython puro, sin SQL"]

    A -.->|"recordado, NO re-ejecutado"| L2["Leccion 2 de este modulo"]
    B -.->|"recordado, NO re-ejecutado"| L2
    C -.->|"recordado, NO re-ejecutado"| L2
    D -.->|"ejecutado / representativo -- L3 a L5"| L2
    E -.->|"ejecutado de verdad -- L6"| L2

El mapa de este módulo

Leccion   Que resuelve
────────  ──────────────────────────────────────────────────────────────
L1        (esta) Por que este modulo existe, y las cinco formas de
          resolver el mismo cambio, lado a lado
L2        Las tres formas que Kiosko YA resolvio -- DuckDB MERGE, dbt
          snapshot, time travel -- recordadas con codigo citado literal,
          sin re-ejecutar ninguna
L3        Instalando Spark con el runtime de Iceberg -- la unica
          dependencia de JVM de toda esta guia, declarada asi de explicita
L4        MERGE INTO en Spark SQL -- la sintaxis general, verificada
          contra la documentacion oficial
L5        Corriendo el cambio de P002 como un MERGE -- el caso concreto
          de Kiosko, sobre el catalogo local de Spark
L6        table.upsert() de PyIceberg -- la alternativa 100% Python,
          ejecutada de verdad, UpsertResult con rows_updated=1
L7        Eligiendo entre MERGE en SQL y upsert en Python -- el criterio,
          no la moda
L8        Proyecto: el upsert nativo de Kiosko, de punta a punta

Las lecciones 2 a 7 siguen una progresión deliberada: primero el pasado (lección 2, sin ejecutar nada nuevo), después la infraestructura que este módulo necesita y que ningún otro módulo de esta guía requiere (lección 3), después la sintaxis general de MERGE INTO (lección 4) antes de aplicarla al caso concreto de P002 (lección 5), después la alternativa en Python puro (lección 6), y por último el criterio para elegir entre las dos (lección 7). La lección 8 integra el resultado ejecutable de este módulo en un solo proyecto.

Una advertencia técnica, declarada desde ahora

Este es el único módulo de las ocho de esta guía que necesita la JVM (Java Virtual Machine). Los módulos 1, 2, 3, 4, 5, 7 y 8 corren enteramente sobre PyIceberg puro —sin Java, sin Spark, sin ningún proceso externo—. Este módulo reutiliza el PySpark 4.2.0 + Java 17 que spark-and-distributed-processing-guide ya dejó instalado en su módulo 1, lección 4 —el mismo JAVA_HOME apuntando al mismo JDK—, específicamente para correr MERGE INTO con SQL real contra una tabla Iceberg. Si no tienes ese entorno configurado, la lección 3 de este módulo te muestra, paso a paso, cómo verificarlo y qué hacer si falta.

Y una segunda advertencia, igual de explícita: el catálogo que usa Spark en este módulo —llamado local, siguiendo la convención literal de la documentación oficial de Iceberg— es un catálogo físicamente distinto del catálogo kiosko que PyIceberg usó en los módulos 1 a 5. Los dos catálogos viven en directorios distintos, con mecanismos de registro distintos (hadoop uno, sql/SQLite el otro), y esta guía nunca afirma que ambos motores lean o escriban la misma tabla física — eso requeriría verificarlo de verdad, y esta guía no lo hace. Vas a ver el cambio de P002 reproducido dos veces, en dos catálogos independientes, cada uno con su propio kiosko.dim_product — no una tabla compartida entre Spark y PyIceberg.

La frontera: qué NO entra en este módulo

Este módulo no orquesta nada: el MERGE INTO de la lección 5 se corre a mano, desde una terminal o un script — disparar ese mismo comando desde un DAG programado, con reintentos gestionados, es terreno de airflow-and-declarative-orchestration-guide, nombrado aquí sin implementarlo. Tampoco es una lección de cómputo distribuido: Spark aparece en este módulo únicamente como el cliente SQL que sabe hablar MERGE INTO contra Iceberg — el particionado, el shuffle, Catalyst, todo eso vive en spark-and-distributed-processing-guide, y este módulo no lo toca. Y el cambio de P002 sigue llegando, en este módulo, como un valor fijo declarado en Python o en una tabla de staging construida a mano —no como un evento real de un sistema transaccional—: el momento exacto en que eso debería cambiar a CDC real se nombra en el módulo 8, lección 7, y pertenece a streaming-with-kafka-and-flink-guide.

Errores comunes

Pensar que este módulo "deshace" el trabajo del módulo 3, porque ahora hay una forma "mejor" de resolver el cambio de P002. Qué pasa: alguien, al ver que este módulo agrega dos técnicas nuevas para el mismo cambio, asume que table.overwrite() + time travel del módulo 3 quedó obsoleto, o que debería haberse hecho así desde el principio. Por qué pasa: es natural asumir que la técnica más reciente que aprendes es la que reemplaza a las anteriores. Cómo detectarlo: si terminas este módulo pensando que ya no tiene sentido usar table.overwrite() para nada, revisa la lección 7 — el criterio para elegir entre las cuatro técnicas de Iceberg (overwrite completo, time travel, MERGE INTO, upsert) depende de la forma en que el cambio llega, no de cuál "suena" más avanzada. Cómo corregirlo: table.overwrite() sigue siendo la herramienta correcta cuando ya tienes, en memoria, el estado completo de una tabla pequeña —exactamente el caso del módulo 3—; MERGE INTO/upsert son la herramienta correcta cuando lo que te llega es un delta parcial, el caso más común en un pipeline de producción real. Ninguna reemplaza a la otra — cada una resuelve un punto de entrada distinto.

Asumir que, porque MERGE INTO y upsert no necesitan columnas de historia, tampoco necesitan pensar en identidad de fila. Qué pasa: alguien intenta correr un MERGE INTO o un upsert sin definir con claridad qué columna identifica de forma única a cada fila del target —el equivalente de una llave primaria—, y descubre, con un error o con un resultado incorrecto, que la operación necesitaba esa información. Por qué pasa: como Iceberg no obliga a declarar una llave primaria formal en el esquema de la tabla —a diferencia de una base de datos relacional clásica—, es fácil asumir que la identidad de fila es opcional. Cómo detectarlo: la lección 6 de este módulo te muestra, con un error real, exactamente este caso: table.upsert() sin indicar la columna de unión falla con un mensaje explícito. Cómo corregirlo: tanto el ON de MERGE INTO como el join_cols de upsert() necesitan, siempre, una columna (o combinación de columnas) que identifique de forma única cada fila del target — en Kiosko, product_id, la misma llave natural que ya usaste en las tres técnicas anteriores del ecosistema.

Ejercicios

Ejercicio 1 — Nombra, de memoria, las tres formas en que Kiosko ya resolvió el cambio de P002 antes de este módulo. Sin mirar atrás, nombra las tres técnicas, con la guía y el módulo de origen de cada una.

Ver solución
  1. MERGE INTO a mano sobre DuckDB, con columnas valid_from/valid_to/is_currentdata-modeling-for-analytics-guide, módulo 4. 2. dbt snapshot, con columnas dbt_valid_from/dbt_valid_to/dbt_scd_id generadas automáticamente — dbt-analytics-engineering-guide, módulo 5. 3. table.overwrite() + time travel (table.scan(snapshot_id=...)), sin ninguna columna de historia — esta misma guía, módulo 3. Las tres llegan al mismo resultado de negocio (margin=10.8 correcto), con infraestructura y disciplina completamente distintas.

Ejercicio 2 — Traza la analogía tú mismo. Con tus propias palabras, y usando la analogía del mostrador del banco de esta lección, explica qué representan WHEN MATCHED y WHEN NOT MATCHED en MERGE INTO.

Ver solución

WHEN MATCHED es el caso "ya eres cliente de este banco": el product_id que llega en la fuente ya existe en la tabla destino, así que la acción correcta es actualizar el registro existente con los valores nuevos —nunca crear uno duplicado—. WHEN NOT MATCHED es el caso "eres cliente nuevo": el product_id de la fuente no tiene ningún registro previo en el destino, así que la acción correcta es darlo de alta con un INSERT. La clave de la analogía es que las dos ramas conviven en el mismo trámite —el mismo statement de MERGE INTO—, sin que nadie tenga que decidir de antemano cuál de las dos aplica: el propio motor lo resuelve comparando la fuente contra el destino.

Ejercicio 3 — Predicción. Antes de leer la lección 3: ¿por qué crees que esta guía elige reutilizar el PySpark de spark-and-distributed-processing-guide en vez de instalar un motor SQL más liviano —como DuckDB, que ya conoces de data-modeling— para correr MERGE INTO sobre una tabla Iceberg?

Ver solución

No hay una única respuesta correcta —es un ejercicio de predicción—, pero la razón real, que la lección 3 confirma: MERGE INTO nativo de Iceberg —el que entiende snapshots, manifest files, y el resto de la anatomía que el módulo 2 de esta guía enseñó— necesita un motor de cómputo que tenga una integración oficial con el formato de tabla Iceberg. Spark tiene esa integración de forma madura y bien documentada (iceberg-spark-runtime, mantenido por el propio proyecto Apache Iceberg); DuckDB tiene su propio soporte de Iceberg, pero mucho más limitado en 2026, y no es el foco de esta guía. La elección no es "Spark es mejor que DuckDB en general" — es que este módulo específico necesita el motor que la industria usa, hoy, para MERGE INTO transaccional contra tablas Iceberg reales.

Resumen y siguiente paso

En esta lección viste por qué existe este módulo: el módulo 3 resolvió el cambio de P002 con overwrite() completo, pero un pipeline real casi nunca recibe "el estado completo de la dimensión" — recibe un delta. Recorriste el mapa de las cinco formas en que Kiosko ya resolvió, o va a resolver, ese mismo cambio, y las dos advertencias técnicas centrales de este módulo: es el único que necesita la JVM, y usa dos catálogos físicamente distintos, declarados así sin ambigüedad.

Antes de avanzar deberías poder: explicar por qué MERGE INTO/upsert no reemplazan a table.overwrite() del módulo 3, sino que resuelven un punto de entrada distinto; y nombrar las cinco técnicas de este módulo, en orden, con su vehículo correspondiente.

La lección 2 no instala nada todavía — recorre, con código citado literal de data-modeling y dbt, las tres formas que Kiosko ya usó para resolver este mismo cambio, sin volver a correr ninguna.

Recursos

  • Apache Iceberg — documentación oficial, "Spark Writes", sección MERGE INTO — la base formal de la sintaxis que este módulo va a usar en las lecciones 4 y 5. iceberg.apache.org/docs/latest/spark-writes. En inglés.
  • PyIceberg — referencia de API, table.upsert() — la base de la lección 6. py.iceberg.apache.org/api. En inglés.
  • DISEÑO de data-modeling-for-analytics-guide — fuente del MERGE INTO de DuckDB que la lección 2 recuerda. src/guides/data-modeling-for-analytics-guide/DISENO.md. En español.
  • DISEÑO de dbt-analytics-engineering-guide — fuente de dbt snapshot que la lección 2 recuerda. src/guides/dbt-analytics-engineering-guide/DISENO.md. En español.
  • DISEÑO de spark-and-distributed-processing-guide — fuente del entorno PySpark 4.2.0 + Java 17 que la lección 3 reutiliza. src/guides/spark-and-distributed-processing-guide/DISENO.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. src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.