Sesión 8 — Ingesta continua con Spark Structured Streaming¶

Esta sesión cierra el recorrido iniciado en S2. El alumnado ya sabe leer y escribir Parquet con Spark, organizar un data lake en zonas raw/silver/ gold, registrar tablas en un Metastore y gestionar una tabla Iceberg con snapshots, evolución y optimización medida. Todo ese trabajo, hasta ahora, ha sido por lotes: se materializa TPC-DS SF1 una vez y se procesa.

La pregunta central es:

¿Qué cambia cuando los datos no llegan de una vez, sino en tandas sucesivas que hay que ir incorporando sin volver a procesar todo lo anterior?

Esta sesión no introduce Kafka ni ningún otro sistema de mensajería: sería una pieza de infraestructura nueva y compleja para un beneficio docente pequeño en este punto del curso. En su lugar, un proceso simple simula la llegada periódica de pedidos nuevos escribiendo ficheros en la zona raw, y Spark Structured Streaming los recoge desde ahí con su propia fuente de ficheros: exactamente el mismo mecanismo de vigilancia de directorio que ya se ha usado en HDFS durante todo el curso, aplicado ahora de forma continua en vez de una sola vez.

Paquetes de Python de esta sesión. Esta sesión necesita pyspark, trino, pandas, fsspec, pyarrow, requests. Respecto a S7 añade fsspec, pyarrow, requests. 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 s8/ 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 8: ingesta continua con Spark Structured Streaming

Del lote único a las tandas

  • La tabla Iceberg de S7 se alimenta de forma continua, sin recargarla
  • Micro-lotes, MERGE INTO y coherencia con el recálculo por lotes
  • Cierre del itinerario S1-S8

Recorrido del curso¶

S2: Parquet como ficheros
  → S3: los mismos Parquet en un almacenamiento de objetos
  → S4: Spark procesa por lotes
  → S5: zonas raw/silver/gold y particionado físico
  → S6: catálogo Hive y Trino
  → S7: tabla Iceberg, evolución y optimización medida
  → S8: la misma tabla Iceberg, alimentada de forma continua

S8 no sustituye el recorrido por lotes: lo completa. La tabla Iceberg de S7 sigue siendo la tabla de referencia; esta sesión añade una segunda vía de escritura, incremental, que debe producir resultados de negocio coherentes con la vía por lotes.

Qué significa "streaming" aquí, y qué no¶

Conviene fijar el vocabulario antes de escribir código:

Concepto Significado en esta sesión Lo que no es
Generador simulado Proceso simple que escribe nuevos ficheros Parquet en una ruta raw de aterrizaje Un sistema de colas o mensajería
Fuente de streaming spark.readStream vigilando esa ruta y detectando ficheros nuevos Una conexión a un broker
Micro-lote (micro-batch) Conjunto de filas nuevas que Structured Streaming procesa en una ejecución del disparador Un evento individual procesado uno a uno
Disparador (trigger) Cada cuánto tiempo, o con qué condición, se procesa un micro-lote Un cron del sistema operativo
Sink con foreachBatch Función que recibe el DataFrame de cada micro-lote y decide qué hacer con él Una escritura directa fila a fila
Merge (Iceberg) Operación MERGE INTO que inserta, actualiza o ignora filas según coincidan con la tabla destino Sobrescribir toda la tabla

Spark Structured Streaming procesa datos no acotados como una sucesión de lotes pequeños, no evento a evento. Eso es exactamente lo que se necesita para ilustrar cómo pasar de raw a silver y de silver a gold sin volver a leer todo lo anterior en cada ejecución.

Objetivos de la sesión¶

Al terminar esta sesión deberías poder:

  • explicar la diferencia entre procesamiento por lotes y por micro-lotes;
  • crear un readStream sobre una ruta HDFS con la fuente de ficheros;
  • distinguir el esquema declarado explícitamente (obligatorio en streaming) de la inferencia de esquema usada en las sesiones por lotes;
  • usar foreachBatch para aplicar una operación arbitraria —en este caso, dos MERGE INTO— a cada micro-lote, sin volver a leer todo el histórico cada vez;
  • explicar qué configuración necesita el catálogo Iceberg de Spark (tipo hive, el mismo Hive Metastore que usa Trino) para que foreachBatch pueda escribir directamente en iceberg.tcdm.*;
  • escribir un MERGE INTO SQL con Spark sobre una tabla Iceberg que inserte pedidos nuevos y actualice pedidos existentes en la misma operación;
  • distinguir los modos de escritura de Iceberg copy-on-write y merge-on-read, y observar su efecto físico (ficheros reescritos frente a ficheros de borrado);
  • encadenar una actualización de silver con una actualización incremental de un agregado gold, sin recalcular todo el histórico;
  • consultar el progreso de una consulta de streaming (lastProgress, número de filas por micro-lote, tiempo de procesamiento);
  • detener una consulta de streaming de forma ordenada y explicar qué garantiza (y qué no) el checkpoint;
  • comprobar que el resultado de negocio tras varias tandas incrementales coincide con el que daría un recálculo completo por lotes.
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: los ficheros de aterrizaje, la consulta con al menos tres micro-lotes y sus métricas, una fila insertada y una corregida por el mismo MERGE, y la diferencia física entre cow y mor
  • Se entrega también una memoria breve (una o dos páginas)

Antes de empezar¶

Dónde se ejecuta este notebook. Igual que en las sesiones anteriores, Jupyter se ejecuta directamente dentro del contenedor namenode, como luser, con acceso de red directo a HDFS, YARN y el Hive Metastore: todos comparten la misma 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 ya está en marcha y el warehouse (PostgreSQL, Hive Metastore y Trino) de S6 también, porque el catálogo Iceberg de esta sesión usa el mismo Hive Metastore y las celdas de sólo lectura usan Trino;
  • S2 generó los Parquet de TPC-DS SF1 en /datalake/raw/tpcds, incluidas web_sales, customer, customer_address, item y date_dim;
  • S7 ya creó y pobló las dos tablas Iceberg que esta sesión alimenta de forma continua: iceberg.tcdm.web_sales (la silver a nivel de línea) y iceberg.tcdm.web_sales_by_year (el agregado gold). Esta sesión no las crea: las comprueba más abajo y falla pronto y con un mensaje claro si no existen o si su esquema no es el esperado.

Si el clúster o el warehouse no están levantados 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.

In [ ]:
!hdfs dfs -ls -d \
    /datalake/raw/tpcds/web_sales \
    /datalake/raw/tpcds/customer \
    /datalake/raw/tpcds/customer_address \
    /datalake/raw/tpcds/item \
    /datalake/raw/tpcds/date_dim \
    /warehouse

Si alguna de esas rutas no aparece, revisa S2 antes de continuar: esta sesión no regenera TPC-DS. La comprobación de que S7 ya creó las tablas Iceberg llega un poco más abajo, en «Verificar que S7 ya creó las tablas Iceberg».

Instalar las dependencias de este notebook¶

El kernel de este notebook instala PySpark sobre sí mismo, igual que en las sesiones anteriores, más el cliente trino de PyPI —el mismo de S6 y S7—, que usaremos en varias celdas de sólo lectura como confirmación independiente de lo que Spark escribe directamente en iceberg.tcdm.* (se explica en la siguiente sección). El generador simulado que escribiremos y ejecutaremos más abajo (generate_web_sales_batch.py) es un programa aparte que habla con HDFS sólo por WebHDFS, con fsspec y PyArrow — el mismo mecanismo que ya usaron los lectores de S2 —, así que sus dependencias también viven en requirements.txt. La celda siguiente lo recrea 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 8 (Spark local con un
# catálogo Iceberg de tipo hive, el mismo Hive Metastore que usa Trino). La
# serie de PySpark aceptada por el curso es >4,<4.2, igual que en S4, S5, S6
# y S7.
pyspark>4,<4.2
# Cliente DB-API 2.0 de Trino: se usa en las celdas de sólo lectura como
# comprobación independiente de lo que Spark escribe directamente en
# iceberg.tcdm.* dentro de foreachBatch (ver "Cómo configuramos el catálogo
# Iceberg de Spark" en el notebook). 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
# generate_web_sales_batch.py habla con WebHDFS por fsspec y lee/escribe
# Parquet con PyArrow, igual que los lectores de la sesión 2; requests es
# la dependencia HTTP que usa fsspec para WebHDFS.
fsspec>=2026.2.0
pyarrow
requests
In [ ]:
%pip install -q -r requirements.txt
In [ ]:
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__}")
A continuación

