Databricks Tips #4: Structured Streaming — watermarks, triggers y las trampas del micro-batch

Databricks Tips
Data Engineering
Streaming
Trigger modes, watermarks, foreachBatch, Auto Loader y las trampas que te hacen perder datos o plata.
Autor
Publicado

17 de marzo de 2026

Cuarta entrega de Databricks Tips. Structured Streaming en Databricks parece simple hasta que perdés datos en producción o te llega una factura de compute que no esperabas.

Trigger modes: la decisión más importante

El trigger define cuándo se ejecuta cada micro-batch. Elegir mal sale caro:

Listado 1: Trigger modes: processingTime vs availableNow según latencia
# MALO en producción: procesa continuamente, consume cluster 24/7
df.writeStream \
  .trigger(processingTime="0 seconds") \
  .start()

# MEJOR para la mayoría: cada 5 minutos
df.writeStream \
  .trigger(processingTime="5 minutes") \
  .start()

# IDEAL para pipelines que no necesitan baja latencia:
# Procesa todo lo pendiente y para
df.writeStream \
  .trigger(availableNow=True) \
  .start()

# Para pipelines event-driven (ej: cada vez que llega un archivo)
df.writeStream \
  .trigger(availableNow=True) \
  .start()
# Combinado con un job que se triggerea por file arrival

La trampa: processingTime="0 seconds" es el default y mucha gente ni lo toca. Un cluster corriendo 24/7 procesando cero datos la mayor parte del tiempo.

Regla práctica:

Latencia requerida Trigger
< 1 minuto processingTime="10 seconds"
1-15 minutos processingTime="5 minutes"
> 15 minutos availableNow=True + job schedule
Batch diario availableNow=True + cron diario

Watermarks: cómo no perder datos tardíos

Sin watermark, Spark mantiene todo el estado en memoria para joins y aggregations. Con tablas grandes, esto explota:

Listado 2: Watermark: limitar el estado en memoria para aggregations
from pyspark.sql.functions import window, col

# Sin watermark: el estado crece infinitamente
events \
  .groupBy(window("event_time", "1 hour")) \
  .count()

# Con watermark: Spark descarta estado de más de 2 horas
events \
  .withWatermark("event_time", "2 hours") \
  .groupBy(window("event_time", "1 hour")) \
  .count()

¿Qué pasa con datos que llegan después del watermark? Se descartan silenciosamente. No hay error, no hay log. Solo desaparecen.

Cómo elegir el watermark:

  1. Medí el delay máximo real de tus datos (percentil 99)
  2. Agregale un margen de seguridad (2x)
  3. Monitorealo: si ves datos faltantes, aumentá el watermark

Auto Loader: la forma correcta de ingestar archivos

spark.readStream.format("cloudFiles") es el Auto Loader. Es superior a spark.readStream.format("parquet") por varias razones:

Listado 3: Auto Loader vs file source: schema evolution y exactly-once
# MALO: file source básico
df = (spark.readStream
  .format("parquet")
  .schema(my_schema)
  .load("s3://bucket/landing/"))

# BIEN: Auto Loader
df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "parquet")
  .option("cloudFiles.schemaLocation", "/checkpoints/schema/")
  .option("cloudFiles.inferColumnTypes", "true")
  .option("cloudFiles.schemaEvolutionMode", "addNewColumns")
  .load("s3://bucket/landing/"))

# Escribir con merge schema
df.writeStream \
  .format("delta") \
  .option("checkpointLocation", "/checkpoints/events/") \
  .option("mergeSchema", "true") \
  .trigger(availableNow=True) \
  .toTable("catalog.bronze.events")

Ventajas del Auto Loader:

  • Notificación vs listing: usa SNS/SQS (AWS) o Event Grid (Azure) para detectar archivos nuevos. En directorios con millones de archivos, listing tarda minutos; notificación es instantáneo.
  • Schema evolution: detecta columnas nuevas automáticamente.
  • Rescue column: columnas que no matchean el schema van a _rescued_data en vez de fallar.
  • Exactly-once: el checkpoint garantiza que nunca procesás un archivo dos veces.
Listado 4: Rescue column: capturar datos que no matchean el schema
# Habilitar rescue column para datos sucios
df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaLocation", "/checkpoints/schema/")
  .option("rescuedDataColumn", "_rescued_data")
  .load("s3://bucket/raw/"))

foreachBatch: el escape hatch

Cuando necesitás lógica que no se puede expresar en streaming puro (upserts, llamadas a APIs, multi-table writes):

Listado 5: foreachBatch con MERGE: upsert idempotente en streaming
def upsert_to_delta(batch_df, batch_id):
    from delta.tables import DeltaTable

    target = DeltaTable.forName(spark, "catalog.silver.customers")

    (target.alias("t")
     .merge(batch_df.alias("s"), "t.customer_id = s.customer_id")
     .whenMatchedUpdateAll()
     .whenNotMatchedInsertAll()
     .execute())

# Stream con upsert
(spark.readStream
  .table("catalog.bronze.customers")
  .writeStream
  .foreachBatch(upsert_to_delta)
  .option("checkpointLocation", "/checkpoints/customers_upsert/")
  .trigger(availableNow=True)
  .start())

Ojo con foreachBatch: si tu función falla a mitad de camino, el batch se reintenta completo. Asegurate de que la operación sea idempotente (MERGE lo es, INSERT no).

Monitoreo: qué mirar en producción

Listado 6: StreamingQueryListener: monitorear batches y detectar lag
# Query listener para métricas
from pyspark.sql.streaming import StreamingQueryListener

class StreamMonitor(StreamingQueryListener):
    def onQueryProgress(self, event):
        progress = event.progress
        print(f"""
        Batch: {progress.batchId}
        Input rows: {progress.numInputRows}
        Processing time: {progress.batchDuration}ms
        Watermark: {progress.eventTime.get('watermark', 'N/A')}
        """)

spark.streams.addListener(StreamMonitor())

Métricas clave:

  • numInputRows = 0 por mucho tiempo → el trigger está corriendo sin datos (plata tirada)
  • batchDuration crece batch a batch → el estado está creciendo, necesitás watermark
  • inputRowsPerSecond vs processedRowsPerSecond → si input > processed, estás acumulando lag

Próxima semana: MLflow en Unity Catalog — registro de modelos, linaje de experimentos y serving.