Módulo 6: Merge Into And Native Upserts
Corriendo el cambio de P002 como un MERGE
Descripción
Esta lección aplica la sintaxis de MERGE INTO de la lección 4 al caso concreto que acompaña toda esta guía: P002 Energy Bar cambia de category='snacks', unit_cost=0.60 a category='health-snacks', unit_cost=0.68. Vas a crear local.kiosko.dim_product desde cero, en el catálogo local de Spark —una tabla nueva, independiente de kiosko.dim_product que PyIceberg usó en los módulos 1 a 5—, cargarla con el estado V1, y aplicar el cambio con un solo MERGE INTO, usando una tabla de staging que contiene únicamente la fila que cambió.
Conexión con el módulo. La lección 3 documentó, con evidencia real, que este entorno no puede ejecutar operaciones de catálogo Iceberg contra Spark 4.2.0. Esta lección hereda esa limitación de forma explícita: todo el SQL de acá está marcado "Qué esperar (representativo)" — la sintaxis exacta que ejecutarías, verificada contra la documentación oficial, con una salida cuya estructura y valores son los que ese MERGE produciría de verdad, pero no una corrida capturada en este entorno.
El material: una tabla nueva, un staging con solo el delta
A diferencia de las tres técnicas de la lección 2 —que reciben el estado completo de la tabla nueva—, esta lección usa una tabla de staging que contiene únicamente la fila que cambió. Es, a propósito, el patrón más común en un pipeline de producción real: un feed de cambios, o una exportación diaria del sistema de catálogo de Kiosko, casi nunca trae "los cuatro productos otra vez" — trae, con precisión, lo que cambió desde ayer.
Paso 1 — Crea local.kiosko.dim_product, y cárgala con V1
-- Qué esperar (representativo): sintaxis verificada contra la documentacion
-- oficial de Iceberg ("Getting Started" / "Spark DDL"), no ejecutada en este entorno.
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);
Fíjate en dos cosas. Primero, local.kiosko.dim_product es una tabla nueva, creada desde cero en el catálogo local —nunca la misma tabla física que kiosko.dim_product de PyIceberg, aunque el nombre de la última parte coincida—. Segundo, empieza en V1 —P002 todavía snacks/0.60— exactamente el mismo punto de partida que la lección 2 del módulo 3 usó para kiosko.dim_product, para que el MERGE que sigue tenga algo real que cambiar.
Paso 2 — Crea el staging, con solo el delta
CREATE TABLE local.kiosko.dim_product_staging (
product_id STRING,
product_name STRING,
category STRING,
unit_cost DOUBLE
) USING iceberg;
-- Solo P002: la fila que cambio de verdad. P001, P003, P004 no aparecen aqui
-- -- no cambiaron, y un pipeline real casi nunca los volveria a exportar.
INSERT INTO local.kiosko.dim_product_staging VALUES
('P002', 'Energy Bar', 'health-snacks', 0.68);
Una sola fila. Esto es, a propósito, distinto de la tabla staging_product que viste en el MERGE de DuckDB (lección 2), que sí traía los cuatro productos completos —porque esa lección quería demostrar la idempotencia del MERGE frente a filas sin cambios—. Acá el punto es otro: mostrar el patrón de "delta real", el que vas a encontrar con más frecuencia en un pipeline de producción.
Paso 3 — El MERGE INTO
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). Con el staging conteniendo únicamente P002, y P002 ya existente en el target, la única rama que se activa es WHEN MATCHED — una sola vez. WHEN NOT MATCHED nunca se ejercita, porque no hay ninguna fila del source sin coincidencia en el target. Verificado contra la documentación oficial de Iceberg (Spark 4.1 en adelante), el resumen del snapshot que produce este MERGE incluiría, con estos valores lógicos para este caso concreto —no una corrida capturada—:
spark.merge-into.num-target-rows-copied = 3 -- P001, P003, P004: sin cambios
spark.merge-into.num-target-rows-updated = 1 -- P002: la unica fila que cambio
spark.merge-into.num-target-rows-inserted = 0 -- ninguna fila nueva
spark.merge-into.num-target-rows-deleted = 0 -- ningun DELETE en este MERGE
spark.merge-into.num-target-rows-matched-updated = 1 -- confirma: el 1 de arriba vino de WHEN MATCHED
Y el resultado de negocio, consultando la tabla después del MERGE:
SELECT * FROM local.kiosko.dim_product ORDER BY product_id;
Qué esperar (representativo):
product_id | product_name | category | unit_cost
-----------+------------------------+---------------+----------
P001 | Bottled Water 600ml | beverages | 0.40
P002 | Energy Bar | health-snacks | 0.68
P003 | Instant Coffee Sachet | beverages | 0.35
P004 | Phone Charger Cable | electronics | 2.10
Cuatro filas, P002 con el valor nuevo, P001/P003/P004 intactos — exactamente el mismo resultado de negocio que las tres técnicas de la lección 2 ya produjeron, con una diferencia mecánica central: nadie tuvo que reconstruir las tres filas sin cambios en Python o en un CSV — el staging solo trajo el delta, y MERGE INTO se encargó de identificar que las otras tres filas del target no tenían ninguna coincidencia que actualizar, dejándolas exactamente como estaban.
Diagrama: una fila entra, una fila cambia
flowchart LR
subgraph staging["local.kiosko.dim_product_staging"]
S1["P002 health-snacks 0.68\n(la unica fila)"]
end
subgraph target["local.kiosko.dim_product ANTES"]
T1["P001 beverages 0.40"]
T2["P002 snacks 0.60"]
T3["P003 beverages 0.35"]
T4["P004 electronics 2.10"]
end
S1 -->|"ON product_id coincide"| T2
T2 -->|"WHEN MATCHED\nUPDATE"| R2["P002 health-snacks 0.68"]
T1 -.->|"sin coincidencia en source\nqueda igual"| R1["P001 beverages 0.40"]
T3 -.->|"sin coincidencia en source\nqueda igual"| R3["P003 beverages 0.35"]
T4 -.->|"sin coincidencia en source\nqueda igual"| R4["P004 electronics 2.10"]
Profundización: por qué el staging de un solo producto es más realista que el estado completo
Vale la pena detenerse en algo que separa esta lección de las tres técnicas de la lección 2: data-modeling y dbt recibían, cada una, una fuente que traía los cuatro productos, con tres de ellos idénticos a la versión anterior — ese patrón funciona bien cuando la fuente completa es pequeña y barata de regenerar (un archivo CSV, un dict de Python), pero no escala cuando la dimensión tiene miles o millones de filas. En ese caso, exportar "todo el catálogo, otra vez" en cada corrida sería enormemente más costoso que exportar solo lo que cambió. El staging de una sola fila que usa esta lección es el patrón que sí escala: cualquier sistema que produzca un feed de cambios —un CDC real, una cola de eventos, un export incremental— naturalmente entrega solo los deltas, y MERGE INTO, con su WHEN MATCHED/WHEN NOT MATCHED, está diseñado exactamente para consumir ese tipo de fuente sin que nadie tenga que reconstruir el resto de la tabla.
Errores comunes
Incluir las tres filas sin cambios en el staging, "por las dudas". Qué pasa: alguien, acostumbrado al patrón de data-modeling/dbt de la lección 2, incluye las cuatro filas completas en dim_product_staging, en vez de solo P002. Por qué pasa: parece más seguro "traer todo" que confiar en que el MERGE va a manejar correctamente una fuente parcial. Cómo detectarlo: no es un error que rompa nada —MERGE INTO maneja perfectamente un staging con las cuatro filas, simplemente encontraría tres coincidencias sin cambios reales de valor y las "actualizaría" con los mismos valores que ya tenían—, pero es una oportunidad perdida de mostrar el patrón que sí escala. Cómo corregirlo: si tu fuente de cambios real solo trae lo que cambió —el caso más común en producción—, deja que el staging refleje eso fielmente; no reconstruyas artificialmente filas sin cambios solo por costumbre.
Olvidar que WHEN NOT MATCHED nunca se activa en este caso concreto, y esperar ver num-target-rows-inserted > 0. Qué pasa: alguien, familiarizado con la estructura general de MERGE INTO de la lección 4 —que incluye WHEN NOT MATCHED THEN INSERT—, espera ver alguna fila insertada en este ejemplo específico, aunque P002 ya existiera en el target. Por qué pasa: ver una cláusula WHEN NOT MATCHED en el SQL hace sentir que, en algún punto, debería ejercitarse. Cómo detectarlo: si esperas num-target-rows-inserted > 0 para este caso concreto, repasa el staging del Paso 2 — contiene únicamente P002, un producto que ya existe en target, así que cada fila del source encuentra coincidencia, y WHEN NOT MATCHED nunca se activa. Cómo corregirlo: la cláusula WHEN NOT MATCHED está ahí para el caso general —un producto nuevo, no visto antes—, no porque este caso específico la necesite; para ver esa rama en acción, tendrías que agregar al staging un product_id que no exista todavía en target.
Ejercicios
Ejercicio 1 — Reescribe el staging para incluir un producto hipotético nuevo, y predice el resultado. Agrega una fila ('P005', 'Reusable Tote Bag', 'accessories', 1.20) al staging de esta lección, junto con la de P002. Sin correrlo, predice: ¿qué valores tendrían num-target-rows-updated y num-target-rows-inserted ahora?
Ver solución
num-target-rows-updated = 1 (sigue siendo solo P002, la única fila que coincide y cambia de valor) y num-target-rows-inserted = 1 (la fila de P005, que no tiene ninguna coincidencia en target, activa WHEN NOT MATCHED THEN INSERT). El resultado final tendría cinco filas en local.kiosko.dim_product, no cuatro. Nota: este P005 es puramente hipotético, para este ejercicio — el catálogo real de Kiosko, en toda esta guía, sigue teniendo exactamente cuatro productos.
Ejercicio 2 — Explica por qué esta lección usa ON target.product_id = source.product_id, sin ninguna condición adicional, a diferencia del ON de DuckDB en la lección 2. En 2-3 frases, explica la diferencia.
Ver solución
El ON de DuckDB en data-modeling necesitaba AND target.is_current = true, porque dim_product_scd puede tener varias filas para el mismo product_id —una por cada versión histórica—, y sin ese filtro adicional, el MERGE encontraría más de una coincidencia para la misma fila de origen (un error, según la regla dura de la lección 4). local.kiosko.dim_product, en esta lección, nunca tiene más de una fila por product_id —no tiene columnas de historia—, así que target.product_id = source.product_id ya es, por sí sola, una condición que garantiza como máximo una coincidencia. No hace falta ningún filtro adicional porque no hay ninguna ambigüedad que resolver.
Ejercicio 3 — Predicción: ¿qué pasaría si corrieras este mismo MERGE INTO una segunda vez, sin cambiar el staging? Sin correrlo, predice: si ejecutas el MERGE del Paso 3 dos veces seguidas, con el mismo staging (solo P002, health-snacks/0.68), ¿qué esperas ver en num-target-rows-updated la segunda vez?
Ver solución
Con la sintaxis exacta de esta lección —WHEN MATCHED THEN UPDATE SET sin ninguna condición adicional que compare valores—, la segunda corrida seguiría reportando num-target-rows-updated = 1, aunque los valores ya sean idénticos: la cláusula WHEN MATCHED, tal como está escrita, no verifica si algo realmente cambió, solo verifica que hay una coincidencia. Esto es distinto del MERGE de DuckDB en la lección 2, que sí incluía una condición explícita (AND (target.unit_cost <> source.unit_cost OR ...)) para evitar "actualizar" filas sin cambios reales — y también distinto de table.upsert() de PyIceberg, que vas a ver en la lección 6, que sí hace esa comparación internamente. Si quisieras que este MERGE INTO fuera igual de selectivo, tendrías que agregar la misma condición explícita al WHEN MATCHED que ya viste en el MERGE de DuckDB.
Resumen y siguiente paso
En esta lección aplicaste la sintaxis de MERGE INTO al caso real de P002: una tabla local.kiosko.dim_product nueva, cargada con V1, y un staging que trae únicamente el delta —una sola fila, P002 en su estado nuevo—. El MERGE actualiza esa única fila coincidente y deja las otras tres intactas, sin ningún INSERT de seguimiento, sin ninguna columna de historia. Toda la salida de esta lección está marcada como representativa, verificada contra la documentación oficial de Iceberg, no ejecutada en este entorno, por la incompatibilidad real que la lección 3 documentó.
Antes de avanzar deberías poder: escribir de memoria el MERGE INTO completo para el caso de P002; explicar por qué un staging de una sola fila es más realista que el estado completo de la tabla; y explicar la diferencia entre este MERGE (que actualiza sin verificar si el valor cambió) y el de DuckDB (que sí lo verifica).
La lección 6 muestra la alternativa 100% Python: table.upsert() de PyIceberg, sin ningún SQL, corrida de verdad en este entorno —sin la incompatibilidad de Spark que limitó estas dos lecciones—.
Recursos
- Apache Iceberg — documentación oficial, "Spark Writes", sección
MERGE INTO, fuente de la sintaxis exacta que esta lección aplica al caso de Kiosko. iceberg.apache.org/docs/latest/spark-writes/#merge-into. En inglés. - Apache Iceberg — documentación oficial, "Spark DDL", sintaxis de
CREATE TABLE ... USING iceberg. iceberg.apache.org/docs/latest/spark-ddl. En inglés. - Esta misma guía, módulo 6, lección 3 — fuente de la incompatibilidad real que explica por qué esta lección es representativa.
03-setting-up-spark-with-the-iceberg-runtime.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.