Sesión 7 — Iceberg: del data lake al lakehouse, evolución y optimización¶

Esta sesión reúne, de forma deliberada, dos bloques: la creación y evolución de una tabla Iceberg, y la medición y optimización de esa misma tabla. Van juntos porque comparten la tabla de trabajo, el dataset TPC-DS SF1 que ya conoces de S2-S6 y buena parte del vocabulario: no tiene sentido medir manifests y snapshots sobre una tabla que todavía no has visto crear ni evolucionar.

Llegas a este punto habiendo trabajado con Parquet, con Spark como motor de transformación, con particiones físicas y con tablas Hive registradas en un Metastore compartido por Trino y Spark (S6). Ya sabes que un catálogo facilita encontrar los datos y compartir su definición entre motores, pero también has visto sus límites: una colección de directorios y entradas de metastore no resuelve por sí sola la evolución de esquema, la consistencia entre escrituras concurrentes ni la limpieza de versiones antiguas.

La pregunta central de esta sesión es doble, y las dos mitades se responden con la misma tabla:

¿Qué añade un formato de tabla como Iceberg para que un data lake pueda comportarse como un lakehouse? Y una vez que existe esa tabla, ¿cómo medimos y mantenemos sus optimizaciones para que sean útiles y no sólo una colección de propiedades activadas a ciegas?

El salto de Hive a Iceberg no consiste en poner otra etiqueta al mismo directorio. Iceberg introduce un modelo de tabla con metadatos versionados, manifests, snapshots y reglas explícitas de evolución. El catálogo —aquí, otra vez el Hive Metastore que ya conoces de S6— sigue siendo necesario para localizar la tabla, pero ya no es él quien define qué es la tabla: eso lo hace el formato Iceberg. Y una vez creada esa tabla, es donde mediremos qué compensa de verdad: particionado, estadísticas, compactación, ordenación física y, sólo como ampliación si resulta reproducible, Z-order.

Paquetes de Python de esta sesión. Esta sesión necesita pyspark, trino, pandas. No añade ningún paquete respecto a S6. Se instalan en el kernel de namenode en la sección «Instalar las dependencias de este notebook»: una celda %%writefile requirements.txt crea el fichero dentro del contenedor —el requirements.txt de tu copia de s7/ está en tu equipo y namenode no lo ve— y otra ejecuta %pip install -r requirements.txt. La imagen de namenode ya trae estos paquetes en su entorno Python del curso, /opt/tcdm/venv, que es el intérprete del kernel y donde instala %pip; lo normal es que %pip sólo confirme que están; ejecuta esas celdas de todos modos: dejan explícito qué necesita la sesión y reparan el entorno si falta algo.

Diapositivas de la sesión¶

Las diapositivas se generan con el paquete jupyter-notebook-slide. El alias %%diapositiva permite mantener en español los tipos usados en la sesión y produce la misma salida HTML en Jupyter, Colab y la referencia HTML publicada.

En cada sesión encontrarás varios tipos de diapositivas que te ayudarán a seguir el hilo:

  • Sección (titulo): abre la sesión y resume qué vamos a ver y qué haremos.
  • A continuación (avance): anuncia el contenido de la parte siguiente, para saber en qué punto del programa estamos.
  • Recapitulación (resumen): condensa la parte anterior en puntos clave para recordarlos y revisarlos después.
  • Pregunta guía (pregunta): plantea cuestiones que orientan lo que viene a continuación; conviene intentar responderlas antes de seguir.
  • Evaluación de la sesión (evaluacion): recuerda cómo se evalúa el trabajo de esta sesión.
In [ ]:
%pip install -q "jupyter-notebook-slide @ git+https://github.com/dsevilla/jupyter-notebook-slide.git"
%load_ext notebook_slide

import notebook_slide as jnbs

# Colores de 26-27/teoria/tcdm.css, adaptados al tema de las sesiones.
jnbs.configure(
    background="#eaf2f8",
    foreground="#1f2933",
    border="#c9d6e1",
    heading="#0c304d",
    subheading="#0c304d",
    link="#174f7a",
    code_background="#eaf2f8",
    code_foreground="#0c304d",
    quote_background="#ffffff",
    font_family="Atkinson Hyperlegible, Inter, Aptos, Segoe UI, Helvetica, Arial, sans-serif",
    font_url="https://fonts.googleapis.com/css2?family=Atkinson+Hyperlegible:ital,wght@0,400;0,700;1,400;1,700&display=swap",
)
jnbs.register_slide_type("avance", "A continuación", "#174f7a")
jnbs.register_slide_type("resumen", "Recapitulación", "#a54467")
jnbs.register_slide_type("pregunta", "Pregunta guía", "#0c304d")
jnbs.register_slide_type("evaluacion", "Evaluación de la sesión", "#b3701a")
jnbs.register_slide_type(
    "titulo",
    "Sección",
    "#174f7a",
    layout="title",
    background="linear-gradient(135deg, #0c304d, #174f7a 62%, #a54467)",
    foreground="#ffffff",
    border="transparent",
    heading="#ffffff",
    subheading="#dceaf4",
)
jnbs.register_alias("diapositiva")

Sesión 7: Iceberg, del data lake al lakehouse

Tres situaciones, un mismo hilo conductor

  • Lake de ficheros, lake catalogado y lakehouse
  • El catálogo localiza, el formato de tabla define
  • La misma tabla sirve para crearla, evolucionarla y medirla

Recorrido del curso¶

S2: Parquet como ficheros
  → S3: Parquet en HDFS y S3-compatible
  → S4: Spark transforma los datos
  → S5: data lake organizado y particionado (medallion raw/silver/gold)
  → S6: tablas Hive con Metastore y Trino
  → S7: tabla Iceberg catalogada, evolución, snapshots y optimización medida
  → S8: ingesta continua con Spark Structured Streaming

La ruta principal de esta sesión es la de HDFS con Hive Metastore e Iceberg, la misma infraestructura de entorno/compose-warehouse-hdfs.yml que ya arrancaste en S6. La configuración S3 con RustFS, Polaris y el contenedor spark-iceberg que también existe en este entorno es una demostración de portabilidad para otro momento del curso, no el primer contacto con Iceberg.

Data lake, data lake catalogado y lakehouse¶

Antes de tocar una sola celda conviene fijar tres situaciones que a menudo se confunden porque comparten Parquet como formato de fichero:

Situación Qué existe Qué problemas quedan
Data lake de ficheros Parquet en HDFS o S3 Rutas, esquemas y cambios se gestionan manualmente
Data lake catalogado Parquet + Hive Metastore Hay nombres y particiones, pero las operaciones tienen límites
Lakehouse Formato de tabla + catálogo + almacenamiento Metadatos, snapshots y cambios se coordinan como una tabla

S2 te dejó en la primera fila: ficheros Parquet localizables por ruta. S6 te llevó a la segunda: los mismos ficheros, con nombre y particiones conocidas por Trino y Spark a través del metastore. Esta sesión te lleva a la tercera fila, y es importante no simplificar el paso: un data lakehouse no exige que desaparezcan Parquet o el propio Hive Metastore. Iceberg sigue escribiendo ficheros Parquet como datos, y sigue apoyándose en el mismo metastore para que Trino localice la tabla por su nombre lógico. Lo que añade es una capa de metadatos de tabla —manifests y snapshots— que decide qué ficheros pertenecen a cada versión de la tabla y cómo se interpretan, algo que ni el sistema de ficheros ni el metastore por sí solos pueden ofrecer.

Objetivos de la sesión¶

Al terminar esta sesión deberías poder:

  • explicar la diferencia entre un fichero Parquet, una tabla Hive externa y una tabla Iceberg;
  • explicar el papel del catálogo y distinguir el Hive Metastore (quién localiza la tabla) del formato de tabla Iceberg (qué es la tabla);
  • crear una tabla Iceberg desde Trino, con una transformación de partición oculta sobre una columna de fecha;
  • leer la misma tabla desde Trino y desde Spark y obtener el mismo resultado de negocio;
  • observar la ubicación física de los datos y los metadatos bajo /warehouse, sin deducirla del nombre lógico;
  • identificar snapshots y relacionarlos con las operaciones de escritura que los producen;
  • realizar una evolución de esquema segura y observable, y comprobar que las filas y snapshots anteriores siguen siendo legibles;
  • comparar el particionado visible de Hive con las transformaciones de partición ocultas de Iceberg sobre la misma columna de negocio;
  • explicar por qué la selección de columnas y el predicate pushdown siguen siendo la primera optimización a comprobar, antes de técnicas más complejas;
  • detectar el problema de los ficheros pequeños y aplicar compactación con EXECUTE optimize;
  • utilizar estadísticas (ANALYZE) para entender la planificación, sin asumir que siempre cambian la decisión del motor;
  • distinguir la ordenación de un resultado SQL (ORDER BY) de un sort order de tabla que describe cómo se organizan físicamente los datos escritos;
  • valorar Z-order como técnica avanzada dependiente del motor y la versión, y saber cuándo documentarla en lugar de fingir una demostración;
  • mantener snapshots con una política explícita de retención y practicar su expiración únicamente sobre tablas de laboratorio;
  • describir qué características convierten al conjunto en un lakehouse y reconocer cuándo una optimización concreta no aporta beneficio real.
Evaluación de la sesión

Cómo se evalúa esta sesión

  • Reunión individual breve con el profesor
  • Se muestra en vivo: SHOW CREATE TABLE con la ubicación bajo /warehouse, el listado data/metadata, dos snapshots y el viaje en el tiempo, y el inventario de ficheros antes y después de compactar y expirar
  • Se entrega también una memoria breve (una o dos páginas)

Antes de empezar¶

Dónde se ejecuta este notebook. Igual que en S2-S6, Jupyter se ejecuta directamente dentro del contenedor namenode, como luser, con acceso de red directo a HDFS, YARN, Trino, PostgreSQL y el Hive Metastore: todos comparten la red Docker hadoop-cluster. Las celdas ejecutables hablan directamente con esos servicios. Las órdenes que necesitan el Docker del host —arrancar o detener contenedores, make— se muestran como texto para ejecutarlas en una terminal de tu equipo, nunca como celdas de este notebook.

Esta sesión asume que:

  • el clúster de S1 y el warehouse de S6 (PostgreSQL, Hive Metastore, Trino) ya están en marcha;
  • S2 generó los Parquet de TPC-DS SF1 en /datalake/raw/tpcds;
  • S5 dejó /datalake/silver/tpcds/web_sales (particionado por sold_year y sold_month) y /datalake/gold/tpcds/web_sales_by_year escritos en HDFS;
  • S6 registró esos mismos datos como tablas Hive externas: hive.tcdm.date_dim, hive.tcdm.web_sales_silver (con sus particiones ya sincronizadas) y hive.tcdm.web_sales_by_year_gold.

Si el clúster de infraestructura no está levantado todavía, desde una terminal de tu equipo situada en la raíz de la distribución de sesiones:

Animación: la orden make -C entorno warehouse-up se teclea en una terminal de tu equipo (host)

cd ~/tcdm-public
git pull
make -C entorno warehouse-up

Este notebook se ejecuta en el Jupyter de namenode, que no arranca solo. Como en S2, arráncalo desde otra terminal de tu equipo y déjala abierta mientras trabajes:

Animación: docker exec -it namenode bash se teclea en tu equipo y abre una terminal dentro del contenedor namenode, donde se ejecutan su - luser y jupyter lab

# en tu equipo: entra en namenode
docker exec -it namenode bash
# ya dentro de namenode, como luser: Jupyter con el token fijo «tcdm»
su - luser
jupyter lab --ip=0.0.0.0 --port=8888 --no-browser --IdentityProvider.token=tcdm

