Módulo 8: Project Kioskos Lakehouse

Fusionando actualizaciones a la manera nativa

Descripción

Esta lección no cambia kiosko.dim_product otra vez — ya quedó, desde la lección 4, con P002 en health-snacks/0.68, y ese estado no se toca. Lo que hace esta lección es una pregunta distinta, con evidencia nueva: si Kiosko tuviera que aplicar este mismo tipo de cambio de forma rutinaria —no como un overwrite() que exige la tabla completa, sino como una actualización parcial, "solo lo que cambió"—, ¿el resultado de negocio sería el mismo? La respuesta, verificada con table.upsert() sobre una tabla de comprobación aparte, es sí — byte a byte, fila por fila, idéntico al resultado que dejó table.overwrite() en la lección 4.

Conexión con el módulo. Esta lección retoma, aplicado dentro del lakehouse ensamblado, el mecanismo que el módulo 6 de esta guía ya enseñó a fondo: table.upsert() como la vía nativa de Iceberg para aplicar cambios parciales, frente al overwrite() que exige recalcular la tabla completa, y frente al MERGE INTO de Spark SQL, que ese mismo módulo documentó como representativo por la incompatibilidad real entre iceberg-spark-runtime-4.0 y Spark 4.1+ (apache/iceberg#15238).

Una analogía: la misma reforma, con dos contratistas distintos, el mismo resultado final

Imagina que dos contratistas distintos reciben el encargo de cambiar el letrero de una sola tienda dentro de un centro comercial de cuatro locales. El primer contratista —el que ya trabajó en la lección 4— resuelve el encargo reconstruyendo la fachada completa del centro comercial: toma los cuatro letreros, cambia el que hace falta, y vuelve a instalar los cuatro juntos. El segundo contratista —el de esta lección— resuelve el mismo encargo de otra forma: sube una escalera, retira solo el letrero que cambió, y lo reemplaza, sin tocar los otros tres. Si ambos hicieron bien su trabajo, un visitante que llega después no puede notar cuál contratista pasó por ahí — el centro comercial se ve idéntico. La diferencia entre los dos no está en el resultado final, está en cuánto trabajo y cuánta información necesitó cada uno: el primero necesitó conocer, y volver a entregar, los cuatro letreros completos; el segundo solo necesitó el letrero que cambió.

Ejemplo trabajado: el mismo cambio de P002, aplicado con table.upsert()

Paso 1 — Una tabla de comprobación, separada de la real

Esta lección no vuelve a tocar kiosko.dim_product — crea una tabla nueva, kiosko.dim_product_merge_check, cargada con la misma V1, para demostrar el mecanismo sin arriesgar el estado que la lección 4 ya dejó verificado:

# kiosko_native_merge_check.py -- modulo 8, leccion 6
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},
]

Paso 2 — La diferencia clave: el insumo que cada mecanismo necesita