Configurar el catálogo Iceberg de Spark

  • Catálogo tipo hive, contra el mismo Thrift que usa Trino
  • Spark escribe los dos MERGE dentro de foreachBatch
  • Trino sólo confirma leyendo, como verificación independiente

Cómo configuramos el catálogo Iceberg de Spark¶

Esta sesión declara un catálogo Iceberg de Spark respaldado por Hive Metastore, exactamente el mismo patrón que ya usan S6 (catálogo Hive nativo) y S7 (catálogo Iceberg): spark.sql.catalog.iceberg de tipo hive, apuntando al mismo thrift://hive-metastore:9083 que usa Trino. A partir de esa configuración, Spark lee y escribe iceberg.tcdm.* directamente con spark.sql(...): el MERGE INTO que cada micro-lote de streaming ejecuta más abajo, en foreachBatch, lo hace Spark, sin ningún motor ni tabla intermedios.

Spark habla con este catálogo con su cliente por defecto porque la versión fijada del curso es Hive Metastore 3.1.3; la justificación completa está en entorno/README.md, sección "Versiones fijadas y compatibilidad". El resto de esta sesión usa spark.sql(...) para las dos operaciones que de verdad importan aquí —el MERGE INTO hacia silver y hacia gold dentro de foreachBatch— y sigue usando run_trino(...) en las celdas que sólo leen, como comprobación independiente de que ambos motores ven exactamente lo mismo.

Recapitulación

Quién escribe y quién comprueba

  • Spark ejecuta los dos MERGE directamente sobre iceberg.tcdm.*
  • Trino lee lo mismo como verificación independiente, no como motor de escritura

Construir la SparkSession¶

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

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

spark: SparkSession = (
    SparkSession.builder.appName("TCDM_S8_Streaming")
    .master("local[2]")
    .config("spark.driver.memory", "1g")
    .config("spark.sql.shuffle.partitions", "8")
    .config("spark.jars.packages", ICEBERG_RUNTIME_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)
    .getOrCreate()
)
sc: SparkContext = spark.sparkContext

print(f"Spark version: {spark.version} | Master: {sc.master}")
print(f"Paquete Iceberg: {ICEBERG_RUNTIME_PACKAGE}")
print(f"Catálogo Iceberg: iceberg (Hive Metastore {HIVE_METASTORE_URIS})")

Abrir la conexión con Trino¶

Reutilizamos exactamente el patrón de S6 y S7: 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. Esta sesión no necesita a Trino para escribir —Spark ya lo hace directamente contra iceberg.tcdm.*, con el catálogo que acabamos de configurar—, pero sí lo usamos más abajo en varias celdas de sólo lectura: cuando run_trino(...) devuelve una fila que Spark acaba de escribir, es la confirmación de que un motor completamente independiente ve exactamente el mismo dato.

In [ ]:
import subprocess

import pandas as pd
import trino
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)


print(f"Cliente trino: {trino.__version__}")
run_trino("SELECT 1 AS ping")

Verificar que S7 ya creó las tablas Iceberg¶

Antes de escribir una sola línea de streaming, comprobamos con Spark —el motor que va a ejecutar los MERGE INTO reales en foreachBatch— que las dos tablas del contrato de esta sesión existen en el catálogo iceberg y tienen las columnas esperadas. Si esta comprobación falla, el mensaje debe bastar para saber que hay que ejecutar S7 antes de continuar, en vez de fallar más adelante con un error críptico de MERGE INTO.

In [ ]:
def described_columns(table_name: str) -> dict[str, str]:
    '''Convertir el esquema de Spark de una tabla en un diccionario columna -> tipo.'''
    assert spark.catalog.tableExists(
        table_name
    ), f"No existe {table_name}. ¿Se ha ejecutado S7 antes que esta sesión?"
    return {field.name: field.dataType.simpleString() for field in spark.table(table_name).schema}


web_sales_columns: dict[str, str] = described_columns("iceberg.tcdm.web_sales")
for name, data_type in web_sales_columns.items():
    print(f"{name}: {data_type}")
In [ ]:
expected_web_sales_columns: set[str] = {
    "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",
}
missing_web_sales_columns: set[str] = expected_web_sales_columns - set(web_sales_columns)
assert not missing_web_sales_columns, (
    "Faltan columnas en iceberg.tcdm.web_sales: "
    f"{missing_web_sales_columns}. ¿Se ha ejecutado S7 antes que esta sesión?"
)
print("iceberg.tcdm.web_sales tiene todas las columnas del contrato de S8.")
In [ ]:
web_sales_by_year_columns: dict[str, str] = described_columns("iceberg.tcdm.web_sales_by_year")
print(web_sales_by_year_columns)

expected_gold_columns: set[str] = {"d_year", "pedidos", "lineas_de_venta", "beneficio_neto"}
missing_gold_columns: set[str] = expected_gold_columns - set(web_sales_by_year_columns)
assert not missing_gold_columns, (
    "Faltan columnas en iceberg.tcdm.web_sales_by_year: "
    f"{missing_gold_columns}. ¿Se ha ejecutado S7 antes que esta sesión?"
)
print("iceberg.tcdm.web_sales_by_year tiene todas las columnas del contrato de S8.")

Por último, inspeccionamos el CREATE TABLE real de iceberg.tcdm.web_sales —esta vez con spark.sql(...), el mismo catálogo que va a escribir en ella— para ver, sin adivinarlo, cómo particionó S7 la tabla. Debe aparecer una transformación oculta de Iceberg sobre ws_sold_date (por ejemplo, month(ws_sold_date)) en vez de columnas de partición visibles como las de la tabla Hive de S6.

In [ ]:
create_table_web_sales: str = spark.sql("SHOW CREATE TABLE iceberg.tcdm.web_sales").collect()[0][0]
print(create_table_web_sales)
A continuación

Arquitectura y generador simulado

  • El generador escribe ficheros Parquet nuevos en /datalake/raw/streaming/web_sales_incremental
  • readStream los recoge y foreachBatch hace dos MERGE INTO sobre las tablas Iceberg de S7
  • Generador y consulta sólo se coordinan a través de esa ruta de aterrizaje
  • Cada lote mezcla correcciones (clave ya existente) y altas (clave nueva) sobre (ws_order_number, ws_item_sk)

Arquitectura de la sesión¶

Generador simulado (Python, se ejecuta varias veces)
        │  escribe ficheros Parquet nuevos
        ▼
/datalake/raw/streaming/web_sales_incremental/batch-<id>.parquet
        │  spark.readStream.format("parquet")
        ▼
Micro-lote de Structured Streaming (Spark)
        │  foreachBatch: vista temporal + MERGE INTO (Spark, catálogo iceberg)
        ▼
iceberg.tcdm.web_sales           (ya creada en S7)
        │  mismo foreachBatch: segundo MERGE INTO (Spark)
        ▼
iceberg.tcdm.web_sales_by_year   (ya creada en S7)

Spark orquesta el streaming y escribe directamente en iceberg.tcdm.*: no hay ningún motor ni tabla intermedios. El generador y la consulta de streaming siguen siendo procesos independientes que se coordinan únicamente a través de la ruta de aterrizaje en HDFS, igual que antes.

El generador simulado¶

generate_web_sales_batch.py construye, cada vez que se ejecuta, un lote de líneas de venta web sintéticas y lo escribe como un fichero Parquet nuevo bajo la ruta de aterrizaje, con un nombre (marca de tiempo + UUID corto) que no se repite entre ejecuciones. No usa Spark: habla con HDFS sólo por WebHDFS con fsspec y PyArrow, igual que los lectores de S2.