La primera orden es la única que se ejecuta en tu equipo; las otras dos, ya dentro del contenedor, arrancan Jupyter con el Python del curso (make -C entorno jupyter hace lo mismo en una sola orden). En Visual Studio Code, conecta este notebook con Select Kernel → Select Another Kernel → Existing Jupyter Server, usando la dirección http://127.0.0.1:8888/lab?token=tcdm y eligiendo el kernel Python 3 (ipykernel). El token es siempre tcdm.

Si algo no funciona —el notebook no conecta, una celda !hdfs no encuentra la orden, un servicio no responde—, consulta «Solución de problemas» en entorno/README.md.

Esta sesión no vuelve a ejecutar S5 ni S6. Si alguna de las tablas Hive anteriores no existe todavía, ejecuta antes esos dos notebooks.

In [ ]:
!hdfs dfs -ls -d /datalake/raw/tpcds/date_dim \
    /datalake/silver/tpcds/web_sales \
    /datalake/gold/tpcds/web_sales_by_year \
    /warehouse

Si alguna de esas rutas no aparece, revisa S2, S5 o S6 antes de continuar. A continuación instalamos las dependencias de este notebook —las mismas que en S6, PySpark y el cliente trino— y comprobamos, ya desde Python, que las tablas Hive que S6 dejó registradas siguen ahí: son la fuente de comparación de toda esta sesión.

Instalar las dependencias de este notebook¶

Como en S4, S5 y S6, el propio kernel instala PySpark sobre sí mismo. La celda siguiente recrea requirements.txt con %%writefile para que este notebook no dependa de ningún otro fichero de la distribución, tanto si lo abres dentro de namenode como si usas «Existing Jupyter Server» desde tu equipo.

In [ ]:
%%writefile requirements.txt
# Dependencias del kernel interactivo de la sesión 7 (Spark local + cliente
# Trino, ambos hablando con el catálogo Iceberg sobre el mismo Hive
# Metastore). La serie de PySpark aceptada por el curso es >4,<4.2, igual que
# en S4, S5 y S6.
pyspark>4,<4.2
# Cliente DB-API 2.0 de Trino: es como este notebook consulta Trino sin usar
# el binario `trino` ni `docker exec` (el kernel corre dentro de namenode, no
# tiene acceso al Docker del host). Se fija la serie mayor para que el tipado
# de columnas y la representación de valores no cambien entre ejecuciones.
trino>=0.330,<1
pandas
In [ ]:
%pip install -q -r requirements.txt
In [ ]:
import re
import subprocess
import sys

import pandas
import pyspark
import trino

print(f"PySpark {pyspark.__version__} con Python {sys.version.split()[0]}.")
print(f"Cliente trino: {trino.__version__}")
print(f"Pandas: {pandas.__version__}")

Abrir la conexión con Trino¶

Reutilizamos exactamente el patrón de S6: el paquete trino de PyPI habla el protocolo HTTP de Trino directamente contra el servicio trino-hdfs de la red hadoop-cluster, sin CLI y sin docker exec. El catálogo por defecto es hive —el mismo de S6—, aunque en el resto del notebook escribiremos siempre los nombres completos (hive.tcdm.<tabla> o iceberg.tcdm.<tabla>) para que cada sentencia se entienda sin depender del contexto por defecto de la conexión.

In [ ]:
import pandas as pd
from trino.dbapi import Connection, Cursor

trino_conn: Connection = trino.dbapi.connect(
    host="trino-hdfs",
    port=8080,
    user="luser",
    catalog="hive",
    schema="tcdm",
)


def run_trino(sql: str) -> pd.DataFrame:
    '''Ejecuta una sentencia en Trino y devuelve el resultado como DataFrame.'''
    cursor: Cursor = trino_conn.cursor()
    cursor.execute(sql)
    rows: list[tuple[object, ...]] = cursor.fetchall()
    columns: list[str] = [description[0] for description in cursor.description or []]
    return pd.DataFrame(rows, columns=columns)


run_trino("SELECT 1 AS ping")
In [ ]:
run_trino("SHOW CATALOGS")