def main() -> None:
    print("=== Kiosko: el mismo cambio de P002, aplicado con table.upsert() -- la via nativa ===\n")

    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}",
    )
    real_dim_product = catalog.load_table("kiosko.dim_product")
    real_rows = sorted(real_dim_product.scan().to_arrow().to_pylist(), key=lambda r: r["product_id"])
    print(f"Paso 1/5 -- kiosko.dim_product (leccion 4, cambiada con overwrite()) sigue como quedo: "
          f"P002={next(r for r in real_rows if r['product_id'] == 'P002')['category']}")

    check_table = catalog.create_table("kiosko.dim_product_merge_check", schema=DIM_PRODUCT_SCHEMA)
    check_table.append(pa.Table.from_pylist(DIM_PRODUCT_V1, schema=PA_SCHEMA))
    print("Paso 2/5 -- kiosko.dim_product_merge_check creada aparte, cargada con la misma V1")

    only_p002_change = [r for r in DIM_PRODUCT_V2 if r["product_id"] == "P002"]
    result = check_table.upsert(pa.Table.from_pylist(only_p002_change, schema=PA_SCHEMA), join_cols=["product_id"])
    print(f"Paso 3/5 -- table.upsert() con SOLO la fila que cambio: rows_updated={result.rows_updated}, "
          f"rows_inserted={result.rows_inserted} -- nadie calculo el delta a mano")

    check_rows = sorted(check_table.scan().to_arrow().to_pylist(), key=lambda r: r["product_id"])
    ops = [row["operation"] for row in check_table.inspect.snapshots().select(["operation"]).to_pylist()]
    print(f"Paso 4/5 -- table.history() de la tabla de verificacion: {ops} -- el upsert se resuelve, "
          f"internamente, como overwrite (parcial) + append (0 filas, ninguna fila nueva que insertar), "
          f"nunca como el overwrite() del TOTAL de filas que pediria hacerlo a mano")

    same_result = real_rows == check_rows
    print(f"Paso 5/5 -- kiosko.dim_product (overwrite, leccion 4) == kiosko.dim_product_merge_check "
          f"(upsert, esta leccion), fila por fila: {same_result}\n")

    print("=== Verificacion final ===\n")
    assert result.rows_updated == 1 and result.rows_inserted == 0
    assert ops == ["append", "overwrite", "append"], "V1 (append) + upsert (overwrite parcial + append vacio)"
    assert same_result, "el resultado de negocio debe ser identico sin importar el mecanismo de escritura"
    p002_check = next(r for r in check_rows if r["product_id"] == "P002")
    assert p002_check["category"] == "health-snacks" and p002_check["unit_cost"] == 0.68

    print("Todas las verificaciones pasaron:")
    print("  - table.upsert() detecto y aplico solo el cambio real de P002, sin tocar P001/P003/P004")
    print("  - el resultado de negocio es identico, fila por fila, al que dejo table.overwrite() en la leccion 4")
    print("  - la diferencia es el INSUMO que cada mecanismo necesita: overwrite() exige la tabla COMPLETA, "
          "upsert() solo exige el DELTA -- un solo producto, una sola fila")


if __name__ == "__main__":
    main()

Qué esperar (verificado corriendo python3 kiosko_native_merge_check.py real, en el mismo directorio que las lecciones 3 a 5, sin borrar kiosko_warehouse/):

=== Kiosko: el mismo cambio de P002, aplicado con table.upsert() -- la via nativa ===

Paso 1/5 -- kiosko.dim_product (leccion 4, cambiada con overwrite()) sigue como quedo: P002=health-snacks
Paso 2/5 -- kiosko.dim_product_merge_check creada aparte, cargada con la misma V1
Paso 3/5 -- table.upsert() con SOLO la fila que cambio: rows_updated=1, rows_inserted=0 -- nadie calculo el delta a mano
Paso 4/5 -- table.history() de la tabla de verificacion: ['append', 'overwrite', 'append'] -- el upsert se resuelve, internamente, como overwrite (parcial) + append (0 filas, ninguna fila nueva que insertar), nunca como el overwrite() del TOTAL de filas que pediria hacerlo a mano
Paso 5/5 -- kiosko.dim_product (overwrite, leccion 4) == kiosko.dim_product_merge_check (upsert, esta leccion), fila por fila: True

=== Verificacion final ===

Todas las verificaciones pasaron:
  - table.upsert() detecto y aplico solo el cambio real de P002, sin tocar P001/P003/P004
  - el resultado de negocio es identico, fila por fila, al que dejo table.overwrite() en la leccion 4
  - la diferencia es el INSUMO que cada mecanismo necesita: overwrite() exige la tabla COMPLETA, upsert() solo exige el DELTA -- un solo producto, una sola fila

Fíjate en el Paso 3: only_p002_change contiene una sola fila, no las cuatro. table.upsert() no necesita que le entregues P001, P003 ni P004 para dejarlos intactos — los detecta como "no coincide con ningún cambio real" y los deja exactamente como estaban. table.overwrite(), en cambio —el mecanismo que usó la lección 4—, exige recibir siempre la tabla completa: si le hubieras pasado solo la fila de P002, las otras tres habrían desaparecido de la tabla.

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