web_sales es una tabla a nivel de línea, no de pedido: su clave real es el par (ws_order_number, ws_item_sk), no ws_order_number en solitario, porque un mismo pedido puede tener varias líneas con productos distintos. El generador calcula, antes de fabricar ninguna fila, los rangos reales de c_customer_sk, ca_address_sk e i_item_sk (con PyArrow, sobre los Parquet reales de customer, customer_address e item) y lee de date_dim las fechas del periodo real de ventas de web_sales, que empieza en 1998 — así las claves generadas siguen siendo coherentes con las tablas de dimensión, sin inventar rangos, y las líneas nuevas caen en los mismos años que las ventas que ya existen.

Cada lote mezcla dos tipos de líneas:

  • correcciones, que reutilizan un par (ws_order_number, ws_item_sk) muestreado del web_sales real de S2 (ya presente en la tabla Iceberg desde que S7 la materializó) con importes nuevos: es la parte "update" del MERGE;
  • altas, con un par (ws_order_number, ws_item_sk) nuevo que no puede colisionar con ningún pedido real (los números de pedido de alta parten de una base muy por encima del rango real de TPC-DS SF1): es la parte "insert".

La fracción de correcciones se controla con --update-fraction (0.3 por defecto). Ejecutamos el generador una primera vez para inspeccionar lo que produce, antes de construir la consulta de streaming.

El generador es un fichero de la sesión (s8/generate_web_sales_batch.py), pero este notebook se ejecuta dentro de namenode, que no ve tu copia de la distribución. La celda siguiente lo descarga del repositorio público dsevilla/tcdm-public (rama 26-27) al directorio de trabajo del kernel, igual que la sesión 3 descarga sus programas. Si el fichero ya está ahí —por ejemplo, porque lo has copiado tú—, se conserva.

In [ ]:
from pathlib import Path
from urllib.request import Request, urlopen

GENERATOR: Path = Path("generate_web_sales_batch.py")
generator_url: str = (
    "https://raw.githubusercontent.com/dsevilla/tcdm-public/26-27/s8/generate_web_sales_batch.py"
)

if GENERATOR.exists():
    print(f"{GENERATOR} ya está en {Path.cwd()}; se conserva.")
else:
    request: Request = Request(generator_url, headers={"User-Agent": "Mozilla/5.0"})
    with urlopen(request, timeout=30) as response:
        GENERATOR.write_bytes(response.read())
    print(f"Descargado {GENERATOR} en {Path.cwd()}")
In [ ]:
import json

first_run: subprocess.CompletedProcess[str] = subprocess.run(
    ["python", "generate_web_sales_batch.py", "--rows", "150"],
    capture_output=True,
    text=True,
    check=True,
)
print(first_run.stdout)

first_report: dict[str, object] = json.loads(
    first_run.stdout.strip().splitlines()[-1].split("\t", 1)[1]
)
batch_reports: list[dict[str, object]] = [first_report]
print(first_report)
In [ ]:
LANDING_PATH: str = "/datalake/raw/streaming/web_sales_incremental"
print(f"Ruta de aterrizaje: {LANDING_PATH}")
In [ ]:
!hdfs dfs -ls {LANDING_PATH}

El fichero que aparece en el listado es el que acabamos de escribir: su nombre incluye la marca de tiempo y un fragmento de UUID, así que volver a ejecutar el generador no lo sobrescribe. first_report trae, entre otras cosas, la ruta exacta del fichero (batch_path) y las claves (ws_order_number, ws_item_sk) que ha usado para las correcciones y para una muestra de las altas — las usaremos más abajo para enseñar una fila insertada y una fila corregida por el mismo MERGE, en vez de adivinarlas.

A continuación

Una fuente que vigila un directorio

  • readStream con un esquema declarado explícitamente
  • El generador escribe ficheros y no conoce a Spark
  • El checkpoint recuerda qué ficheros ya se procesaron

La consulta de streaming¶

A diferencia de la lectura por lotes de S4/S5 (spark.read.parquet(...)), aquí declaramos el esquema del readStream a mano: Spark podría inferirlo leyendo una muestra, pero esa inferencia no es fiable en un flujo continuo —depende de qué ficheros haya presentes en cada arranque y puede cambiar entre ejecuciones—, así que se declara explícitamente en vez de encarecer y fragilizar el arranque. Declaramos el mismo esquema que ya construye el generador, columna a columna.

In [ ]:
from pyspark.sql.types import (
    DateType,
    DecimalType,
    IntegerType,
    LongType,
    StructField,
    StructType,
)

web_sales_incremental_schema: StructType = StructType(
    [
        StructField("ws_order_number", LongType(), nullable=False),
        StructField("ws_item_sk", LongType(), nullable=False),
        StructField("ws_bill_customer_sk", LongType(), nullable=True),
        StructField("ws_bill_addr_sk", LongType(), nullable=True),
        StructField("ws_sold_date_sk", LongType(), nullable=True),
        StructField("ws_sold_date", DateType(), nullable=True),
        StructField("ws_quantity", IntegerType(), nullable=True),
        StructField("ws_list_price", DecimalType(7, 2), nullable=True),
        StructField("ws_ext_discount_amt", DecimalType(7, 2), nullable=True),
        StructField("ws_ext_sales_price", DecimalType(7, 2), nullable=True),
        StructField("ws_net_paid", DecimalType(7, 2), nullable=True),
        StructField("ws_net_profit", DecimalType(7, 2), nullable=True),
        StructField("sold_year", IntegerType(), nullable=True),
        StructField("sold_month", IntegerType(), nullable=True),
    ]
)
In [ ]:
from pyspark.sql import DataFrame

raw_stream: DataFrame = (
    spark.readStream.format("parquet").schema(web_sales_incremental_schema).load(LANDING_PATH)
)

print("readStream configurado sobre:", LANDING_PATH)

raw_stream no ha leído ninguna fila todavía: como en las transformaciones por lotes de S4, es perezoso. Lo que sí ha hecho es fijar el esquema y la ruta que vigilará la fuente de ficheros una vez arranque la consulta.

A continuación

Qué hace foreachBatch con cada micro-lote

  • Registra el micro-lote como vista temporal micro_batch_src
  • Anota qué año tenía ya cada clave que va a tocar: una corrección puede cambiar una línea de año
  • MERGE a silver por (ws_order_number, ws_item_sk): altas y correcciones a la vez
  • MERGE a gold sólo de los años tocados: los nuevos y los que las filas tenían antes
  • Un año que se queda sin filas se borra de gold

foreachBatch: vista temporal + dos MERGE INTO¶

La función hace, en este orden:

  1. Registra el micro-lote como vista temporal (micro_batch_src) con createOrReplaceTempView: es lo único que hace falta para que spark.sql(...) pueda referenciarlo como origen de un MERGE INTO.
  2. Antes de fusionar nada, pregunta —con spark.sql(...), contra el mismo catálogo iceberg— qué año (sold_year) tenía ya cada clave que el micro-lote va a tocar, cruzando la vista temporal con iceberg.tcdm.web_sales tal y como está antes del MERGE. Es necesario porque el generador simulado reasigna una fecha nueva y aleatoria, dentro del periodo real de ventas, a cada fila que fabrica, tanto si es una corrección como si es un alta (lo comprobarás si lees generate_web_sales_batch.py): una corrección puede mover una línea de venta real de, por ejemplo, 2001 a 1999. Si touched_years sólo contuviera el año nuevo de esa fila, el MERGE de gold recalcularía 1999 (correcto) pero nunca volvería a tocar 2001 —el año que la fila acaba de abandonar—, y el agregado de 2001 se quedaría con una fila fantasma para siempre.
  3. Fusiona micro_batch_src contra iceberg.tcdm.web_sales, con la clave compuesta real de la tabla (ws_order_number y ws_item_sk): un mismo micro-lote puede traer altas y correcciones a la vez, sin que el programa tenga que distinguirlas de antemano. Como la vista temporal tiene exactamente las mismas columnas que la tabla, el MERGE puede usar la abreviatura de la extensión de Iceberg en Spark SQL, UPDATE SET */INSERT *, en vez de listar las catorce columnas a mano.
  4. Fusiona iceberg.tcdm.web_sales_by_year, recalculando sólo la unión de los años que el micro-lote toca ahora y los años que las filas corregidas tenían antes —no toda la tabla—, a partir de la propia iceberg.tcdm.web_sales ya actualizada por el paso anterior. Esta vez la fuente del MERGE no es una vista sobre el micro-lote: es una subconsulta con una tabla de valores literal (VALUES ...) y una agregación directa sobre iceberg.tcdm.web_sales.