hive e iceberg deben aparecer los dos: son dos catálogos Trino distintos que comparten el mismo Hive Metastore (thrift://hive-metastore:9083) como se ve en entorno/trino-hdfs/etc/catalog/hive.properties y en entorno/trino-hdfs/etc/catalog/iceberg.properties (connector.name=iceberg, iceberg.catalog.type=hive_metastore). El catálogo hive interpreta las tablas como Parquet/ORC/CSV clásicos; el catálogo iceberg interpreta esas mismas entradas de metastore como tablas Iceberg cuando su propiedad table_type lo indica. Son dos lentes sobre el mismo almacén de metadatos, no dos almacenes distintos.

In [ ]:
tablas_s6: pd.DataFrame = run_trino("SHOW TABLES FROM hive.tcdm")
print(tablas_s6)

nombres_s6: set[str] = set(tablas_s6["Table"])
for tabla_esperada in ("date_dim", "web_sales_silver", "web_sales_by_year_gold"):
    assert (
        tabla_esperada in nombres_s6
    ), f"falta hive.tcdm.{tabla_esperada}: ejecuta S6 antes de continuar"
print("Las tres tablas Hive de S6 están registradas.")
In [ ]:
filas_silver: pd.DataFrame = run_trino("SELECT count(*) AS filas FROM hive.tcdm.web_sales_silver")
filas_gold: pd.DataFrame = run_trino(
    "SELECT count(*) AS filas FROM hive.tcdm.web_sales_by_year_gold"
)
print("hive.tcdm.web_sales_silver:", filas_silver.loc[0, "filas"], "filas")
print("hive.tcdm.web_sales_by_year_gold:", filas_gold.loc[0, "filas"], "filas")

assert (
    int(filas_silver.loc[0, "filas"]) > 0
), "web_sales_silver está vacía: repite la sincronización de particiones de S6"
assert int(filas_gold.loc[0, "filas"]) > 0, "web_sales_by_year_gold está vacía: repite S5/S6"
A continuación

Qué añade Iceberg sobre Hive

  • Metadatos de tabla versionados y un snapshot nuevo en cada escritura
  • Evolución de esquema y de particionado sin reescribir el histórico
  • Manifests que planifican la lectura, frente a un MSCK REPAIR TABLE

Por qué no basta el modelo Hive¶

En S6, la tabla hive.tcdm.web_sales_silver se apoya en directorios (sold_year=2000/sold_month=1/...) y en entradas de partición que el metastore conoce sólo después de un paso explícito de sincronización (sync_partition_metadata). Ese modelo resuelve "encontrar los datos", pero dejaba abiertas varias preguntas que no llegamos a plantear en S6: ¿qué pasa si una escritura deja varios ficheros a medio escribir? ¿Cómo sabe un lector qué versión de la tabla está viendo si dos escrituras se solapan? ¿Cómo se añade una columna sin reescribir todo el histórico? ¿Cómo se limpia una versión antigua sin arriesgarse a borrar datos que otro lector todavía necesita?

Iceberg no cambia el almacenamiento (sigue siendo Parquet en HDFS) ni elimina el catálogo (Trino y Spark lo siguen localizando por su nombre en el Hive Metastore). Lo que añade, por encima de ambos, es:

  • metadatos de tabla versionados: cada estado de la tabla tiene un fichero metadata.json propio;
  • cada operación de escritura produce un nuevo snapshot, con su propio conjunto de ficheros de datos;
  • los lectores ven siempre una versión coherente de la tabla, nunca un estado a medio escribir;
  • el esquema puede evolucionar con reglas explícitas (añadir columnas, renombrarlas, ampliar tipos) sin reescribir los datos antiguos;
  • el esquema de particionado puede cambiar sin tener que reorganizar de inmediato todos los ficheros ya escritos;
  • los manifests permiten planificar una lectura a partir de metadatos, sin tener que listar a ciegas todo el almacenamiento como haría un MSCK REPAIR TABLE.

Iceberg frente a Delta Lake y Hudi¶

Iceberg no es el único formato de tabla abierto. Delta Lake y Apache Hudi resuelven el mismo problema —decidir qué ficheros Parquet forman cada versión de una tabla— y los tres ofrecen escrituras atómicas, consulta de versiones anteriores y evolución de esquema. Se diferencian en cómo guardan ese estado y en el uso para el que se diseñaron:

Apache Iceberg Delta Lake Apache Hudi
Origen Netflix Databricks Uber
Dónde guarda el estado de la tabla Árbol de metadatos: metadata.json, lista de manifests y manifests Registro de transacciones _delta_log/: un fichero JSON por escritura y checkpoints Parquet periódicos Línea de tiempo .hoodie/: una entrada por cada acción sobre la tabla
Particionado Oculto, mediante transformaciones como month(ws_sold_date) Columnas de partición declaradas, como en Hive Columnas de partición declaradas, como en Hive
Actualizaciones y borrados Copy-on-write o merge-on-read, configurable en cada tabla Copy-on-write, o merge-on-read mediante vectores de borrado Copy-on-write o merge-on-read, según el tipo de tabla elegido al crearla, con un índice por clave de registro
Conector de Trino Lectura y escritura Lectura y escritura Sólo lectura
Uso habitual Tablas analíticas compartidas entre varios motores Plataformas centradas en Spark y Databricks Ingesta incremental con muchas actualizaciones por clave

Esta sesión y la siguiente trabajan con Iceberg porque la misma tabla se crea desde Trino y se modifica desde Spark, y los dos motores escriben en ese formato. Los modos copy-on-write y merge-on-read de la tabla vuelven en S8, donde se compara su efecto físico sobre una tabla Iceberg.

Recapitulación

El modelo Hive encuentra los datos, pero no los versiona

  • El metastore localiza nombres, esquemas y particiones
  • El formato de tabla es quien decide qué ficheros forman cada versión
A continuación

Crear dos tablas Iceberg con CTAS

  • La fuente Hive externa (/datalake) queda intacta
  • La tabla Iceberg es nueva y administrada, bajo /warehouse
  • month(ws_sold_date) se declara como partición oculta

Crear la tabla Iceberg¶

Empezamos por Trino: es el cliente SQL más sencillo para crear y examinar una tabla Iceberg, sin la sobrecarga de arrancar antes una SparkSession. Más adelante leeremos la misma tabla desde Spark y compararemos.

Vamos a crear dos tablas Iceberg mediante CREATE TABLE ... AS SELECT (CTAS) a partir de las tablas Hive externas que ya registró S6. Un CTAS deja intacta la tabla fuente —hive.tcdm.web_sales_silver sigue apuntando a /datalake/silver/tpcds/web_sales— y crea una tabla nueva y gestionada bajo /warehouse, con sus propios ficheros de datos y sus propios metadatos. Es exactamente el mismo patrón "fuente externa → tabla administrada" que ya viste en S6 con web_sales_by_year_managed_demo, sólo que ahora la tabla administrada no es una tabla Hive más: es una tabla Iceberg.

Estas dos tablas son las que usará también S8 para practicar ingesta continua con Spark Structured Streaming, así que sus nombres y columnas se mantienen exactamente como las dejamos aquí durante el resto del curso.

Por eso los dos CTAS llevan IF NOT EXISTS: si repites esta sesión —por ejemplo, tras reiniciar el kernel— las tablas se conservan tal como estén, incluidos los cambios que S8 haya hecho en ellas, en lugar de fallar o de recrearse. Las comprobaciones que comparan cada tabla con su fuente Hive sólo se aplican cuando la tabla se acaba de crear.

iceberg.tcdm.web_sales¶

Mantenemos las mismas columnas de negocio de web_sales_silver, incluida ws_sold_date (tipo DATE): la necesitamos intacta porque es sobre ella donde aplicaremos la transformación de partición oculta de Iceberg, month(ws_sold_date). La clave compuesta de esta tabla sigue siendo (ws_order_number, ws_item_sk) —es una tabla a nivel de línea de venta, un mismo pedido puede tener varias líneas— y el CTAS que sigue no agrupa ni deduplica filas, así que esa clave permanece única después de crear la tabla.

In [ ]:
run_trino("CREATE SCHEMA IF NOT EXISTS iceberg.tcdm")

tablas_iceberg_previas: set[str] = set(run_trino("SHOW TABLES FROM iceberg.tcdm")["Table"])
web_sales_recien_creada: bool = "web_sales" not in tablas_iceberg_previas
web_sales_by_year_recien_creada: bool = "web_sales_by_year" not in tablas_iceberg_previas

ctas_web_sales_sql: str = '''
CREATE TABLE IF NOT EXISTS iceberg.tcdm.web_sales
WITH (partitioning = ARRAY['month(ws_sold_date)']) AS
SELECT
    ws_order_number,
    ws_item_sk,
    ws_bill_customer_sk,
    ws_bill_addr_sk,
    ws_sold_date_sk,
    ws_sold_date,
    ws_quantity,
    ws_list_price,
    ws_ext_discount_amt,
    ws_ext_sales_price,
    ws_net_paid,
    ws_net_profit,
    sold_year,
    sold_month
FROM hive.tcdm.web_sales_silver
'''
run_trino(ctas_web_sales_sql)
if web_sales_recien_creada:
    print("iceberg.tcdm.web_sales creada.")
else:
    print("iceberg.tcdm.web_sales ya existía: se conserva tal como estaba.")

WITH (partitioning = ARRAY['month(ws_sold_date)']) es la sintaxis real que expone el conector Iceberg de Trino para declarar una transformación de partición: month deriva un valor de partición a partir de ws_sold_date en el momento de escribir, sin que tengamos que crear ni mantener nosotros mismos columnas auxiliares como sold_year/sold_month. Particionado de Iceberg recoge el resto de transformaciones disponibles (year, day, hour, bucket, truncate...).

Comprobamos que el recuento de filas y la clave compuesta se conservan exactamente igual que en la tabla Hive de origen.

In [ ]:
filas_hive: pd.DataFrame = run_trino("SELECT count(*) AS filas FROM hive.tcdm.web_sales_silver")
filas_iceberg: pd.DataFrame = run_trino("SELECT count(*) AS filas FROM iceberg.tcdm.web_sales")
print("hive.tcdm.web_sales_silver:", filas_hive.loc[0, "filas"], "filas")
print("iceberg.tcdm.web_sales:    ", filas_iceberg.loc[0, "filas"], "filas")
if web_sales_recien_creada:
    assert int(filas_hive.loc[0, "filas"]) == int(filas_iceberg.loc[0, "filas"])
else:
    print("La tabla ya existía: puede incluir cambios posteriores a su creación (S8).")

clave_compuesta: pd.DataFrame = run_trino('''
    SELECT
        (SELECT count(*) FROM iceberg.tcdm.web_sales) AS filas,
        (SELECT count(*) FROM (
            SELECT DISTINCT ws_order_number, ws_item_sk FROM iceberg.tcdm.web_sales
        )) AS claves_distintas
''')
print(clave_compuesta)
assert int(clave_compuesta.loc[0, "filas"]) == int(
    clave_compuesta.loc[0, "claves_distintas"]
), "la clave (ws_order_number, ws_item_sk) ha dejado de ser única: revisa el CTAS"
print("La clave (ws_order_number, ws_item_sk) sigue siendo única en iceberg.tcdm.web_sales.")

iceberg.tcdm.web_sales_by_year¶

El agregado de negocio no necesita ninguna transformación de partición: es una tabla pequeña, pensada para leerse entera, exactamente igual que su fuente hive.tcdm.web_sales_by_year_gold.

In [ ]:
ctas_web_sales_by_year_sql: str = '''
CREATE TABLE IF NOT EXISTS iceberg.tcdm.web_sales_by_year AS
SELECT d_year, pedidos, lineas_de_venta, beneficio_neto
FROM hive.tcdm.web_sales_by_year_gold
'''
run_trino(ctas_web_sales_by_year_sql)
if web_sales_by_year_recien_creada:
    print("iceberg.tcdm.web_sales_by_year creada.")
else:
    print("iceberg.tcdm.web_sales_by_year ya existía: se conserva tal como estaba.")
In [ ]:
comparacion_gold: pd.DataFrame = run_trino('''
    SELECT count(*) AS filas, sum(beneficio_neto) AS beneficio_total
    FROM iceberg.tcdm.web_sales_by_year
''')
comparacion_gold_hive: pd.DataFrame = run_trino('''
    SELECT count(*) AS filas, sum(beneficio_neto) AS beneficio_total
    FROM hive.tcdm.web_sales_by_year_gold
''')
print("iceberg.tcdm.web_sales_by_year:    ", comparacion_gold.to_dict("records")[0])
print("hive.tcdm.web_sales_by_year_gold:  ", comparacion_gold_hive.to_dict("records")[0])
if web_sales_by_year_recien_creada:
    assert comparacion_gold.loc[0, "filas"] == comparacion_gold_hive.loc[0, "filas"]
else:
    print("La tabla ya existía: puede incluir cambios posteriores a su creación (S8).")

Un metastore también puede guardar un comentario de tabla para una tabla Iceberg. Lo añadimos con COMMENT ON TABLE y lo comprobamos en el propio SHOW CREATE TABLE — es la contrapartida ejecutable de la buena práctica «documentar el propietario, el esquema y la ubicación» que veremos al final de la sesión, la misma que ya demostramos en S6 con las tablas Hive.

In [ ]:
run_trino('''
    COMMENT ON TABLE iceberg.tcdm.web_sales IS
    'Ventas web de TPC-DS SF1 a nivel de línea, creada por CTAS desde hive.tcdm.web_sales_silver (ver S7).'
''')
run_trino("SHOW CREATE TABLE iceberg.tcdm.web_sales")

A partir de aquí, iceberg.tcdm.web_sales e iceberg.tcdm.web_sales_by_year son las dos tablas de referencia de esta sesión —y las que reutilizará S8—. El resto de tablas Iceberg que crearemos más adelante para provocar ficheros pequeños, experimentar con ordenación física o repetir la evolución de esquema con más libertad se llamarán siempre iceberg.tcdm.web_sales_lab_*, precisamente para que quede claro de un vistazo cuáles son desechables y cuáles no.

Recapitulación

Registrar no es copiar la fuente

  • hive.tcdm.web_sales_silver no se ha tocado en ningún momento
  • iceberg.tcdm.web_sales e iceberg.tcdm.web_sales_by_year son las tablas que reutilizará S8
  • El prefijo web_sales_lab_ marca lo desechable de esta sesión

Observar la definición y la ubicación¶

No debemos suponer dónde ha quedado físicamente la tabla: el catálogo puede elegir un directorio con un sufijo único, como ya viste en S6 con web_sales_by_year_managed_demo. Lo confirmamos con SHOW CREATE TABLE y extraemos la ubicación real con una expresión regular, igual que comprobamos en Python en toda esta sesión.

In [ ]:
ddl_web_sales: pd.DataFrame = run_trino("SHOW CREATE TABLE iceberg.tcdm.web_sales")
ddl_web_sales_texto: str = ddl_web_sales.iloc[0, 0]
print(ddl_web_sales_texto)
In [ ]:
def ubicacion_de_tabla_iceberg(ddl_texto: str) -> str:
    '''Extrae la ruta HDFS declarada en la propiedad location de un SHOW CREATE TABLE.'''
    coincidencia: re.Match[str] | None = re.search(r"location = '([^']+)'", ddl_texto)
    assert coincidencia is not None, "SHOW CREATE TABLE no contiene una propiedad location"
    return coincidencia.group(1)


ubicacion_web_sales: str = ubicacion_de_tabla_iceberg(ddl_web_sales_texto)
print("Ubicación real de iceberg.tcdm.web_sales:", ubicacion_web_sales)
assert ubicacion_web_sales.startswith(
    "hdfs://namenode:9000/warehouse/tcdm.db/web_sales-"
), "la ubicación no sería la esperada bajo /warehouse/tcdm.db/"

El patrón es /warehouse/tcdm.db/web_sales-<identificador>, no simplemente /warehouse/tcdm.db/web_sales: Trino añade un sufijo único para evitar colisiones si en algún momento se crea y se borra una tabla con el mismo nombre. Miramos ahora dentro de esa ruta en HDFS, distinguiendo el directorio data (los ficheros Parquet) del directorio metadata (manifests, snapshots y ficheros *.metadata.json).

Como en S4, la celda !orden siguiente interpola una variable de Python escribiendo {variable} dentro de la orden de shell: así evitamos reconstruir la ruta a mano y nos aseguramos de listar exactamente la ubicación que acaba de confirmar Trino.

In [ ]:
!hdfs dfs -ls -R -h {ubicacion_web_sales}

Deberías distinguir tres cosas en ese listado:

  • data/ws_sold_date_month=2000-01/...parquet y directorios equivalentes para cada mes: son los ficheros de datos. El nombre del directorio contiene el valor de la transformación de partición (month), pero no es una columna que aparezca en el esquema de la tabla ni que el usuario tenga que mantener: es Iceberg quien decide, escribe y luego lee ese directorio a partir de ws_sold_date. Por eso se llama partición oculta.
  • metadata/*.metadata.json: un fichero por cada versión del esquema y la configuración de la tabla.
  • metadata/*.avro y metadata/snap-*.avro: los manifests y la lista de manifests de cada snapshot.

Compáralo mentalmente con lo que ya conoces de hive.tcdm.web_sales_silver en S6: allí los directorios sold_year=2000/sold_month=1/ son columnas de partición reales, declaradas en el CREATE TABLE y sincronizadas explícitamente con sync_partition_metadata. Aquí no hemos declarado ninguna columna de partición adicional ni hemos tenido que sincronizar nada: Iceberg ya sabe qué ficheros pertenecen a cada mes porque lo registra en sus propios manifests, no en el listado de particiones del Hive Metastore.

A continuación

Un segundo cliente del mismo catálogo

  • El catálogo iceberg de Spark, tipo hive, habla con el mismo Thrift que Trino
  • Ninguna ruta física en el código: sólo el nombre de la tabla
  • Misma comparación de negocio, ahora desde Trino y desde Spark

Leer desde Trino y Spark¶

iceberg.tcdm.web_sales ya existe y ya hemos comprobado su contenido con Trino. Toca ahora construir una SparkSession capaz de leer la misma tabla, igual que en S6 hicimos con hive.tcdm.web_sales_silver.

Configurar Spark para Iceberg¶

Partimos de la misma base de S6 —spark.hadoop.hive.metastore.uris, spark.sql.warehouse.dir, spark.sql.catalogImplementation=hive y enableHiveSupport()— y añadimos lo necesario para que Spark entienda el formato de tabla Iceberg y lo resuelva por catálogo, exactamente igual que ya hace con las tablas Hive:

  • spark.jars.packages con la coordenada Maven org.apache.iceberg:iceberg-spark-runtime-4.1_2.13:1.11.0. El sufijo 4.1_2.13 no es arbitrario: identifica la versión de Spark (4.1) y de Scala (2.13) para las que se compiló ese conector, y son exactamente las que fija entorno/spark/Dockerfile (apache/spark:4.1.3-scala2.13-java21-python3-ubuntu, ARG ICEBERG_VERSION=1.11.0) para el contenedor spark-iceberg de la ruta S3 de este mismo entorno. La versión de PySpark que instala requirements.txt (>4,<4.2) resuelve hoy a la 4.1.3, así que usamos la misma coordenada Iceberg aquí.

  • spark.sql.extensions con org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions, el mismo valor que ya declara entorno/spark/spark-defaults.conf para el contenedor spark-iceberg: sin esta extensión, Spark no reconoce sentencias específicas de Iceberg como CALL a procedimientos del sistema o ciertas cláusulas de MERGE INTO que usará S8.

  • un catálogo con nombre iceberg, de tipo hive, apuntando exactamente al mismo Hive Metastore que ya usa Trino para su propio catálogo iceberg:

    • spark.sql.catalog.iceberg = org.apache.iceberg.spark.SparkCatalog
    • spark.sql.catalog.iceberg.type = hive
    • spark.sql.catalog.iceberg.uri = thrift://hive-metastore:9083
    • spark.sql.catalog.iceberg.warehouse = hdfs://namenode:9000/warehouse

Con esto, Spark resuelve iceberg.tcdm.web_sales por su nombre de catálogo exactamente igual que ya resuelve tcdm.web_sales_silver por el catálogo Hive nativo: no hace falta ninguna ruta física ni ningún paso adicional.

El catálogo Iceberg de Spark y el Hive Metastore¶

El catálogo spark.sql.catalog.iceberg de Spark habla con el mismo Hive Metastore que usa Trino porque este entorno fija la versión 3.1.3; la justificación completa —incluyendo por qué no se elige sin más una versión más reciente de la rama 4.x— está en entorno/README.md, sección "Versiones fijadas y compatibilidad".

In [ ]:
from pyspark import SparkContext
from pyspark.sql import DataFrame, SparkSession

HIVE_METASTORE_URIS: str = "thrift://hive-metastore:9083"
WAREHOUSE_DIR: str = "hdfs://namenode:9000/warehouse"
ICEBERG_VERSION: str = "1.11.0"
ICEBERG_SPARK_PACKAGE: str = f"org.apache.iceberg:iceberg-spark-runtime-4.1_2.13:{ICEBERG_VERSION}"

spark: SparkSession = (
    SparkSession.builder.appName("TCDM_S7_Local")
    .master("local[2]")
    .config("spark.driver.memory", "1g")
    .config("spark.sql.shuffle.partitions", "8")
    .config("spark.hadoop.hive.metastore.uris", HIVE_METASTORE_URIS)
    .config("spark.sql.warehouse.dir", WAREHOUSE_DIR)
    .config("spark.sql.catalogImplementation", "hive")
    .config("spark.jars.packages", ICEBERG_SPARK_PACKAGE)
    .config(
        "spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
    )
    .config("spark.sql.catalog.iceberg", "org.apache.iceberg.spark.SparkCatalog")
    .config("spark.sql.catalog.iceberg.type", "hive")
    .config("spark.sql.catalog.iceberg.uri", HIVE_METASTORE_URIS)
    .config("spark.sql.catalog.iceberg.warehouse", WAREHOUSE_DIR)
    .enableHiveSupport()
    .getOrCreate()
)
sc: SparkContext = spark.sparkContext

print(f"Spark version: {spark.version} | Master: {sc.master}")
print(f"Metastore URIs: {HIVE_METASTORE_URIS}")
print(f"Warehouse dir: {WAREHOUSE_DIR}")
print(f"Paquete Iceberg: {ICEBERG_SPARK_PACKAGE}")
print("Catálogo 'iceberg' configurado como tipo hive contra el mismo metastore.")

La primera ejecución de esta celda tarda más de lo habitual: Spark tiene que descargar el jar de iceberg-spark-runtime desde Maven Central a través de Ivy. Las siguientes ejecuciones reutilizan la caché local.

Comprobamos primero que el catálogo Hive nativo de Spark sigue viendo las tablas de S6, igual que en esa sesión.

In [ ]:
spark.sql("SHOW TABLES IN tcdm").show()

web_sales_silver_spark: DataFrame = spark.table("tcdm.web_sales_silver")
web_sales_silver_spark.printSchema()
print("Filas en tcdm.web_sales_silver (Spark):", web_sales_silver_spark.count())

Leer la tabla Iceberg por su nombre de catálogo¶

Con el catálogo spark.sql.catalog.iceberg ya configurado, Spark resuelve iceberg.tcdm.web_sales exactamente igual que resuelve tcdm.web_sales_silver por el catálogo Hive nativo: sin construir ninguna ruta a mano ni depender de dónde ha decidido colocar Trino los ficheros físicamente.

In [ ]:
web_sales_spark: DataFrame = spark.table("iceberg.tcdm.web_sales")
web_sales_spark.printSchema()

Repetimos, con Spark, la misma comparación de negocio que ya distinguía S2: líneas de venta (count(*)) frente a pedidos distintos (count(DISTINCT ws_order_number)), y comprobamos que Trino y Spark devuelven exactamente el mismo resultado sobre la misma tabla Iceberg.

In [ ]:
resultado_trino: pd.DataFrame = run_trino('''
    SELECT count(*) AS lineas_de_venta, count(DISTINCT ws_order_number) AS pedidos
    FROM iceberg.tcdm.web_sales
''')
resultado_spark: pd.DataFrame = web_sales_spark.selectExpr(
    "count(*) AS lineas_de_venta", "count(DISTINCT ws_order_number) AS pedidos"
).toPandas()

print("Trino:", resultado_trino.to_dict("records")[0])
print("Spark:", resultado_spark.to_dict("records")[0])

assert resultado_trino.loc[0, "lineas_de_venta"] == resultado_spark.loc[0, "lineas_de_venta"]
assert resultado_trino.loc[0, "pedidos"] == resultado_spark.loc[0, "pedidos"]
print("Trino y Spark coinciden sobre iceberg.tcdm.web_sales.")

Particionado Hive frente a particionado Iceberg¶

En S6, hive.tcdm.web_sales_silver tiene columnas de partición visibles: sold_year y sold_month aparecen en el DESCRIBE, en cada fila devuelta por SELECT *, y sus valores están codificados en los nombres de los directorios (sold_year=2000/sold_month=1/). El metastore sólo conoce esas particiones después de un paso explícito de descubrimiento (sync_partition_metadata), y si alguien cambiara el esquema de particionado tendría que reescribir esa organización de directorios y repetir el descubrimiento.

En iceberg.tcdm.web_sales, la partición es una transformación (month(ws_sold_date)) sobre una columna que ya existía: no hemos añadido ninguna columna nueva al esquema lógico de la tabla, y no hace falta ningún paso de sincronización porque Iceberg registra en sus propios manifests, en el momento de cada escritura, qué fichero pertenece a qué valor de partición. Quien filtra por ws_sold_date no necesita saber que esa columna está particionada por mes: el motor decide solo qué ficheros puede descartar.

Comparamos el plan de una misma consulta —filtrar por el primer mes de 2000— en las dos tablas.

In [ ]:
explain_hive_particionada: pd.DataFrame = run_trino('''
    EXPLAIN SELECT count(*) FROM hive.tcdm.web_sales_silver
    WHERE sold_year = 2000 AND sold_month = 1
''')
print("\n".join(explain_hive_particionada.iloc[:, 0].tolist()))
In [ ]:
explain_iceberg_transformada: pd.DataFrame = run_trino('''
    EXPLAIN SELECT count(*) FROM iceberg.tcdm.web_sales
    WHERE ws_sold_date >= DATE '2000-01-01' AND ws_sold_date < DATE '2000-02-01'
''')
print("\n".join(explain_iceberg_transformada.iloc[:, 0].tolist()))

En el primer plan, la poda ocurre sobre las columnas de partición declaradas (sold_year, sold_month), que el usuario tuvo que escribir en el WHERE sabiendo que existían como columnas de partición. En el segundo, filtramos únicamente por ws_sold_date —la columna de negocio, la misma que ya usábamos en S2 y S4— y el planificador de Trino traduce ese rango de fechas al valor de partición month que corresponde, sin que hayamos tenido que mencionar ninguna columna de partición en la consulta. Busca en ambos planes la parte del TableScan/IcebergScan para confirmar que el filtrado ocurre antes de leer los ficheros de datos y no como un Filter posterior.

Repetimos la comparación con Spark, ahora que ya tenemos las dos tablas disponibles en la sesión.

In [ ]:
spark.sql('''
    SELECT count(*) FROM tcdm.web_sales_silver
    WHERE sold_year = 2000 AND sold_month = 1
''').explain("formatted")
In [ ]:
web_sales_spark.filter(
    "ws_sold_date >= DATE'2000-01-01' AND ws_sold_date < DATE'2000-02-01'"
).selectExpr("count(*)").explain("formatted")

No se debe concluir de aquí que Iceberg "hace innecesario" pensar en particiones: una transformación mal elegida —por ejemplo, particionar por ws_order_number o por bucket con un número de cubos excesivo— puede seguir produciendo demasiados ficheros pequeños, exactamente el problema que mediremos más adelante en esta misma sesión. Lo que cambia es quién mantiene la relación entre el valor lógico y el directorio físico: en Hive, el usuario; en Iceberg, el propio formato de tabla.

Recapitulación

Partición visible frente a partición oculta

  • En Hive el usuario declara sold_year/sold_month y sincroniza el metastore
  • En Iceberg el formato registra month(ws_sold_date) directamente en los manifests
  • Una transformación mal elegida sigue creando demasiados ficheros pequeños
Pregunta guía
  • ¿Por qué no experimentamos con INSERT ni ALTER TABLE directamente sobre iceberg.tcdm.web_sales?
  • ¿Qué contrato con S8 protege la clave compuesta (ws_order_number, ws_item_sk)?

Snapshots y evolución¶

iceberg.tcdm.web_sales ya tiene un primer snapshot: el que produjo el propio CTAS. Lo consultamos con la tabla de metadatos $snapshots, la misma convención <tabla>$metadato que ya usaste en S6 para $partitions.

In [ ]:
snapshots_web_sales: pd.DataFrame = run_trino('''
    SELECT snapshot_id, committed_at, operation
    FROM iceberg.tcdm."web_sales$snapshots"
    ORDER BY committed_at
''')
snapshots_web_sales

Un único snapshot, con operation = 'append' (así etiqueta Iceberg un CTAS: es, en el fondo, una escritura que añade el primer conjunto de ficheros). Si repites esta sesión después de S8, verás además los snapshots que dejó el streaming. Antes de seguir experimentando con inserciones, cambios de esquema y viajes en el tiempo, vamos a decidir dónde hacerlo.

Por qué no seguimos experimentando sobre iceberg.tcdm.web_sales a partir de aquí. El ejemplo más simple de una nueva tanda de datos sería reinsertar un subconjunto ya existente de la propia tabla. Pero iceberg.tcdm.web_sales tiene un contrato con S8: su clave compuesta (ws_order_number, ws_item_sk) debe seguir siendo única, porque S8 hará un MERGE INTO sobre exactamente esa pareja de columnas. Reinsertar filas que ya existen —aunque sea sólo para ver aparecer un segundo snapshot— rompería esa unicidad para el resto del curso. Por eso, para la parte de la sesión que experimenta con nuevas versiones, columnas nuevas y viajes en el tiempo, trabajamos sobre una copia de laboratorio, iceberg.tcdm.web_sales_lab_evolution, creada con el mismo CTAS y que podemos manipular con total libertad sin afectar a la tabla que reutilizará S8. iceberg.tcdm.web_sales queda así con un único snapshot, exactamente el que dejó su creación, y su clave compuesta intacta.

In [ ]:
run_trino("DROP TABLE IF EXISTS iceberg.tcdm.web_sales_lab_evolution")
run_trino(
    "CREATE TABLE iceberg.tcdm.web_sales_lab_evolution AS SELECT * FROM iceberg.tcdm.web_sales"
)
print("iceberg.tcdm.web_sales_lab_evolution creada como copia de laboratorio.")

snapshots_lab_inicial: pd.DataFrame = run_trino('''
    SELECT snapshot_id, operation FROM iceberg.tcdm."web_sales_lab_evolution$snapshots"
''')
snapshots_lab_inicial

Insertar una nueva tanda de datos¶

Insertamos un pedido nuevo, con una pareja (ws_order_number, ws_item_sk) que no existía todavía, para que el snapshot que se genera represente una escritura real de datos nuevos y no una duplicación de la clave. En lugar de inventar un número de pedido a mano —que podría coincidir por casualidad con uno real de TPC-DS—, calculamos uno que no puede colisionar: max(ws_order_number) + 1 sobre la propia tabla.

Esta fila de laboratorio copia ws_sold_date_sk de una fila real sin recalcularla para la fecha 2000-06-15 que fijamos a mano: la clave surrogate y la fecha derivada no se corresponden. No afecta a esta tabla de laboratorio —sólo sirve para observar el efecto del INSERT en los snapshots—, pero no la tomes como ejemplo del contrato ws_sold_date/ws_sold_date_sk que sí se respeta en iceberg.tcdm.web_sales.

In [ ]:
nuevo_order_number = run_trino(
    "SELECT max(ws_order_number) + 1 AS siguiente FROM iceberg.tcdm.web_sales_lab_evolution"
).loc[0, "siguiente"]
print("Nuevo ws_order_number, garantizado libre:", nuevo_order_number)

run_trino(f'''
    INSERT INTO iceberg.tcdm.web_sales_lab_evolution
    SELECT
        {nuevo_order_number} AS ws_order_number,
        ws_item_sk, ws_bill_customer_sk, ws_bill_addr_sk, ws_sold_date_sk,
        DATE '2000-06-15' AS ws_sold_date,
        ws_quantity, ws_list_price, ws_ext_discount_amt, ws_ext_sales_price,
        ws_net_paid, ws_net_profit,
        2000 AS sold_year, 6 AS sold_month
    FROM iceberg.tcdm.web_sales_lab_evolution
    LIMIT 1
''')

snapshots_tras_insert: pd.DataFrame = run_trino('''
    SELECT snapshot_id, committed_at, operation
    FROM iceberg.tcdm."web_sales_lab_evolution$snapshots"
    ORDER BY committed_at
''')
snapshots_tras_insert
In [ ]:
assert (
    len(snapshots_tras_insert) == len(snapshots_lab_inicial) + 1
), "debería haber un snapshot más tras el INSERT"
snapshot_inicial_id: int = int(snapshots_tras_insert.loc[0, "snapshot_id"])
snapshot_tras_insert_id: int = int(
    snapshots_tras_insert.loc[len(snapshots_tras_insert) - 1, "snapshot_id"]
)
print("Snapshot inicial:      ", snapshot_inicial_id)
print("Snapshot tras INSERT:  ", snapshot_tras_insert_id)

Evolución de esquema¶

Añadimos una columna nueva y compatible —canal_venta, para anotar en el futuro por qué canal llegó cada línea de venta—. Es una evolución segura: no cambia ni renombra ninguna columna existente, y las filas ya escritas simplemente no tienen valor en la columna nueva.

In [ ]:
run_trino("ALTER TABLE iceberg.tcdm.web_sales_lab_evolution ADD COLUMN canal_venta VARCHAR")
run_trino("DESCRIBE iceberg.tcdm.web_sales_lab_evolution")
In [ ]:
filas_tras_evolucion: pd.DataFrame = run_trino('''
    SELECT count(*) AS filas, count(canal_venta) AS filas_con_canal
    FROM iceberg.tcdm.web_sales_lab_evolution
''')
print(filas_tras_evolucion)
assert (
    int(filas_tras_evolucion.loc[0, "filas_con_canal"]) == 0
), "las filas antiguas no deberían tener valor en la columna nueva"
print("Las filas anteriores a la evolución de esquema se siguen leyendo, con canal_venta a NULL.")

Consultar una versión anterior¶

Con el identificador del snapshot inicial —el que capturamos antes del INSERT—, Trino permite viajar en el tiempo con FOR VERSION AS OF. Es la prueba de que Iceberg no ha "editado" la tabla en el sitio: ha añadido un snapshot nuevo, y el anterior sigue siendo legible tal cual era, con su esquema y sus filas de entonces.

In [ ]:
conteo_snapshot_inicial: pd.DataFrame = run_trino(f'''
    SELECT count(*) AS filas FROM iceberg.tcdm.web_sales_lab_evolution
    FOR VERSION AS OF {snapshot_inicial_id}
''')
conteo_actual: pd.DataFrame = run_trino(
    "SELECT count(*) AS filas FROM iceberg.tcdm.web_sales_lab_evolution"
)

print("Filas en el snapshot inicial:", conteo_snapshot_inicial.loc[0, "filas"])
print("Filas en la versión actual:  ", conteo_actual.loc[0, "filas"])
assert int(conteo_actual.loc[0, "filas"]) == int(conteo_snapshot_inicial.loc[0, "filas"]) + 1

También se puede viajar en el tiempo por fecha con FOR TIMESTAMP AS OF <timestamp>, útil cuando no se ha guardado el identificador del snapshot pero sí se sabe aproximadamente cuándo se escribió la versión que se busca. Ambas cláusulas son sintaxis real de Trino sobre el conector Iceberg, no HiveQL.

La misma exploración desde Spark¶

Leemos la misma copia de laboratorio directamente por su nombre de catálogo, iceberg.tcdm.web_sales_lab_evolution —Spark ve automáticamente tanto el INSERT como el ALTER TABLE, porque los dos ya están comprometidos en el catálogo— y consultamos su tabla de metadatos de snapshots con la sintaxis de Spark para tablas de catálogo: <catalogo>.<esquema>.<tabla>.snapshots.

In [ ]:
lab_evolution_spark: DataFrame = spark.table("iceberg.tcdm.web_sales_lab_evolution")
lab_evolution_spark.printSchema()
print("Filas (Spark):", lab_evolution_spark.count())
In [ ]:
snapshots_spark: DataFrame = spark.table("iceberg.tcdm.web_sales_lab_evolution.snapshots")
snapshots_spark.select("snapshot_id", "operation").show(truncate=False)

Y viajamos en el tiempo desde Spark con la opción versionAsOf del lector Iceberg —la misma opción nativa que usa Spark para cualquier tabla que soporte viajar en el tiempo, no una opción propia de Iceberg— apuntando siempre al mismo identificador de tabla del catálogo. Iceberg admite un identificador de snapshot o el nombre de una rama/etiqueta como valor de versionAsOf; las opciones snapshot-id/as-of-timestamp que definía el propio conector Iceberg han quedado obsoletas a favor de esta sintaxis común de Spark.

In [ ]:
spark_snapshot_inicial: DataFrame = spark.read.option("versionAsOf", snapshot_inicial_id).table(
    "iceberg.tcdm.web_sales_lab_evolution"
)
print("Filas en el snapshot inicial (Spark):", spark_snapshot_inicial.count())
assert spark_snapshot_inicial.count() == int(conteo_snapshot_inicial.loc[0, "filas"])
print("Trino y Spark coinciden al viajar en el tiempo al mismo snapshot.")
Recapitulación

Cada escritura añade un snapshot

  • El CTAS inicial y cada INSERT producen snapshots append
  • ALTER TABLE ADD COLUMN no reescribe el histórico de datos
  • Trino (FOR VERSION AS OF) y Spark (versionAsOf) leen exactamente el mismo estado antiguo
A continuación

Medir sin cronómetro

  • Proyección y filtros, tres variantes de particionado, estadísticas, ficheros pequeños y ordenación
  • La evidencia es el plan de ejecución y el inventario de ficheros, no el reloj

Metodología de medición¶

A partir de aquí la sesión cambia de pregunta. Ya existe una tabla Iceberg, ya sabemos crearla, leerla desde dos motores y hacerla evolucionar. Toca decidir qué hace falta para que sea eficiente, y comprobarlo, no sólo suponerlo.

TPC-DS SF1 es lo bastante grande para trabajar con datos realistas durante el curso, pero puede ser demasiado pequeño —y este entorno, al ejecutarse en contenedores Docker compartiendo máquina con el resto del portátil o del servidor de sesiones, demasiado variable— para que el tiempo de pared de una consulta sea una medida fiable. Ya lo advertían S4, S5 y S6 con sus propias consultas, y aquí no es distinto: no vamos a apoyar ninguna conclusión de esta sesión en "esta consulta tardó menos". En su lugar, para cada experimento dejamos constancia de:

  • la consulta y los predicados exactos usados;
  • el número y tamaño de los ficheros de datos involucrados;
  • el número de particiones y de manifests;
  • el plan lógico y físico (EXPLAIN en Trino, .explain("formatted") en Spark);
  • las estadísticas disponibles antes y después de cada cambio;
  • el resultado de negocio, para confirmar que la optimización no ha alterado los números;
  • la configuración relevante de Spark o Trino que hace posible la comparación.

Proyección y filtros¶

Antes de tocar particiones, compactación u ordenación, la primera comprobación es la más básica: leer sólo las columnas necesarias y aplicar el filtro lo antes posible. Retomamos la comparación de S2 y S4 entre SELECT * y una proyección estrecha con filtro, ahora sobre iceberg.tcdm.web_sales.

In [ ]:
explain_select_estrella: pd.DataFrame = run_trino("EXPLAIN SELECT * FROM iceberg.tcdm.web_sales")
print("\n".join(explain_select_estrella.iloc[:, 0].tolist()))
In [ ]:
explain_proyeccion_filtro: pd.DataFrame = run_trino('''
    EXPLAIN SELECT ws_order_number, ws_net_paid
    FROM iceberg.tcdm.web_sales
    WHERE ws_sold_date >= DATE '2000-01-01' AND ws_sold_date < DATE '2000-02-01'
''')
print("\n".join(explain_proyeccion_filtro.iloc[:, 0].tolist()))

Fíjate en la lista de columnas (outputSymbols o equivalente) del TableScan/IcebergScan de cada plan: el segundo debería mostrar sólo ws_order_number y ws_net_paid, además del predicado sobre ws_sold_date. Esa reducción ocurre gracias al formato columnar de Parquet, a las estadísticas por fichero que ya guarda Iceberg en sus manifests (mínimo, máximo y nulos por columna) y al propio planificador de Trino: no es una técnica exclusiva de Iceberg, pero Iceberg la aprovecha mejor porque sus manifests ya incluyen estadísticas por fichero sin tener que abrir cada Parquet para leer su footer. No tiene sentido saltar directamente a compactación u ordenación si todavía se están leyendo todas las columnas o todas las particiones de la tabla.

Repetimos la comparación con Spark.

In [ ]:
web_sales_spark.explain("formatted")
In [ ]:
web_sales_spark.select("ws_order_number", "ws_net_paid").filter(
    "ws_sold_date >= DATE'2000-01-01' AND ws_sold_date < DATE'2000-02-01'"
).explain("formatted")

Particiones de datos¶

Ya comparamos antes el particionado visible de Hive con la transformación oculta de Iceberg. Cerramos la comparación añadiendo una tercera variante: el Parquet sin particionar de la capa raw, tal y como S2 lo dejó en /datalake/raw/tpcds/web_sales. Registramos esa ruta como una tabla externa efímera —sólo para esta comparación— para poder consultarla con SQL.

In [ ]:
run_trino("DROP TABLE IF EXISTS hive.tcdm.web_sales_raw_sin_particionar")
run_trino('''
    CREATE TABLE hive.tcdm.web_sales_raw_sin_particionar (
        ws_sold_date_sk BIGINT,
        ws_order_number BIGINT,
        ws_item_sk BIGINT
    ) WITH (
        format = 'PARQUET',
        external_location = 'hdfs://namenode:9000/datalake/raw/tpcds/web_sales'
    )
''')
print(
    "hive.tcdm.web_sales_raw_sin_particionar registrada (sólo columnas necesarias para esta comparación)."
)

Comparamos ahora el plan de una consulta equivalente —contar líneas de venta del año 2000— sobre las tres variantes. La tabla raw no tiene ninguna columna de fecha ya calculada como ws_sold_date, así que el filtro tiene que apoyarse en el rango de ws_sold_date_sk que corresponde a ese año (podemos obtenerlo de date_dim, la tabla que ya registró S6).

In [ ]:
rango_2000: pd.DataFrame = run_trino('''
    SELECT min(d_date_sk) AS primero, max(d_date_sk) AS ultimo
    FROM hive.tcdm.date_dim WHERE d_year = 2000
''')
primer_sk_2000: int = int(rango_2000.loc[0, "primero"])
ultimo_sk_2000: int = int(rango_2000.loc[0, "ultimo"])
print(f"Claves d_date_sk de 2000: {primer_sk_2000}..{ultimo_sk_2000}")
In [ ]:
explain_raw: pd.DataFrame = run_trino(f'''
    EXPLAIN SELECT count(*) FROM hive.tcdm.web_sales_raw_sin_particionar
    WHERE ws_sold_date_sk BETWEEN {primer_sk_2000} AND {ultimo_sk_2000}
''')
print("=== Parquet sin particionar (raw) ===")
print("\n".join(explain_raw.iloc[:, 0].tolist()))
In [ ]:
explain_hive_2000: pd.DataFrame = run_trino('''
    EXPLAIN SELECT count(*) FROM hive.tcdm.web_sales_silver WHERE sold_year = 2000
''')
print("=== Hive particionada por sold_year/sold_month (silver, S6) ===")
print("\n".join(explain_hive_2000.iloc[:, 0].tolist()))
In [ ]:
explain_iceberg_2000: pd.DataFrame = run_trino('''
    EXPLAIN SELECT count(*) FROM iceberg.tcdm.web_sales
    WHERE ws_sold_date >= DATE '2000-01-01' AND ws_sold_date < DATE '2001-01-01'
''')
print("=== Iceberg con transformación month(ws_sold_date) ===")
print("\n".join(explain_iceberg_2000.iloc[:, 0].tolist()))

En el primer plan no hay poda de particiones: el catálogo no sabe nada de la distribución por fecha de web_sales en raw, así que Trino considera todos sus ficheros — aunque sí puede seguir usando las estadísticas de footer de cada Parquet para descartar row groups dentro de un fichero, que no es lo mismo que descartar directorios o ficheros enteros sin abrirlos. En el segundo plan, la poda ocurre sobre las particiones Hive sold_year/sold_month que ya conoces de S6. En el tercero, la poda ocurre sobre la transformación month(ws_sold_date), sin que hayamos tenido que declarar ninguna columna de partición.

Una advertencia que conviene no perder de vista: particionar por una clave de alta cardinalidad —un ws_order_number o un bucket con demasiados cubos para el volumen de datos de SF1— puede producir tantas particiones pequeñas que el coste de planificación (leer y combinar muchos manifests pequeños) supere al beneficio de la poda. La transformación temporal que hemos elegido aquí (month) es razonable porque el número de meses de TPC-DS SF1 es pequeño y coincide con cómo se suelen filtrar los datos de ventas en las consultas de negocio.

Retiramos la tabla temporal raw que hemos usado sólo para esta comparación: es una definición de catálogo efímera, y DROP TABLE sobre una tabla externa no toca los Parquet de /datalake/raw.

In [ ]:
run_trino("DROP TABLE hive.tcdm.web_sales_raw_sin_particionar")
print("Definición temporal retirada; /datalake/raw/tpcds/web_sales no se ha visto afectado.")

Estadísticas¶

Las estadísticas ayudan al planificador a estimar filas, tamaños y el coste de un JOIN, pero conviene no confundir varias capas: las estadísticas de los footers de Parquet (mínimo, máximo, nulos por columna, por fichero), las estadísticas de tabla/partición del Hive Metastore que ya viste en S6 con ANALYZE, y las estadísticas de tabla propias de Iceberg, que agregan la información de los manifests. ANALYZE sobre una tabla Iceberg en Trino alimenta esta última capa.

No prejuzgamos el resultado: comparamos el plan antes y después, y explicamos si la información adicional cambia realmente la decisión del motor, en lugar de asumir que siempre lo hace.

In [ ]:
stats_antes: pd.DataFrame = run_trino("SHOW STATS FOR iceberg.tcdm.web_sales")
stats_antes

Antes de ejecutar ANALYZE, capturamos también el plan de un JOIN con date_dim —la misma unión que ya usábamos en S2 para relacionar ws_sold_date_sk con el año de negocio—, para tener con qué comparar después.

In [ ]:
explain_join_antes_texto: str = "\n".join(run_trino('''
    EXPLAIN SELECT w.ws_order_number, d.d_year
    FROM iceberg.tcdm.web_sales w
    JOIN hive.tcdm.date_dim d ON w.ws_sold_date_sk = d.d_date_sk
    WHERE d.d_year = 2000
''').iloc[:, 0].tolist())
print(explain_join_antes_texto)
In [ ]:
run_trino("ANALYZE iceberg.tcdm.web_sales")
stats_despues: pd.DataFrame = run_trino("SHOW STATS FOR iceberg.tcdm.web_sales")
stats_despues
In [ ]:
explain_join_despues_texto: str = "\n".join(run_trino('''
    EXPLAIN SELECT w.ws_order_number, d.d_year
    FROM iceberg.tcdm.web_sales w
    JOIN hive.tcdm.date_dim d ON w.ws_sold_date_sk = d.d_date_sk
    WHERE d.d_year = 2000
''').iloc[:, 0].tolist())
print(explain_join_despues_texto)

print()
print(
    "Planes idénticos antes y después de ANALYZE:",
    explain_join_antes_texto == explain_join_despues_texto,
)

Compara stats_antes y stats_despues: antes de ANALYZE, columnas como distinct_values_count o data_size suelen venir vacías o son estimaciones muy toscas basadas sólo en el tamaño de los ficheros; después, cada columna trae una estimación calculada sobre los datos reales, y aparece row_count en la fila final. Ahora bien: en una tabla tan pequeña como la de esta sesión (SF1, unos cientos de miles de líneas de venta), es perfectamente posible que explain_join_antes_texto y explain_join_despues_texto resulten idénticos —Trino puede elegir el mismo tipo de JOIN con o sin estadísticas si el volumen de datos ya es pequeño de por sí—. Si es tu caso, esa es precisamente la lección: las estadísticas informan al optimizador, no lo obligan a cambiar de plan, y ANALYZE no es una operación que "siempre acelera la consulta" sino una que "siempre mejora la información disponible", que son cosas distintas.

Recapitulación

ANALYZE informa, no obliga

  • SHOW STATS antes de ANALYZE: estimaciones toscas o vacías
  • Después: una estimación por columna y row_count en la fila final
  • Que el plan del JOIN no cambie también es una lección válida
A continuación

Mantenimiento físico de la tabla

  • Muchos INSERT pequeños dejan muchos ficheros de datos, manifests y snapshots
  • ALTER TABLE ... EXECUTE optimize compacta, pero no borra las versiones antiguas
  • expire_snapshots sí las borra: sólo sobre la tabla de laboratorio, y Trino exige por defecto 7 días de retención
  • Después: ordenación física y Z-order

Pequeños ficheros y compactación¶

Muchas escrituras pequeñas aumentan el coste de listar, planificar y abrir ficheros; en Iceberg también aumentan el número de manifests y el volumen de metadatos que hay que leer para planificar una consulta. Para hacer visible el problema sin arriesgar ninguna de las dos tablas de referencia, creamos una tabla de laboratorio y la llenamos con varias tandas deliberadamente pequeñas: iceberg.tcdm.web_sales_lab_smallfiles.

In [ ]:
run_trino("DROP TABLE IF EXISTS iceberg.tcdm.web_sales_lab_smallfiles")
run_trino('''
    CREATE TABLE iceberg.tcdm.web_sales_lab_smallfiles (
        ws_order_number BIGINT,
        ws_item_sk BIGINT,
        ws_net_paid DECIMAL(7, 2)
    )
''')

for _tanda in range(8):
    run_trino('''
        INSERT INTO iceberg.tcdm.web_sales_lab_smallfiles
        SELECT ws_order_number, ws_item_sk, ws_net_paid
        FROM iceberg.tcdm.web_sales LIMIT 3
    ''')
print("8 tandas de 3 filas insertadas por separado en iceberg.tcdm.web_sales_lab_smallfiles.")
print("Cada INSERT independiente produce su propio fichero de datos y su propio snapshot,")
print("que es exactamente lo que queremos observar aquí: no importa que las filas se repitan,")
print("esto es una tabla de laboratorio pensada sólo para medir ficheros y manifests.")

Cada INSERT produce su propio fichero de datos y su propio snapshot. Medimos antes de compactar: número de ficheros de datos, número de manifests y número de snapshots.

In [ ]:
def inventario_tabla_iceberg(esquema: str, tabla: str) -> dict[str, int]:
    '''Recoge un pequeño inventario de una tabla Iceberg vía las tablas de metadatos de Trino.'''
    n_ficheros = run_trino(f'SELECT count(*) AS n FROM iceberg.{esquema}."{tabla}$files"').loc[
        0, "n"
    ]
    n_manifests = run_trino(f'SELECT count(*) AS n FROM iceberg.{esquema}."{tabla}$manifests"').loc[
        0, "n"
    ]
    n_snapshots = run_trino(f'SELECT count(*) AS n FROM iceberg.{esquema}."{tabla}$snapshots"').loc[
        0, "n"
    ]
    return {
        "ficheros": int(n_ficheros),
        "manifests": int(n_manifests),
        "snapshots": int(n_snapshots),
    }


inventario_antes_compactar: dict[str, int] = inventario_tabla_iceberg(
    "tcdm", "web_sales_lab_smallfiles"
)
print("Antes de compactar:", inventario_antes_compactar)

Trino ofrece ALTER TABLE ... EXECUTE optimize para reescribir ficheros pequeños en tablas Iceberg. El conector Iceberg de Trino documenta este procedimiento; Spark ofrece la acción equivalente a través del procedimiento rewrite_data_files de su catálogo Iceberg —se cita más abajo, en la sección de Z-order, porque es la única de las dos vías que admite esa estrategia—. Aquí seguimos con Trino porque es el motor que ya ha creado y sincronizado esta tabla de laboratorio, y la compactación por tamaño (strategy => 'binpack', la que activa EXECUTE optimize sin más parámetros) se comporta igual desde cualquiera de los dos catálogos.

In [ ]:
resultado_optimize: pd.DataFrame = run_trino(
    "ALTER TABLE iceberg.tcdm.web_sales_lab_smallfiles EXECUTE optimize"
)
resultado_optimize
In [ ]:
inventario_despues_compactar: dict[str, int] = inventario_tabla_iceberg(
    "tcdm", "web_sales_lab_smallfiles"
)
print("Antes de compactar:  ", inventario_antes_compactar)
print("Después de compactar:", inventario_despues_compactar)

assert inventario_despues_compactar["ficheros"] < inventario_antes_compactar["ficheros"]
print(
    f"De {inventario_antes_compactar['ficheros']} ficheros de datos a {inventario_despues_compactar['ficheros']}."
)

EXECUTE optimize no borra las versiones antiguas: reescribe los datos en ficheros más grandes y añade un snapshot nuevo con operation = 'replace', pero los ficheros pequeños originales —y los snapshots que los referenciaban— siguen existiendo hasta que se expiran explícitamente. Lo comprobamos en la propia tabla de snapshots.

In [ ]:
snapshots_lab_smallfiles: pd.DataFrame = run_trino('''
    SELECT snapshot_id, operation FROM iceberg.tcdm."web_sales_lab_smallfiles$snapshots"
    ORDER BY committed_at
''')
snapshots_lab_smallfiles

Manifests y snapshots¶

Los manifests y los snapshots no son ficheros temporales que se puedan ignorar: forman parte del modelo de datos de Iceberg y de la trazabilidad de la tabla. La operación de mantenimiento que sí elimina versiones antiguas es expire_snapshots, y la aplicamos únicamente sobre la tabla de laboratorio que acabamos de compactar, nunca sobre iceberg.tcdm.web_sales ni sobre iceberg.tcdm.web_sales_by_year.

Al intentar expirar snapshots con una retención de cero días, Trino rechaza la operación:

Retention specified (0.00d) is shorter than the minimum retention
configured in the system (7.00d). Minimum retention can be changed with
iceberg.expire-snapshots.min-retention configuration property or
iceberg.expire_snapshots_min_retention session property

Ese límite de 7 días no es un capricho: existe para que un mantenimiento lanzado sin pensar no borre por accidente una versión que otro proceso —posiblemente uno de larga duración, como un INSERT OVERWRITE distribuido o un lector con una consulta abierta desde hace horas— todavía necesita leer. En un entorno docente, donde queremos ver el efecto en la misma sesión, hace falta bajar deliberadamente ese mínimo con una propiedad de sesión, sabiendo que en una tabla real de producción este mismo gesto debería ir acompañado de una política de retención explícita y de la certeza de que ningún lector depende ya de las versiones que se van a borrar.

In [ ]:
run_trino("SET SESSION iceberg.expire_snapshots_min_retention = '0s'")
resultado_expire: pd.DataFrame = run_trino('''
    ALTER TABLE iceberg.tcdm.web_sales_lab_smallfiles
    EXECUTE expire_snapshots(retention_threshold => '0d')
''')
resultado_expire
In [ ]:
snapshots_tras_expirar: pd.DataFrame = run_trino('''
    SELECT snapshot_id, operation FROM iceberg.tcdm."web_sales_lab_smallfiles$snapshots"
''')
print(f"Snapshots antes de expirar: {len(snapshots_lab_smallfiles)}")
print(f"Snapshots después de expirar: {len(snapshots_tras_expirar)}")
assert len(snapshots_tras_expirar) < len(snapshots_lab_smallfiles)
snapshots_tras_expirar

Tras expire_snapshots, sólo queda el snapshot vigente: los que apuntaban a los ocho ficheros pequeños ya no son alcanzables desde el historial de la tabla. remove_orphan_files es el procedimiento complementario para el caso menos habitual de ficheros huérfanos —escritos por una operación interrumpida a medio camino— que ningún snapshot llegó a referenciar nunca; comparte la misma protección de retención mínima que expire_snapshots por el mismo motivo.

A continuación

Ordenación física y Z-order

  • ORDER BY ordena el resultado de una consulta, no los ficheros en disco
  • Un sort order de tabla organiza los datos dentro de cada fichero y aprovecha el mínimo/máximo por fichero
  • Es una pista para el escritor y el planificador, no una garantía sobre el orden de lectura
  • Z-order combina varias columnas de filtrado: aquí queda como ampliación documentada, no ejecutada

Ordenación física¶

Conviene distinguir con cuidado tres cosas que suenan parecidas:

  • ORDER BY en una consulta ordena únicamente el resultado que se entrega al cliente; no cambia en nada cómo están organizados los ficheros en disco.
  • Un sort order de tabla en Iceberg describe cómo se espera que se organicen los datos dentro de cada fichero cuando se escriben, ayudando a lecturas selectivas por esa columna gracias a las estadísticas de mínimo/máximo por fichero.
  • Ordenación local frente a global: ordenar dentro de cada tarea (local) es mucho más barato que garantizar un orden global entre todos los ficheros de la tabla, y normalmente basta con la primera para el beneficio que buscamos.

La documentación de DDL de Iceberg en Spark es explícita en que el orden de escritura no garantiza el orden de las filas que devuelve una consulta posterior: un sort order es una pista para el escritor y para el planificador, no un contrato sobre el resultado.

Creamos una tercera tabla de laboratorio con un sort order declarado sobre ws_item_sk, usando la propiedad sorted_by que expone el conector Iceberg de Trino.

In [ ]:
run_trino("DROP TABLE IF EXISTS iceberg.tcdm.web_sales_lab_sorted")
run_trino('''
    CREATE TABLE iceberg.tcdm.web_sales_lab_sorted
    WITH (sorted_by = ARRAY['ws_item_sk']) AS
    SELECT ws_order_number, ws_item_sk, ws_sold_date, ws_net_paid
    FROM iceberg.tcdm.web_sales
''')

ddl_lab_sorted: str = run_trino("SHOW CREATE TABLE iceberg.tcdm.web_sales_lab_sorted").iloc[0, 0]
print(ddl_lab_sorted)
assert "sorted_by" in ddl_lab_sorted

SHOW CREATE TABLE confirma la propiedad sorted_by. Para contrastar con el otro extremo del espectro —ordenar sólo el resultado, sin tocar nada físico— ejecutamos un ORDER BY normal sobre la tabla de referencia: es exactamente el mismo tipo de operación que ya usábamos en S2 para presentar resultados, y no persiste ningún cambio en la tabla.

In [ ]:
run_trino('''
    SELECT ws_order_number, ws_net_paid FROM iceberg.tcdm.web_sales
    ORDER BY ws_net_paid DESC LIMIT 5
''')

Z-order¶

Z-order es una técnica que combina varias columnas de filtrado en un único criterio de organización física, útil cuando las consultas filtran a veces por una columna y a veces por otra, no siempre por la misma.

Esta sesión lo presenta sin ejecutarlo, porque con las versiones del laboratorio no hay una demostración reproducible:

  • el conector Iceberg de Trino no ofrece ninguna estrategia de Z-order: ALTER TABLE ... EXECUTE optimize sólo compacta por tamaño de fichero;
  • Iceberg sí lo define para Spark, en el procedimiento rewrite_data_files con strategy => 'sort' y un sort_order de tipo zorder(columna1, columna2, ...), pero con dos restricciones: no admite columnas DECIMAL —las monetarias de web_sales lo son— y en este entorno falla sobre una tabla cuyo sort order ya declaró Trino, como web_sales_lab_sorted.

La llamada tendría esta forma:

CALL iceberg.system.rewrite_data_files(
    table => 'tcdm.<tabla>',
    strategy => 'sort',
    sort_order => 'zorder(<columna1>, <columna2>)'
)

Como cualquier otra optimización de esta sesión, Z-order sólo merece la pena cuando hay un problema medido que resolver: filtros frecuentes por varias columnas sobre una tabla lo bastante grande como para que descartar ficheros por sus mínimos y máximos se note.

Recapitulación

Compactar y expirar tienen coste oculto

  • EXECUTE optimize reduce el número de ficheros y añade un snapshot replace
  • expire_snapshots se ha ejecutado sólo sobre la tabla de laboratorio
  • La retención mínima (7 días por defecto) protege a los lectores de versiones antiguas

Experimentos con una misma tabla¶

Reunimos en una sola tabla lo que hemos ido observando en las celdas anteriores. No inventamos aquí ningún número nuevo: cada fila hace referencia a una variable de Python o a un resultado ya calculado más arriba, siguiendo la misma idea que ya usó S5 en su propia tabla comparativa: mejor señalar exactamente de dónde sale cada dato que fabricar una cifra de tiempo de pared poco fiable en este entorno.

Variante Qué observamos Evidencia en este notebook
Parquet no particionado (raw) Sin poda posible; hay que recorrer todo web_sales explain_raw, sección "Particiones de datos"
Hive particionada por año/mes (web_sales_silver, S6) Directorios visibles, poda por columnas de partición declaradas explain_hive_2000
Iceberg con transformación temporal (web_sales) Particionado gestionado y oculto, poda sobre la columna de negocio explain_iceberg_2000
Iceberg con ficheros pequeños (web_sales_lab_smallfiles) Más ficheros y manifests con el mismo volumen de datos inventario_antes_compactar
Iceberg compactada Menos ficheros, mismo resultado de negocio inventario_despues_compactar
Iceberg con sort order (web_sales_lab_sorted) Propiedad sorted_by declarada, sin garantía de orden en lectura ddl_lab_sorted
Iceberg y Z-order No se ejecuta: Trino no lo ofrece y rewrite_data_files de Spark no admite columnas DECIMAL ni la tabla con sort order declarado por Trino sección "Z-order" más arriba

Y la comprobación de negocio que atraviesa toda la sesión: el número de líneas de venta no es lo mismo que el número de pedidos distintos.

In [ ]:
print("Líneas de venta vs pedidos distintos, iceberg.tcdm.web_sales:")
print(resultado_trino.to_dict("records")[0])
print()
print("Inventario de la tabla de laboratorio de ficheros pequeños:")
print("  antes de compactar:  ", inventario_antes_compactar)
print("  después de compactar:", inventario_despues_compactar)
print("  snapshots antes de expirar:  ", len(snapshots_lab_smallfiles))
print("  snapshots después de expirar:", len(snapshots_tras_expirar))

Ampliación opcional: leer una tabla Iceberg sin catálogo¶

Todo este notebook lee las tablas Iceberg por su nombre de catálogo (iceberg.tcdm.<tabla>), tal como haría cualquier despliegue normal de Spark con un Hive Metastore compatible. Existe, sin embargo, una vía alternativa de lectura que no depende en ningún momento del catálogo: resolver directamente el fichero metadata.json vigente de la tabla y cargarlo por su ruta física en HDFS. Es una vía útil en cualquier entorno donde Spark no tenga configurado un catálogo Iceberg —por ejemplo, un despliegue sólo con Trino, o una herramienta de inspección que no quiere depender de ningún catálogo—.

In [ ]:
def ruta_metadata_json_vigente(catalogo: str, esquema: str, tabla: str) -> str:
    '''Devuelve la ruta HDFS del metadata.json vigente de una tabla Iceberg.

    No depende del Hive Metastore para nada salvo para preguntarle a Trino
    (que sí puede leer la tabla) cuál es su ubicación en HDFS.
    '''
    ddl_texto: str = run_trino(f"SHOW CREATE TABLE {catalogo}.{esquema}.{tabla}").iloc[0, 0]
    ubicacion: str = ubicacion_de_tabla_iceberg(ddl_texto)

    listado: str = subprocess.run(
        ["hdfs", "dfs", "-ls", f"{ubicacion}/metadata"],
        capture_output=True,
        text=True,
        check=True,
    ).stdout

    ficheros_metadata: list[str] = sorted(
        linea.split()[-1]
        for linea in listado.splitlines()
        if linea.strip().endswith(".metadata.json")
    )
    assert ficheros_metadata, f"no se ha encontrado ningún metadata.json bajo {ubicacion}/metadata"
    return ficheros_metadata[-1]


def leer_tabla_iceberg_por_ruta(catalogo: str, esquema: str, tabla: str) -> DataFrame:
    '''Carga con Spark una tabla Iceberg por su metadata.json vigente, sin pasar por el catálogo.'''
    ruta_metadata: str = ruta_metadata_json_vigente(catalogo, esquema, tabla)
    return spark.read.format("iceberg").load(ruta_metadata)


web_sales_by_year_por_ruta: DataFrame = leer_tabla_iceberg_por_ruta(
    "iceberg", "tcdm", "web_sales_by_year"
)
web_sales_by_year_por_ruta.printSchema()
In [ ]:
filas_por_ruta: int = web_sales_by_year_por_ruta.count()
filas_por_catalogo: int = spark.table("iceberg.tcdm.web_sales_by_year").count()
print("Filas leyendo por ruta física:", filas_por_ruta)
print("Filas leyendo por catálogo:  ", filas_por_catalogo)
assert filas_por_ruta == filas_por_catalogo
print("Las dos vías de lectura devuelven el mismo resultado sobre iceberg.tcdm.web_sales_by_year.")
A continuación

Buenas prácticas del lakehouse

  • Fuente separada de las tablas derivadas: se crea con CTAS, no reescribiendo la fuente
  • Particiones según las consultas reales, nunca de alta cardinalidad «por si acaso»
  • Compactar y expirar sólo con un problema medido y una política de retención, y primero en laboratorio
  • Validar el resultado de negocio antes y después de cada cambio
  • Separar lo reversible de lo que borra, y las tablas web_sales_lab_ de las oficiales

Buenas prácticas del lakehouse¶

Las reglas que ya aprendiste en S5 sobre organizar un data lake siguen siendo válidas; Iceberg añade algunas propias de su mantenimiento como formato de tabla:

  • mantener los datos de origen separados de las tablas derivadas: aquí, hive.tcdm.web_sales_silver no se ha tocado en ningún momento;
  • documentar el propietario, el esquema y la ubicación de cada tabla Iceberg, igual que ya hacíamos con las tablas Hive de S6;
  • no sobrescribir la fuente para "convertirla" a Iceberg sin una estrategia explícita de copia o migración: por eso hemos creado iceberg.tcdm.web_sales con un CTAS, no reescribiendo web_sales_silver;
  • elegir las transformaciones de partición según las consultas reales, no "por si acaso": month(ws_sold_date) responde a cómo se filtran de hecho las ventas por periodo;
  • controlar el número y el tamaño de los ficheros de datos, y compactar cuando el inventario lo justifique, no de forma automática;
  • mantener una política explícita de snapshots y de retención antes de expirar nada, y aplicarla primero sobre tablas de laboratorio;
  • validar el resultado de negocio antes y después de cualquier cambio, como hemos hecho en cada sección de este notebook;
  • comprobar que todos los motores que van a usar la tabla comparten realmente el mismo catálogo, y no dar por hecho que "compatible" significa "sin incompatibilidades conocidas entre versiones concretas";
  • separar las operaciones reversibles (compactar, reescribir manifests) de las que borran metadatos o datos (expirar snapshots, eliminar ficheros huérfanos);
  • no compactar ni reordenar sin conocer primero el problema concreto que se quiere resolver;
  • no crear particiones de alta cardinalidad sólo para "tener partición";
  • no borrar snapshots sin una política de retención explícita, y nunca sobre las tablas que otras sesiones o procesos todavía necesitan;
  • no dar por validada una mejora a partir de una única ejecución en un entorno tan variable como este;
  • aislar con nombres claros las tablas de laboratorio de las tablas oficiales del curso, como hemos hecho con el prefijo web_sales_lab_.

Preguntas para interpretar la experiencia¶

  • ¿Qué diferencia observable hay entre el CREATE TABLE de una tabla Hive externa (S6) y el CTAS que ha creado iceberg.tcdm.web_sales en cuanto a qué ocurre en HDFS y qué metadatos se generan?
  • ¿Por qué la ubicación real de iceberg.tcdm.web_sales no coincide exactamente con /warehouse/tcdm.db/web_sales, y qué papel juega ese sufijo añadido?
  • ¿Qué diferencia hay entre cómo Hive y cómo Iceberg deciden qué directorio corresponde a cada valor de partición? ¿Quién tiene que declarar y sincronizar esa relación en cada caso?
  • ¿Qué demuestra que iceberg.tcdm.web_sales_lab_evolution conserva su snapshot inicial intacto después del INSERT y del ALTER TABLE ADD COLUMN? ¿Qué tendrías que hacer para perder esa capacidad de volver atrás?
  • ¿Por qué se ha decidido no experimentar con inserciones repetidas ni evolución de esquema directamente sobre iceberg.tcdm.web_sales, y crear en su lugar una copia de laboratorio?
  • ¿Qué diferencia hay entre ORDER BY en una consulta y el sorted_by que declaramos al crear iceberg.tcdm.web_sales_lab_sorted? ¿Qué garantiza cada uno y qué no garantiza ninguno de los dos?
  • ¿Por qué Trino impide expirar snapshots con una retención de cero días por defecto? ¿Qué riesgo concreto evita ese límite de 7 días?
  • Z-order sólo puede pedirse desde Spark con CALL iceberg.system.rewrite_data_files(...), nunca desde Trino. ¿Por qué? ¿Qué dos restricciones impiden demostrarlo en esta sesión, y qué tipo de consultas justificarían aplicarlo?
  • Antes de compactar iceberg.tcdm.web_sales_lab_smallfiles, ¿qué problema concreto estabas resolviendo? ¿Cómo lo mediste, y cómo comprobaste que seguía existiendo el mismo resultado de negocio después de compactar?

Evidencias para la siguiente sesión¶

Antes de la siguiente sesión tendrás una reunión individual breve con el profesor para revisar el trabajo de esta sesión. Esa reunión combina una demostración en vivo sobre tu propio ordenador y una memoria escrita breve.

Qué mostrar en el ordenador durante la reunión¶

Ten preparadas estas evidencias:

  1. El SHOW CREATE TABLE de iceberg.tcdm.web_sales y de iceberg.tcdm.web_sales_by_year, con su ubicación efectiva bajo /warehouse.
  2. El listado de HDFS de esa ubicación, distinguiendo los directorios data (con las particiones ocultas ws_sold_date_month=...) de metadata (los ficheros .metadata.json y los manifests .avro).
  3. Al menos dos snapshots de iceberg.tcdm.web_sales_lab_evolution, correspondientes a operaciones distintas (la creación y el INSERT posterior), y la confirmación de que el snapshot inicial sigue siendo consultable con FOR VERSION AS OF.
  4. La comprobación de que Trino y Spark devuelven el mismo resultado de negocio sobre iceberg.tcdm.web_sales (líneas de venta y pedidos distintos), leído por Spark directamente por su nombre de catálogo (iceberg.tcdm.web_sales), igual que las tablas Hive de S6.
  5. Los planes de EXPLAIN (Trino) y .explain("formatted") (Spark) antes y después de aplicar un filtro por partición, y la comparación de las tres variantes de particionado.
  6. La tabla de SHOW STATS FOR iceberg.tcdm.web_sales antes y después de ANALYZE, con tu propia valoración de si cambió o no el plan del JOIN con date_dim.
  7. El inventario de ficheros, manifests y snapshots de iceberg.tcdm.web_sales_lab_smallfiles antes y después de EXECUTE optimize, y antes y después de expire_snapshots.
  8. Que sabes explicar, con tus propias palabras, la diferencia entre ORDER BY, un sort order de tabla y la ordenación local frente a global.
  9. Que sabes explicar qué es Z-order, desde qué motor se puede pedir y por qué esta sesión no lo ejecuta (columnas DECIMAL no admitidas; tabla con un sort order ya declarado).
  10. Una conclusión propia, de una o dos frases, sobre cuándo una optimización de las vistas en esta sesión ayuda de verdad y cuándo no merece la pena aplicarla en TPC-DS SF1.
  11. Que sabes qué tablas de laboratorio se borran después de la reunión y cuáles deben conservarse para S8, tal como se describe a continuación.

Memoria escrita (una o dos páginas)¶

Trae también un documento breve —una o dos páginas, no hace falta más— que no se limite a pegar capturas de las salidas anteriores: debe explicar con tus propias palabras la diferencia entre un fichero Parquet, una tabla Hive externa y una tabla Iceberg; qué aporta el catálogo y qué aporta el formato de tabla; qué es un snapshot y qué ocurre con los anteriores tras un INSERT o un ALTER TABLE; en qué se diferencia una partición oculta de Iceberg de una partición visible de Hive; y tu conclusión razonada sobre qué optimización de las vistas ayuda de verdad en SF1 y cuál no.

Qué es importante de cara al examen final¶

El examen no pide recordar la sintaxis exacta de una sentencia. Debes poder explicar:

  • qué distingue un data lake de ficheros, un data lake catalogado y un lakehouse;
  • qué metadatos mantiene Iceberg (metadata.json, manifests y snapshots) y qué permite cada uno;
  • cómo funcionan la evolución de esquema y el viaje en el tiempo, y por qué no exigen reescribir los datos;
  • qué es una transformación de partición oculta y qué cambia respecto a las columnas de partición de Hive;
  • por qué la selección de columnas y los filtros se comprueban antes que cualquier otra optimización;
  • el problema de los ficheros pequeños, la compactación y la expiración de snapshots con una política de retención;
  • para qué sirven las estadísticas y por qué ANALYZE informa al planificador sin obligarlo a cambiar de plan;
  • la diferencia entre ORDER BY, un sort order de tabla y Z-order.

Detener y limpiar¶

Estas operaciones se muestran como texto, no como celdas de código, para que un «ejecutar todo» de este notebook no las dispare por accidente. Las tablas de laboratorio de esta sesión —web_sales_lab_evolution, web_sales_lab_smallfiles y web_sales_lab_sorted— son recursos desechables: consérvalas hasta la reunión de evidencias, que las necesita, y bórralas después. Nunca apliques el mismo razonamiento a iceberg.tcdm.web_sales ni a iceberg.tcdm.web_sales_by_year: son las dos tablas que reutilizará S8, y borrarlas o expirar sus snapshots obligaría a repetir esta sesión entera.

Situación Orden Consecuencia
Eliminar las tablas de laboratorio DROP TABLE iceberg.tcdm.web_sales_lab_evolution, DROP TABLE iceberg.tcdm.web_sales_lab_smallfiles, DROP TABLE iceberg.tcdm.web_sales_lab_sorted (desde el cliente Trino de este notebook, o docker exec trino-hdfs trino --user luser --execute "..." desde una terminal) Borra tanto la entrada del metastore como los ficheros Parquet y los metadatos Iceberg de cada tabla: son tablas administradas, su ciclo de vida de datos depende del catálogo.
Conservar iceberg.tcdm.web_sales e iceberg.tcdm.web_sales_by_year No ejecutar ningún DROP TABLE ni expire_snapshots sobre ellas S8 las necesita con exactamente el esquema y la clave compuesta que tienen ahora.
Parar temporalmente el warehouse docker compose -f entorno/compose-warehouse-hdfs.yml stop Conserva PostgreSQL, Hive Metastore y Trino (y su contenido, incluidas las tablas Iceberg); se reanuda con start. No afecta al clúster Hadoop.
Eliminar el warehouse por completo (destructivo) docker compose -f entorno/compose-warehouse-hdfs.yml down -v --remove-orphans Borra el volumen de PostgreSQL: se pierde todo el contenido del Hive Metastore, incluidas las definiciones de iceberg.tcdm.web_sales y iceberg.tcdm.web_sales_by_year. No borra /datalake ni /warehouse en HDFS: los ficheros Parquet y de metadatos seguirían ahí, ahora sin ninguna definición de catálogo que los describa. Tras volver a levantar el warehouse haría falta repetir los CREATE TABLE/CTAS de esta sesión antes de continuar con S8.

No confundas ese down -v del warehouse con el down/volume rm de S1 (clúster Hadoop) o con el down -v de la sesión 3 (RustFS): cada Compose tiene su propio volumen y afecta sólo a sus propios servicios.

Recapitulación

Sesión 7

  • Iceberg coordina metadatos, snapshots y evolución de esquema y particionado
  • Optimizar es medir antes y después, no aplicar una técnica a ciegas
  • Z-order queda documentado como ampliación, no como paso ejecutado
  • iceberg.tcdm.web_sales e iceberg.tcdm.web_sales_by_year quedan intactas para S8

Siguiente paso¶

S7 establece el modelo de lakehouse sobre la misma tabla que lo hace evolucionar y mide qué decisiones de particionado, tamaño de ficheros, estadísticas, ordenación y mantenimiento mejoran realmente las consultas —y cuáles no—. También ha dejado constancia de que "lakehouse" no es una propiedad binaria de un formato de tabla: las versiones de motor, cliente y catálogo se comprueban contra el clúster real, no se suponen compatibles porque «deberían» serlo.

S8 cierra el recorrido incorporando ingesta continua: un proceso simulado generará nuevos pedidos que un trabajo de Spark Structured Streaming irá fusionando desde raw hacia silver y gold sobre esta misma tabla iceberg.tcdm.web_sales, usando sus mecanismos de merge (copy-on-write y merge-on-read) en lugar de una recarga completa, y apoyándose exactamente en la clave compuesta (ws_order_number, ws_item_sk) que hemos mantenido intacta durante toda esta sesión.