Databricks Tips #4: Structured Streaming — watermarks, triggers y las trampas del micro-batch
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:
# 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 arrivalLa 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:
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:
- Medí el delay máximo real de tus datos (percentil 99)
- Agregale un margen de seguridad (2x)
- 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:
# 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_dataen vez de fallar. - Exactly-once: el checkpoint garantiza que nunca procesás un archivo dos veces.
# 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):
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
# 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 = 0por mucho tiempo → el trigger está corriendo sin datos (plata tirada)batchDurationcrece batch a batch → el estado está creciendo, necesitás watermarkinputRowsPerSecondvsprocessedRowsPerSecond→ si input > processed, estás acumulando lag
Próxima semana: MLflow en Unity Catalog — registro de modelos, linaje de experimentos y serving.