Un año tocado puede quedarse sin ninguna fila (si la única línea que tenía se corrigió hacia otro año): la subconsulta de origen incluye explícitamente todos los años tocados, con COALESCE(..., 0) para los que ya no tienen filas, y el MERGE borra la fila de gold correspondiente (WHEN MATCHED AND source.lineas_de_venta = 0 THEN DELETE) en lugar de dejarla obsoleta. Sin este DELETE, el MERGE no tocaría esa fila porque no la evaluaría en absoluto: sólo compara filas que sí aparecen en la subconsulta de origen, y una consulta sin COALESCE que agrupe únicamente lo que queda en web_sales no produciría ninguna fila para un año vacío.

In [ ]:
def merge_into_silver_and_gold(micro_batch_df: DataFrame, batch_id: int) -> None:
    '''Aplicar un micro-lote de streaming a la silver y a la gold Iceberg.

    Spark ejecuta aquí los dos MERGE INTO directamente contra el catálogo
    iceberg (Hive Metastore), el mismo que ya verificamos más arriba: no hay
    ningún motor ni tabla intermedios.
    '''

    if micro_batch_df.isEmpty():
        print(f"[foreachBatch] Micro-lote {batch_id}: 0 filas, nada que fusionar.")
        return

    # foreachBatch entrega micro_batch_df con su propia SparkSession "clonada":
    # aunque comparte el mismo catálogo que la sesión global `spark`, es un
    # objeto Python distinto, y una vista temporal es visible sólo dentro de
    # la sesión que la registró. Por eso todo este cuerpo usa
    # `batch_spark`, obtenida del propio DataFrame, en vez de la variable
    # global `spark` capturada por el closure.
    batch_spark: SparkSession = micro_batch_df.sparkSession

    row_count: int = micro_batch_df.count()
    print(f"[foreachBatch] Micro-lote {batch_id}: {row_count} filas recibidas.")

    micro_batch_df.createOrReplaceTempView("micro_batch_src")

    # Años que las claves de este micro-lote tenían ANTES del MERGE: hace
    # falta capturarlos ahora, antes de fusionar, porque una corrección
    # puede cambiar el año de una fila (ver la celda anterior).
    previous_years_df: DataFrame = batch_spark.sql('''
        SELECT DISTINCT target.sold_year AS sold_year
        FROM iceberg.tcdm.web_sales AS target
        JOIN micro_batch_src AS source
          ON target.ws_order_number = source.ws_order_number
             AND target.ws_item_sk = source.ws_item_sk
        ''')
    previous_years: set[int] = {row["sold_year"] for row in previous_years_df.collect()}

    batch_spark.sql('''
        MERGE INTO iceberg.tcdm.web_sales AS target
        USING micro_batch_src AS source
        ON target.ws_order_number = source.ws_order_number
           AND target.ws_item_sk = source.ws_item_sk
        WHEN MATCHED THEN UPDATE SET *
        WHEN NOT MATCHED THEN INSERT *
        ''')

    new_years: set[int] = {
        row["sold_year"] for row in micro_batch_df.select("sold_year").distinct().collect()
    }
    touched_years: list[int] = sorted(previous_years | new_years)
    years_list_sql: str = ", ".join(str(year) for year in touched_years)
    years_values_sql: str = ", ".join(f"({year})" for year in touched_years)
    print(f"[foreachBatch] Micro-lote {batch_id}: años tocados = {touched_years}")

    batch_spark.sql(f'''
        MERGE INTO iceberg.tcdm.web_sales_by_year AS target
        USING (
            SELECT year_list.sold_year AS d_year,
                   COALESCE(agg.pedidos, 0) AS pedidos,
                   COALESCE(agg.lineas_de_venta, 0) AS lineas_de_venta,
                   COALESCE(agg.beneficio_neto, 0) AS beneficio_neto
            FROM (VALUES {years_values_sql}) AS year_list (sold_year)
            LEFT JOIN (
                SELECT sold_year,
                       count(DISTINCT ws_order_number) AS pedidos,
                       count(*) AS lineas_de_venta,
                       sum(ws_net_profit) AS beneficio_neto
                FROM iceberg.tcdm.web_sales
                WHERE sold_year IN ({years_list_sql})
                GROUP BY sold_year
            ) AS agg ON year_list.sold_year = agg.sold_year
        ) AS source
        ON target.d_year = source.d_year
        WHEN MATCHED AND source.lineas_de_venta = 0 THEN DELETE
        WHEN MATCHED THEN UPDATE SET
            pedidos = source.pedidos,
            lineas_de_venta = source.lineas_de_venta,
            beneficio_neto = source.beneficio_neto
        WHEN NOT MATCHED AND source.lineas_de_venta > 0 THEN INSERT (d_year, pedidos, lineas_de_venta, beneficio_neto)
            VALUES (source.d_year, source.pedidos, source.lineas_de_venta, source.beneficio_neto)
        ''')

El segundo MERGE es la diferencia práctica entre mantener un agregado de forma incremental y recalcularlo por completo en cada micro-lote: el coste de la actualización depende de cuántos años ha tocado el micro-lote, no del tamaño total de silver. Con este generador, que reparte las fechas por todo el periodo de ventas, cada micro-lote toca casi todos los años; en una ingesta real, donde los pedidos nuevos son recientes, serían uno o dos.

Arrancar la consulta¶

El checkpoint guarda qué ficheros de la fuente ya se han procesado, para que reiniciar la consulta más adelante no vuelva a aplicar el mismo micro-lote dos veces.

Recapitulación

Un micro-lote, dos fusiones

  • MERGE a silver con la clave compuesta (ws_order_number, ws_item_sk)
  • El segundo MERGE sólo toca los años del micro-lote actual, y también los años que dejaron de tener filas
  • Actualizar gold no significa recalcularlo entero
In [ ]:
from pyspark.sql.streaming import StreamingQuery

CHECKPOINT_PATH: str = "/user/luser/s8/checkpoints/web_sales_incremental"

query: StreamingQuery = (
    raw_stream.writeStream.foreachBatch(merge_into_silver_and_gold)
    .trigger(processingTime="5 seconds")
    .option("checkpointLocation", CHECKPOINT_PATH)
    .start()
)

print("Consulta de streaming arrancada. ID:", query.id)

La consulta arranca con el fichero del primer lote ya esperando en la ruta de aterrizaje. processAllAvailable() bloquea hasta que se procesa todo lo que había disponible en el momento de llamarlo — no espera al siguiente disparador ni se queda esperando datos que nunca llegan —, así que es la forma correcta de sincronizar este notebook con la consulta sin dejar ninguna celda bloqueada indefinidamente.

In [ ]:
query.processAllAvailable()
print(query.lastProgress)

Primeros micro-lotes¶

Simulamos la periodicidad ejecutando el generador varias veces seguidas desde este mismo notebook, con una pausa corta entre ejecuciones. En un entorno real este mismo programa se lanzaría desde cron, un systemd timer o un job programado; aquí basta con dejar clara esa equivalencia sin construirla.

El bucle tiene un número de iteraciones fijo (LOOP_ITERATIONS = 3, más el lote inicial de la sección anterior, cuatro en total) y cada vuelta llama a processAllAvailable() en vez de a awaitTermination() sin límite: esta celda termina siempre en un tiempo acotado, no se queda sirviendo la consulta para siempre.

In [ ]:
import time

LOOP_ITERATIONS: int = 3