Para que esta lección documente el mecanismo completo de "la vía nativa" que enseñó el módulo 6, y no solo la mitad que corre en este entorno, esta es la sintaxis exacta del MERGE INTO, aplicada al mismo cambio de P002, marcada explícitamente como lo que es — la misma incompatibilidad real (apache/iceberg#15238) que documentó el módulo 6:

-- Qué esperar (representativo) -- sintaxis identica a la del modulo 6, no ejecutada
-- en este entorno por la incompatibilidad real entre iceberg-spark-runtime-4.0 y Spark 4.1+.

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 ya confirmaron table.overwrite() (lección 4) y table.upsert() (esta lección) — P002 con category='health-snacks', unit_cost=0.68, P001/P003/P004 intactos. La diferencia, documentada a fondo en el módulo 6: este MERGE INTO corre distribuido, sobre Spark, contra el catálogo local (Hadoop) — un catálogo físicamente distinto del kiosko (SQL/SQLite) que usa el resto de este capstone —, y necesita la JVM; table.upsert() corre en un solo proceso Python, sin ninguna de las dos.

Diagrama: tres caminos, un mismo destino

flowchart TB
    START["Cambio de P002:\nsnacks/0.60 -> health-snacks/0.68"]
    START --> OW["table.overwrite()\n(leccion 4)\nexige la tabla COMPLETA"]
    START --> UP["table.upsert()\n(esta leccion)\nsolo exige el DELTA"]
    START --> MI["MERGE INTO Spark SQL\n(modulo 6, representativo)\ncatalogo distinto, necesita JVM"]
    OW --> RESULT["kiosko.dim_product:\nP002 = health-snacks/0.68\nP001/P003/P004 intactos"]
    UP --> RESULT
    MI -.->|"mismo resultado esperado,\nno ejecutado en este entorno"| RESULT

Profundización: por qué "el mismo resultado" no significa "el mismo costo"

El Paso 5 de esta lección confirma que overwrite() y upsert() producen exactamente el mismo estado final — pero esa igualdad esconde una diferencia de costo real que el módulo 6 ya cuantificó en su lección 7 ("Choosing between SQL MERGE and Python upsert"): sobre una tabla de cuatro filas, como dim_product, la diferencia entre "reescribir todo" y "reescribir solo el delta" es invisible — ambos mecanismos terminan en milisegundos. Pero la misma decisión, aplicada a una tabla con millones de filas y solo un puñado de cambios reales por día —el tipo de tabla que un lakehouse de producción de verdad maneja—, deja de ser una elección estética: overwrite() obligaría a leer, procesar y volver a escribir todas las filas para cambiar un puñado; upsert() procesa únicamente el delta. Esta lección demuestra la equivalencia de resultado a pequeña escala, precisamente porque a esta escala es fácil verificar, fila por fila, que ningún mecanismo introdujo un error — la decisión de cuál usar en producción depende del volumen, no de cuál "se ve más simple" en un ejemplo de cuatro filas.

Errores comunes

Modificar kiosko.dim_product (la tabla real) en vez de kiosko.dim_product_merge_check. Qué pasa: alguien, al adaptar este código, llama catalog.load_table("kiosko.dim_product") en vez de crear la tabla de comprobación aparte, y corre upsert() sobre la tabla real que la lección 4 ya dejó verificada. Por qué pasa: es más corto escribir "la tabla real" que crear una tabla nueva solo para una comprobación. Cómo detectarlo: si el assert p002_check["category"] == "health-snacks" de esta lección falla con un error de tabla no encontrada, o si el kiosko.dim_product de una lección posterior muestra un historial de snapshots distinto al que dejó la lección 4, mezclaste las dos tablas. Cómo corregirlo: mantén kiosko.dim_product_merge_check como una tabla completamente separada — su único propósito es demostrar la equivalencia de resultado, nunca reemplazar el estado que la lección 4 ya verificó y que la lección 8 va a dar por sentado.

Asumir que ops == ["append", "overwrite", "append"] significa que upsert() "en realidad hace un overwrite completo por dentro". Qué pasa: alguien, al ver la palabra "overwrite" en la lista de operaciones, concluye que table.upsert() no es distinto de table.overwrite() a nivel interno, y que toda la comparación de esta lección es una diferencia solo de nombre. Por qué pasa: PyIceberg reutiliza la misma palabra "overwrite" para describir el tipo de operación de snapshot, tanto si se llama explícitamente table.overwrite() como si es el resultado interno de un upsert() parcial. Cómo detectarlo: si tu conclusión es que ambos mecanismos leen y reescriben el mismo volumen de datos, revisa qué pyarrow.Table recibió cada uno como argumento — table.overwrite() en la lección 4 recibió las cuatro filas de V2; table.upsert() en esta lección recibió una sola fila. Cómo corregirlo: la palabra "overwrite" en table.history() describe el tipo de snapshot (reemplaza archivos existentes en vez de solo agregar), no el volumen de datos involucrado — la diferencia real entre los dos mecanismos está en cuántas filas tuvo que recibir cada uno como insumo, no en cómo se etiqueta la operación resultante.

Ejercicios

Ejercicio 1 — Corre el script tú mismo, en el mismo directorio que las lecciones 3 a 5. Confirma que ves los cinco pasos completarse y el mensaje final con las cuatro verificaciones.

Ver solución

Si corriste las lecciones 3, 4 y 5 primero, en el mismo directorio, la salida debería reproducir exactamente la estructura de esta lección: cinco pasos numerados, seguidos de la verificación final confirmando rows_updated=1, rows_inserted=0, y same_result=True. Si kiosko.dim_product no existe todavía, corriste esta lección antes que la 4.

Ejercicio 2 — Modifica el script para pasarle a table.upsert() las cuatro filas de DIM_PRODUCT_V2, en vez de solo la de P002, y observa si el resultado cambia. Corre el script modificado y compara rows_updated/rows_inserted contra el original.

Ver solución

Con las cuatro filas de V2 como argumento, result.rows_updated sigue siendo 1 (solo P002 cambió de verdad) y result.rows_inserted sigue siendo 0 — el resultado de negocio es idéntico. Este ejercicio confirma algo importante sobre table.upsert(): no importa cuántas filas le entregues, siempre detecta cuáles cambiaron de verdad comparando contra el estado actual, fila por fila. La ventaja práctica de pasarle solo el delta (como hace el script original) no es de corrección —ambas formas llegan al mismo resultado—, es de eficiencia: no hace falta que quien escribe el pipeline sepa de antemano cuál fila cambió, pero tampoco hace falta desperdiciar trabajo reenviando las que no cambiaron si ya se conocen.

Ejercicio 3 — Explica, en tus propias palabras, por qué esta lección demuestra la equivalencia de resultado con una tabla de comprobación aparte, en vez de simplemente re-verificar kiosko.dim_product con un assert adicional. En 2-3 frases, justifica esta decisión de diseño.

Ver solución

Si esta lección corriera upsert() directamente sobre kiosko.dim_product, el resultado final seguiría siendo health-snacks/0.68 — pero ya no habría forma de distinguir "esto quedó así por el overwrite() de la lección 4" de "esto quedó así por el upsert() de esta lección", porque ambos producirían el mismo estado sobre la misma tabla. Crear kiosko.dim_product_merge_check como una tabla aparte, cargada de forma independiente con la misma V1, permite comparar los dos mecanismos de forma aislada, fila por fila, con evidencia de que llegan al mismo resultado por caminos distintos — la misma disciplina de "verificar el camino, no solo el destino" que ya aplicaron los proyectos de cierre de los módulos 4 y 5 de esta guía.

Resumen y siguiente paso

En esta lección aplicaste el mismo cambio de P002 con table.upsert(), sobre una tabla de comprobación separada, y confirmaste —con assert, fila por fila— que el resultado es idéntico al que dejó table.overwrite() en la lección 4. Documentaste, junto a esa evidencia, la sintaxis completa del MERGE INTO vía Spark que el módulo 6 ya marcó como representativo por la incompatibilidad real entre iceberg-spark-runtime-4.0 y Spark 4.1+. El lakehouse de Kiosko tiene, ahora, sus cinco tablas completas, y la confirmación de que dos vías de escritura distintas —overwrite() completo, upsert() parcial— convergen en el mismo resultado de negocio.

Antes de avanzar deberías poder: explicar la diferencia de insumo entre table.overwrite() y table.upsert(), y por qué esa diferencia importa más a escala que sobre una tabla de cuatro filas; y nombrar el catálogo físicamente distinto que necesitaría el MERGE INTO de Spark para correr.

La lección 7 cierra el ecosistema completo: nombra, una por una, las siete guías hermanas de data-engineering-ecosystem y qué resuelve cada una sobre lo que este lakehouse deja pendiente.

Recursos

  • PyIceberg — documentación oficial (quickstart), el flujo de append() y upsert() que integra esta lección. 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 esta lección. 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 esta lección es representativa. github.com/apache/iceberg/issues/15238. En inglés.
  • Esta misma guía, módulo 6, lección 7 — fuente del criterio completo de cinco factores para elegir entre MERGE INTO SQL y table.upsert() Python. ../module-06-merge-into-and-native-upserts/es/07-choosing-between-sql-merge-and-python-upsert.md. En español.
  • DISEÑO de esta guía — el mapa completo de los ocho módulos, incluida la lección 7 que sigue. src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.