for iteration in range(1, LOOP_ITERATIONS + 1):
    run: subprocess.CompletedProcess[str] = subprocess.run(
        ["python", "generate_web_sales_batch.py", "--rows", "150", "--update-fraction", "0.4"],
        capture_output=True,
        text=True,
        check=True,
    )
    report: dict[str, object] = json.loads(run.stdout.strip().splitlines()[-1].split("\t", 1)[1])
    batch_reports.append(report)
    print(f"--- Generador, vuelta {iteration}: {report['batch_path']} ({report['rows']} filas) ---")

    query.processAllAvailable()
    print(f"[lastProgress tras la vuelta {iteration}]")
    print(query.lastProgress)

    time.sleep(2)

print(f"Lotes procesados en total: {len(batch_reports)}")

lastProgress es un diccionario con, entre otras claves, numInputRows (filas del micro-lote), durationMs (desglose de tiempos por fase: addBatch, getBatch, queryPlanning...) y sink. Es la respuesta concreta a "¿cuánto tarda en incorporarse un lote nuevo, desde que el generador lo escribe hasta que aparece en el agregado gold?": la suma de triggerExecution de cada micro-lote.

Mientras la consulta está activa, la interfaz web de Spark (http://localhost:4040, presentada en S4) añade la pestaña Structured Streaming, con una fila por consulta y, dentro de ella, las gráficas de filas de entrada y de duración de cada micro-lote. Los MERGE que lanza foreachBatch aparecen además como consultas en SQL / DataFrame.

In [ ]:
progress_summary: list[dict[str, object]] = [
    {
        "batchId": progress["batchId"],
        "numInputRows": progress["numInputRows"],
        "durationMs_total": progress["durationMs"].get("triggerExecution"),
    }
    for progress in query.recentProgress
]
for entry in progress_summary:
    print(entry)
Pregunta guía
  • ¿Cómo distingue el mismo MERGE una alta de una corrección si el programa no se lo dice?
  • ¿Qué papel juega la clave compuesta (ws_order_number, ws_item_sk)?

MERGE INTO hacia silver, explicado¶

Con los cuatro lotes ya fusionados, buscamos en iceberg.tcdm.web_sales una fila que sabemos que es una corrección (una clave que ya existía en el web_sales real de S2, con un importe nuevo) y una fila que sabemos que es un alta (una clave que no existía antes de esta sesión). Las dos claves salen de batch_reports, no se adivinan.

In [ ]:
update_reports_with_keys: list[dict[str, object]] = [
    report for report in batch_reports if report["updated_keys"]
]
assert update_reports_with_keys, "Ningún lote generó correcciones; revisa --update-fraction"

example_update_report: dict[str, object] = update_reports_with_keys[-1]
example_update_order_number, example_update_item_sk = example_update_report["updated_keys"][0]

print(
    f"Clave corregida en el último lote con correcciones: ({example_update_order_number}, {example_update_item_sk})"
)

run_trino(f'''
    SELECT *
    FROM iceberg.tcdm.web_sales
    WHERE ws_order_number = {example_update_order_number}
      AND ws_item_sk = {example_update_item_sk}
    ''')
In [ ]:
example_insert_order_number, example_insert_item_sk = first_report["inserted_keys_sample"][0]
print(f"Clave nueva del primer lote: ({example_insert_order_number}, {example_insert_item_sk})")

run_trino(f'''
    SELECT *
    FROM iceberg.tcdm.web_sales
    WHERE ws_order_number = {example_insert_order_number}
      AND ws_item_sk = {example_insert_item_sk}
    ''')

La primera consulta devuelve una fila cuya clave ya existía en el web_sales real de S2 (el mismo (ws_order_number, ws_item_sk)), pero con los importes que generó el micro-lote correspondiente: es la rama WHEN MATCHED THEN UPDATE SET * del MERGE. La segunda devuelve una fila cuya clave no puede existir en TPC-DS SF1 (los números de pedido de alta parten de una base muy por encima del rango real): es la rama WHEN NOT MATCHED THEN INSERT *. El mismo MERGE INTO, ejecutado una sola vez por micro-lote, ha resuelto las dos situaciones sin que foreachBatch tuviera que distinguirlas de antemano.

De silver a gold, de forma incremental¶

iceberg.tcdm.web_sales_by_year se ha ido actualizando en el mismo foreachBatch que la silver, un MERGE por micro-lote que sólo toca los años presentes en ese micro-lote. La comprobamos ahora completa.

In [ ]:
run_trino("SELECT * FROM iceberg.tcdm.web_sales_by_year ORDER BY d_year")

Ningún paso de este notebook ha recalculado iceberg.tcdm.web_sales_by_year desde toda la silver: cada micro-lote sólo ha vuelto a agregar los años que él mismo tocaba (touched_years, impreso más arriba por cada vuelta del bucle). Con TPC-DS SF1 —unos pocos años de web_sales— la diferencia de coste entre las dos estrategias no se nota en el tiempo de pared de esta sesión, pero sí se notaría con un histórico de muchos años: recalcular todo el agregado en cada micro-lote crecería con el tamaño de silver, mientras que este MERGE filtrado crece sólo con el número de años tocados por el propio micro-lote.

Recapitulación

Agregar no es recalcular todo

  • Gold se actualiza sólo por los años que toca el micro-lote
  • En SF1 la diferencia no se nota en el reloj, pero crece con el número de años de un histórico largo, no con el tamaño de silver
A continuación

Dos formas de escribir lo mismo

  • Copy-on-write reescribe los ficheros afectados
  • Merge-on-read acumula ficheros de borrado
  • Trino sólo implementa merge-on-read; Spark permite elegir y comparar

Copy-on-write frente a merge-on-read¶

Iceberg permite configurar, por tabla, cómo se materializan las filas afectadas por un MERGE:

  • copy-on-write (write.merge.mode = 'copy-on-write'): Iceberg reescribe por completo los ficheros de datos afectados.
  • merge-on-read (write.merge.mode = 'merge-on-read'): Iceberg escribe únicamente los cambios — ficheros de borrado y, cuando aplica, ficheros de datos nuevos — y delega en el lector la tarea de aplicarlos al leer.

Qué motor permite elegir. El conector Iceberg de Trino sólo implementa merge-on-read para UPDATE/DELETE/MERGE: siempre escribe ficheros de borrado, nunca reescribe ficheros de datos completos. Por eso no deja declarar write.merge.mode, write.update.mode ni write.delete.mode; si se intentan pasar a través de extra_properties, responde:

TrinoUserError: Illegal keys in extra_properties: [write.merge.mode]

Es una característica del conector (trinodb/trino#17272), no de este clúster. El escritor de Iceberg de Spark sí respeta write.merge.mode, y el catálogo iceberg de Spark configurado más arriba habla con el mismo Hive Metastore que Trino: creamos las dos tablas de laboratorio con spark.sql("CREATE TABLE iceberg.tcdm... TBLPROPERTIES(...)"), en el mismo catálogo que usa el resto de la sesión. Como estas tablas sólo sirven para esta comparación, se llaman iceberg.tcdm.web_sales_lab_cow y iceberg.tcdm.web_sales_lab_mor, y se borran en la sección "Detener y limpiar" al final del notebook.

Cada tabla parte de un único año de iceberg.tcdm.web_sales (el que tenga más líneas), leído directamente con Spark —el mismo catálogo ya sabe leer esta tabla, así que no hace falta pasar por Trino ni por pandas— y cargado con INSERT. No de la tabla completa: el mecanismo que se observa no depende del volumen de datos.

Ambas tablas se crean con 'format-version' = '2': los ficheros de borrado a nivel de fila que necesita merge-on-read sólo existen en el formato de tabla Iceberg v2.

In [ ]:
from pyspark.sql import Row

lab_year_counts: list[Row] = (
    spark.table("iceberg.tcdm.web_sales")
    .groupBy("sold_year")
    .count()
    .withColumnRenamed("count", "lineas")
    .orderBy("lineas", ascending=False)
    .collect()
)
LAB_YEAR: int = int(lab_year_counts[0]["sold_year"])
print(
    f"Año elegido para las tablas de laboratorio: {LAB_YEAR} ({lab_year_counts[0]['lineas']} líneas)"
)
In [ ]:
LAB_TABLE_DDL_COLUMNS: str = '''
    ws_order_number BIGINT,
    ws_item_sk BIGINT,
    ws_bill_customer_sk BIGINT,
    ws_bill_addr_sk BIGINT,
    ws_sold_date_sk BIGINT,
    ws_sold_date DATE,
    ws_quantity INT,
    ws_list_price DECIMAL(7, 2),
    ws_ext_discount_amt DECIMAL(7, 2),
    ws_ext_sales_price DECIMAL(7, 2),
    ws_net_paid DECIMAL(7, 2),
    ws_net_profit DECIMAL(7, 2),
    sold_year INT,
    sold_month INT
'''

spark.sql("DROP TABLE IF EXISTS iceberg.tcdm.web_sales_lab_cow")
spark.sql(f'''
    CREATE TABLE iceberg.tcdm.web_sales_lab_cow ({LAB_TABLE_DDL_COLUMNS})
    USING iceberg
    PARTITIONED BY (months(ws_sold_date))
    TBLPROPERTIES (
        'format-version' = '2',
        'write.merge.mode' = 'copy-on-write',
        'write.update.mode' = 'copy-on-write',
        'write.delete.mode' = 'copy-on-write'
    )
''')

spark.sql("DROP TABLE IF EXISTS iceberg.tcdm.web_sales_lab_mor")
spark.sql(f'''
    CREATE TABLE iceberg.tcdm.web_sales_lab_mor ({LAB_TABLE_DDL_COLUMNS})
    USING iceberg
    PARTITIONED BY (months(ws_sold_date))
    TBLPROPERTIES (
        'format-version' = '2',
        'write.merge.mode' = 'merge-on-read',
        'write.update.mode' = 'merge-on-read',
        'write.delete.mode' = 'merge-on-read'
    )
''')
print(
    "Tablas de laboratorio creadas: web_sales_lab_cow (copy-on-write), web_sales_lab_mor (merge-on-read)."
)

Cargamos cada tabla de laboratorio con las filas de LAB_YEAR directamente desde iceberg.tcdm.web_sales, con el mismo catálogo Spark: no hace falta pasar por Trino ni por pandas, así que no hay que lidiar con la conversión de tipos ni con valores nulos representados como NaN.

In [ ]:
from pyspark.sql.functions import col as spark_col

lab_source_df: DataFrame = (
    spark.table("iceberg.tcdm.web_sales")
    .filter(spark_col("sold_year") == LAB_YEAR)
    .select([field.name for field in web_sales_incremental_schema.fields])
)

lab_source_df.write.mode("append").saveAsTable("iceberg.tcdm.web_sales_lab_cow")
lab_source_df.write.mode("append").saveAsTable("iceberg.tcdm.web_sales_lab_mor")

lab_source_rows: int = lab_source_df.count()
cow_rows: int = spark.table("iceberg.tcdm.web_sales_lab_cow").count()
mor_rows: int = spark.table("iceberg.tcdm.web_sales_lab_mor").count()
print(f"Filas cargadas: cow={cow_rows}, mor={mor_rows}")
assert cow_rows == mor_rows == lab_source_rows

Para aplicar la misma secuencia de cambios a las dos tablas construimos, con Spark, un pequeño lote sintético que reutiliza tres claves que ya están en las tablas de laboratorio (corrección, rama "update") y añade dos claves nuevas (alta, rama "insert") — el mismo tipo de mezcla que ya generaba generate_web_sales_batch.py, pero fabricado aquí directamente para garantizar que las claves de corrección existen en este año concreto.

Estas filas de laboratorio no reutilizan la clave de fecha ni los importes reales: ws_sold_date_sk y los importes quedan a None, porque sólo sirven para comparar el efecto físico del mismo MERGE entre las tablas copy-on-write y merge-on-read — no para representar ventas reales.

In [ ]:
import random
from datetime import date

from pyspark.sql import Row

lab_sample_keys: list[Row] = (
    spark.table("iceberg.tcdm.web_sales_lab_cow")
    .select("ws_order_number", "ws_item_sk")
    .limit(3)
    .collect()
)

lab_rng: random.Random = random.Random(2027)
lab_rows: list[dict[str, object]] = []

for key_row in lab_sample_keys:
    lab_rows.append(
        {
            "ws_order_number": key_row["ws_order_number"],
            "ws_item_sk": key_row["ws_item_sk"],
            "ws_bill_customer_sk": lab_rng.randint(1, 100_000),
            "ws_bill_addr_sk": lab_rng.randint(1, 50_000),
            "ws_sold_date_sk": None,
            "ws_sold_date": date(LAB_YEAR, 6, 15),
            "ws_quantity": lab_rng.randint(1, 20),
            "ws_list_price": None,
            "ws_ext_discount_amt": None,
            "ws_ext_sales_price": None,
            "ws_net_paid": None,
            "ws_net_profit": None,
            "sold_year": LAB_YEAR,
            "sold_month": 6,
        }
    )

for offset in range(2):
    lab_rows.append(
        {
            "ws_order_number": 950_000_000 + offset,
            "ws_item_sk": lab_rng.randint(1, 18_000),
            "ws_bill_customer_sk": lab_rng.randint(1, 100_000),
            "ws_bill_addr_sk": lab_rng.randint(1, 50_000),
            "ws_sold_date_sk": None,
            "ws_sold_date": date(LAB_YEAR, 6, 15),
            "ws_quantity": lab_rng.randint(1, 20),
            "ws_list_price": None,
            "ws_ext_discount_amt": None,
            "ws_ext_sales_price": None,
            "ws_net_paid": None,
            "ws_net_profit": None,
            "sold_year": LAB_YEAR,
            "sold_month": 6,
        }
    )

lab_merge_batch: DataFrame = spark.createDataFrame(lab_rows, schema=web_sales_incremental_schema)
lab_merge_batch.createOrReplaceTempView("lab_merge_batch_view")
lab_merge_batch.show(truncate=False)

Antes de fusionar, listamos el contenido físico de cada tabla de laboratorio: son la línea base con la que compararemos después del MERGE.

In [ ]:
from pyspark.sql import Row


def table_location(table_name: str) -> str:
    '''Ubicación física real de una tabla del catálogo iceberg (Hive Metastore).

    A diferencia de un catálogo hadoop, donde la ruta se puede calcular sin
    preguntarle a nadie, aquí el Hive Metastore puede añadir un sufijo único
    a la ubicación de cada tabla (como ya viste en S7 con
    "SHOW CREATE TABLE"): hay que preguntárselo al catálogo con
    "DESCRIBE TABLE EXTENDED", la versión Spark del mismo tipo de consulta.
    '''
    extended_rows: list[Row] = spark.sql(f"DESCRIBE TABLE EXTENDED {table_name}").collect()
    location_row: Row = next(row for row in extended_rows if row["col_name"] == "Location")
    return location_row["data_type"]


cow_location: str = table_location("iceberg.tcdm.web_sales_lab_cow")
mor_location: str = table_location("iceberg.tcdm.web_sales_lab_mor")
print(f"cow: {cow_location}")
print(f"mor: {mor_location}")
In [ ]:
!hdfs dfs -ls -R {cow_location}/data
In [ ]:
!hdfs dfs -ls -R {mor_location}/data

Ahora aplicamos el mismo MERGE INTO —misma clave, mismo lote— a las dos tablas.

In [ ]:
spark.sql('''
    MERGE INTO iceberg.tcdm.web_sales_lab_cow AS target
    USING lab_merge_batch_view AS source
    ON target.ws_order_number = source.ws_order_number
       AND target.ws_item_sk = source.ws_item_sk
    WHEN MATCHED THEN UPDATE SET *
    WHEN NOT MATCHED THEN INSERT *
''')

spark.sql('''
    MERGE INTO iceberg.tcdm.web_sales_lab_mor AS target
    USING lab_merge_batch_view AS source
    ON target.ws_order_number = source.ws_order_number
       AND target.ws_item_sk = source.ws_item_sk
    WHEN MATCHED THEN UPDATE SET *
    WHEN NOT MATCHED THEN INSERT *
''')
print("MERGE aplicado a las dos tablas de laboratorio.")
In [ ]:
!hdfs dfs -ls -R {cow_location}/data
In [ ]:
!hdfs dfs -ls -R {mor_location}/data

Compara los dos listados posteriores al MERGE con los que tomamos antes: en web_sales_lab_cow deben aparecer ficheros de datos nuevos que sustituyen a los que contenían las tres filas corregidas (el resto de ficheros, los que no tenían ninguna fila afectada, se conservan igual). En web_sales_lab_mor los ficheros de datos originales deben seguir intactos y aparecer, además, uno o varios ficheros nuevos de un tipo distinto — ficheros de borrado (delete files), cuyo nombre termina en -deletes.parquet—, más los ficheros de datos con las filas nuevas. Esa es la diferencia física que promete la teoría: copy-on-write paga el coste en la escritura, reescribiendo; merge-on-read lo aplaza a la lectura, acumulando ficheros de borrado que en algún momento habría que compactar — exactamente el mismo problema de pequeños ficheros y compactación que ya se vio en S7.

Recapitulación

El coste se paga al escribir o al leer

  • En cow aparecen ficheros nuevos que sustituyen a los afectados
  • En mor los ficheros originales siguen y se añaden ficheros de borrado
  • Una tabla mor con muchos MERGE es la que justifica compactar (enlace con S7)

Observar el progreso y detener la consulta¶

Antes de detener la consulta original (la que alimenta las tablas compartidas del contrato), repasamos las tres formas de consultar su estado.

In [ ]:
print("query.status:", query.status)
print("query.lastProgress:", query.lastProgress)
print(f"Micro-lotes acumulados en recentProgress: {len(query.recentProgress)}")

query.status distingue si la consulta está esperando datos nuevos o procesando; lastProgress describe el último micro-lote; recentProgress guarda el historial de esta ejecución (con un límite interno de entradas). Detenemos la consulta de forma ordenada:

In [ ]:
query.stop()
print(f"Consulta detenida. isActive = {query.isActive}")

Reiniciar desde el mismo checkpoint¶

Detener la consulta no borra el checkpoint: si se reinicia la misma consulta con la misma ruta de checkpoint y no hay ficheros nuevos que procesar, Structured Streaming no debería reprocesar nada. Lo comprobamos arrancando una segunda consulta idéntica, sin generar ningún lote nuevo antes.

In [ ]:
query_restarted: StreamingQuery = (
    raw_stream.writeStream.foreachBatch(merge_into_silver_and_gold)
    .trigger(processingTime="5 seconds")
    .option("checkpointLocation", CHECKPOINT_PATH)
    .start()
)

query_restarted.processAllAvailable()
restarted_progress: dict[str, object] | None = query_restarted.lastProgress
print(restarted_progress)

query_restarted.stop()
print(f"Segunda consulta detenida. isActive = {query_restarted.isActive}")

Si restarted_progress es None o su numInputRows es 0, es la señal de que el checkpoint ha hecho su trabajo: la consulta reiniciada sabe qué ficheros ya se procesaron en la ejecución anterior y no ha vuelto a fusionarlos. Borrar el checkpoint (hdfs dfs -rm -r /user/luser/s8/checkpoints/web_sales_incremental) obligaría a reprocesar desde el principio — útil para repetir el ejercicio desde cero, pero debe hacerse de forma explícita, nunca como efecto colateral de otra limpieza; se deja documentado, no como celda, en la sección "Detener y limpiar" de más abajo.

A continuación

La prueba de que el streaming no cambia el resultado

  • El checkpoint evita reprocesar los mismos ficheros
  • El agregado incremental debe coincidir, año a año, con el recálculo completo por lotes

Comprobar la coherencia con el recorrido por lotes¶

La comprobación más importante de esta sesión: el agregado gold mantenido de forma incremental por streaming debe coincidir, año a año, con el que se obtendría recalculándolo por lotes desde silver, con la misma consulta que ya usaron S5 y S7 (pedidos distintos, líneas de venta y beneficio neto).

In [ ]:
batch_recomputed_gold: pd.DataFrame = run_trino('''
    SELECT sold_year AS d_year,
           count(DISTINCT ws_order_number) AS pedidos,
           count(*) AS lineas_de_venta,
           sum(ws_net_profit) AS beneficio_neto
    FROM iceberg.tcdm.web_sales
    GROUP BY sold_year
    ORDER BY d_year
    ''')

incremental_gold: pd.DataFrame = run_trino(
    "SELECT * FROM iceberg.tcdm.web_sales_by_year ORDER BY d_year"
)

comparison: pd.DataFrame = batch_recomputed_gold.merge(
    incremental_gold, on="d_year", suffixes=("_batch", "_incremental")
)
comparison["beneficio_neto_batch"] = comparison["beneficio_neto_batch"].astype(float)
comparison["beneficio_neto_incremental"] = comparison["beneficio_neto_incremental"].astype(float)
comparison["diferencia_beneficio"] = (
    comparison["beneficio_neto_batch"] - comparison["beneficio_neto_incremental"]
).abs()

assert (
    len(comparison) == len(batch_recomputed_gold) == len(incremental_gold)
), "Faltan años en una de las dos versiones"
assert (comparison["pedidos_batch"] == comparison["pedidos_incremental"]).all()
assert (comparison["lineas_de_venta_batch"] == comparison["lineas_de_venta_incremental"]).all()
assert (comparison["diferencia_beneficio"] < 0.01).all()

print("El agregado incremental coincide, año a año, con el recálculo completo por lotes.")
comparison

S7 demostró que (ws_order_number, ws_item_sk) es la clave única de iceberg.tcdm.web_sales. Esta sesión reescribe esa misma tabla con MERGE en cada micro-lote, así que cerramos el contrato comprobándolo de nuevo, y de paso miramos los snapshots que ha ido dejando el streaming — uno por micro-lote, igual que en S7.

In [ ]:
key_uniqueness: pd.DataFrame = run_trino('''
    WITH totals AS (
        SELECT count(*) AS filas FROM iceberg.tcdm.web_sales
    ),
    distinct_keys AS (
        SELECT count(*) AS claves FROM (
            SELECT DISTINCT ws_order_number, ws_item_sk FROM iceberg.tcdm.web_sales
        )
    )
    SELECT filas, claves FROM totals, distinct_keys
''')
assert int(key_uniqueness.loc[0, "filas"]) == int(key_uniqueness.loc[0, "claves"]), (
    "La clave compuesta (ws_order_number, ws_item_sk) debe seguir siendo única "
    "en iceberg.tcdm.web_sales tras los MERGE de esta sesión"
)
key_uniqueness
In [ ]:
snapshots_web_sales_tras_streaming: pd.DataFrame = run_trino('''
    SELECT snapshot_id, committed_at, operation
    FROM iceberg.tcdm."web_sales$snapshots"
    ORDER BY committed_at
''')
snapshots_web_sales_tras_streaming

Esta comprobación es la evidencia de que el streaming no ha introducido un motor de cálculo distinto: sigue siendo SQL sobre la misma tabla —ejecutado por Trino, no por Spark, como comprobación independiente de lo que Spark escribió—, con más frecuencia y sobre menos datos por ejecución. Si alguno de los assert anteriores fallara, el error estaría en la lógica del foreachBatch (por ejemplo, en el cálculo de touched_years), no en los datos de origen.

Recapitulación

Mismo SQL, más a menudo y sobre menos datos

  • Si el recálculo por lotes y el incremental difieren, el fallo está en touched_years, no en los datos de origen
  • El resultado de negocio es la evidencia, no el número de micro-lotes procesados

Qué no se ha visto todavía¶

Esta sesión no ha introducido:

  • Kafka ni ningún otro sistema de mensajería: el generador simulado escribe ficheros directamente en HDFS, y readStream los recoge con la misma fuente de ficheros que ya se conocía de las sesiones por lotes;
  • procesamiento evento a evento ni garantías de latencia de segundos: Spark Structured Streaming, incluso en su modo de disparador más frecuente, sigue procesando micro-lotes, no eventos individuales;
  • watermarking ni agregaciones con ventanas de tiempo sobre datos desordenados: quedan como ampliación si el tiempo del curso lo permite;
  • un panel de control que visualice los resultados: eso se aborda en una fase posterior explícita, construida sobre esta misma tabla gold, con Metabase como añadido al Compose del clúster y un mecanismo para declarar paneles "como código" mediante su API REST.

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 dos consultas de streaming de este notebook ya se han detenido con query.stop() más arriba; lo que queda pendiente es limpieza de recursos de laboratorio. Las tablas de laboratorio y los ficheros de aterrizaje forman parte de las evidencias: consérvalos hasta la reunión y límpialos después. Las órdenes spark.sql(...) van en una celda de este notebook, igual que las de hdfs dfs ... precedidas de !; docker compose, en una terminal de tu equipo desde la raíz de la distribución.

Situación Orden Consecuencia
Eliminar las tablas de laboratorio copy-on-write/merge-on-read spark.sql("DROP TABLE iceberg.tcdm.web_sales_lab_cow") y lo mismo para web_sales_lab_mor Borra las dos tablas del catálogo iceberg (metadatos en el Hive Metastore y ficheros bajo /warehouse). No afecta a iceberg.tcdm.web_sales ni a iceberg.tcdm.web_sales_by_year: son tablas distintas.
Vaciar la ruta de aterrizaje del streaming hdfs dfs -rm -r /datalake/raw/streaming/web_sales_incremental Borra todos los lotes Parquet generados por generate_web_sales_batch.py en esta sesión. No afecta a /datalake/raw/tpcds.
Forzar que la consulta reprocese todo desde el principio hdfs dfs -rm -r /user/luser/s8/checkpoints/web_sales_incremental Borra el checkpoint: la próxima vez que se arranque la consulta sobre la misma ruta de aterrizaje, volverá a fusionar todos los ficheros que sigan ahí. Útil para repetir el ejercicio desde cero, nunca como efecto colateral de otra limpieza.
Volver al estado en que S7 dejó las tablas Con la consulta detenida: spark.sql("DROP TABLE iceberg.tcdm.web_sales") y spark.sql("DROP TABLE iceberg.tcdm.web_sales_by_year"); borrar la ruta de aterrizaje y el checkpoint (las dos filas anteriores); y volver a ejecutar en S7 las celdas de «Crear la tabla Iceberg» Deshace todo lo que esta sesión ha fusionado: S7 recrea las dos tablas desde las tablas Hive de S6, que no han cambiado. Si no se borran el aterrizaje y el checkpoint, la siguiente ejecución de esta sesión volvería a aplicar los lotes antiguos.
Parar temporalmente el warehouse docker compose -f entorno/compose-warehouse-hdfs.yml stop Conserva PostgreSQL, Hive Metastore y Trino (y su contenido); se reanuda con start. No afecta al clúster Hadoop ni a iceberg.tcdm.web_sales.

No confundas ninguna de estas operaciones con docker compose ... down -v sobre el warehouse: ese comando sí es destructivo sobre volúmenes completos, como se documenta en el propio Compose del warehouse; el equivalente en S1 (clúster Hadoop) es el borrado de volúmenes con docker volume rm/make clean-hdfs, no un down -v.

Preguntas para interpretar la experiencia¶

  • ¿Por qué readStream exige un esquema explícito cuando spark.read.parquet de las sesiones anteriores podía inferirlo?
  • ¿Qué recibe exactamente la función que se pasa a foreachBatch? ¿En qué se parece a un DataFrame de las sesiones por lotes y en qué no?
  • ¿Por qué el MERGE INTO hacia iceberg.tcdm.web_sales puede usar la abreviatura UPDATE SET */INSERT *, y qué tendría que dejar de ser cierto sobre la vista temporal del micro-lote para que esa abreviatura fallara?
  • ¿Por qué el MERGE INTO de esta sesión compara ws_order_number y ws_item_sk, y qué pasaría si sólo comparara ws_order_number?
  • ¿Qué garantiza el checkpoint de una consulta de streaming, y qué no garantiza? ¿Qué evidencia de este notebook lo demuestra?
  • ¿Por qué el segundo MERGE de foreachBatch filtra por los años tocados por el micro-lote en vez de recalcular todo iceberg.tcdm.web_sales_by_year cada vez?
  • ¿Por qué Trino sólo implementa merge-on-read para escribir en Iceberg, y qué consecuencia tiene eso para una tabla que recibe muchos MERGE?
  • ¿Qué diferencia física observaste en HDFS entre web_sales_lab_cow y web_sales_lab_mor después del mismo MERGE? ¿Cuál de las dos tablas necesitaría antes una compactación, y por qué?
  • ¿Por qué la comprobación de coherencia con el recálculo por lotes es la evidencia más importante de esta sesión, más que el número de micro-lotes procesados?
  • ¿Qué parte de esta sesión cambiaría si, en vez de un generador simulado, los datos llegaran realmente por Kafka? ¿Qué parte del foreachBatch seguiría siendo exactamente igual?

Evidencias de la sesión¶

Esta es la última sesión del itinerario, así que no hay una sesión siguiente: la reunión individual breve con el profesor para revisar este trabajo se fija aparte. Como en las anteriores, 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. Varias ejecuciones del generador, con los ficheros nuevos visibles en /datalake/raw/streaming/web_sales_incremental (hdfs dfs -ls).
  2. La consulta de streaming en marcha, con al menos tres micro-lotes procesados y sus métricas (lastProgress/recentProgress).
  3. Una fila insertada y una fila corregida en iceberg.tcdm.web_sales por el mismo MERGE, distinguibles con sus claves.
  4. iceberg.tcdm.web_sales_by_year actualizado de forma incremental, coincidente con el recálculo por lotes de control (la tabla comparison de este notebook).
  5. El listado HDFS de web_sales_lab_cow y web_sales_lab_mor antes y después del mismo MERGE, mostrando la diferencia entre ficheros de datos reescritos y ficheros de borrado.
  6. Una consulta de streaming detenida de forma ordenada (query.stop()) y reiniciada desde su checkpoint sin reprocesar micro-lotes ya aplicados.

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 qué cambia entre procesar por lotes y por micro-lotes; qué hace foreachBatch con cada micro-lote y por qué necesita dos MERGE; por qué la clave del MERGE es el par (ws_order_number, ws_item_sk); qué garantiza el checkpoint y qué no; qué diferencia física hay entre copy-on-write y merge-on-read; y por qué el agregado mantenido de forma incremental coincide con el recálculo completo.

Qué es importante de cara al examen final¶

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

  • qué es un micro-lote, qué es un disparador y por qué una fuente de ficheros en streaming exige declarar el esquema;
  • qué recibe foreachBatch y por qué permite reutilizar en streaming las operaciones de las sesiones por lotes;
  • cómo resuelve un único MERGE INTO las altas, las correcciones y los borrados, y qué papel tiene la clave de comparación;
  • cómo se mantiene un agregado de forma incremental, tocando sólo lo que cambia, y por qué hay que tener en cuenta el valor anterior de las filas corregidas;
  • qué recuerda un checkpoint y qué ocurre al reiniciar una consulta con y sin él;
  • la diferencia entre copy-on-write y merge-on-read, y su relación con la compactación vista en S7;
  • por qué el resultado de negocio, y no el número de micro-lotes, es la evidencia de que la ingesta continua es correcta.
Recapitulación

Sesión 8

  • Ficheros + readStream + foreachBatch + MERGE: la misma tabla Iceberg de S7, alimentada de forma continua
  • El itinerario S1-S8 cierra aquí
  • El panel de control con Metabase queda como fase posterior, sobre la misma tabla gold

Siguiente paso¶

El itinerario docente completo (S1-S8) enseña a levantar el clúster, conocer los datos como ficheros, moverlos a un almacenamiento de objetos, procesarlos con Spark, organizarlos en zonas, catalogarlos, gestionarlos como una tabla lakehouse y, finalmente, alimentarlos de forma continua. La fase siguiente, ya fuera de esta numeración, añade un panel de control (Metabase) como servicio adicional del Compose, con una conexión a Trino y un mecanismo para declarar nuevos paneles mediante ficheros de configuración y la API REST de Metabase, consultando la misma tabla gold que esta sesión mantiene al día.