Saltar al contenido
Lakehouse LabLakehouse LabPreparación Databricks Data Engineer
Módulo 13 · Professional

Streaming

Contenido abierto

Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.

Professional

Structured Streaming, triggers y checkpoints

Construye consultas incrementales con estado durable y semántica de recuperación clara.

Lectura pública
Al terminar podrás
  • Explicar microbatches y progreso
  • Configurar triggers y checkpoints
  • Diseñar sinks idempotentes
Ver fuentes y revisión

Metadatos editoriales

Última revisión
21 jul 2026
Nivel
Professional
Ruta relacionada
streaming
Dominios blueprint
Data Ingestion and Acquisition · Structured Streaming
Estado
Revisión editorial interna
Fuentes principales
Structured Streaming checkpoints · Databricks · Monitor Structured Streaming queries · Databricks
Reportar un error
01
Modelo mental

Modelo incremental de Structured Streaming

Structured Streaming ejecuta una consulta incremental como una secuencia de microbatches y conserva el progreso necesario para continuar tras un fallo.

Objetivo
Structured Streaming ejecuta una consulta incremental como una secuencia de microbatches y conserva el progreso necesario para continuar tras un fallo.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

Una lectura con `readStream` describe un flujo que todavía no está en ejecución. Spark construye el plan, descubre qué datos nuevos están disponibles en cada trigger y solo entonces materializa un microbatch. La misma API de DataFrames permite reutilizar transformaciones batch, pero una operación válida en batch puede ser inviable en streaming si exige conservar estado sin límite.

En un pipeline de pedidos conviene separar tres decisiones: cómo se descubre la entrada, qué transformaciones son incrementales y cómo se confirma la salida. El SLA no lo decide `readStream`; lo determinan el trigger, el volumen por lote, la capacidad y el comportamiento del sink ante reintentos.

Modelo mental

Imagina Structured Streaming como un motor de consultas incrementales que avanza una frontera de datos confirmados, no como un bucle que vuelve a ejecutar una consulta batch completa. `readStream` declara una fuente cuyo final se mueve; las transformaciones construyen un plan lógico y `writeStream` crea una consulta activa. En cada microbatch, Spark determina un rango nuevo de entrada, ejecuta solo ese delta y coordina el resultado con el checkpoint. La API parece idéntica a DataFrames batch porque comparte álgebra, pero la viabilidad cambia: una proyección es stateless, mientras que contar por cliente obliga a recordar información entre lotes. Por eso una solución de producción se diseña alrededor de tres fronteras: qué progreso ofrece la fuente, cuánto estado requiere el plan y cómo confirma el sink. El trigger regula cuándo se intenta avanzar; no garantiza por sí mismo latencia, capacidad ni exactamente una vez.

Consulta incremental

Plan que procesa únicamente el rango de entrada aún no confirmado y conserva continuidad entre ejecuciones.

Evita razonar como si cada trigger recalculara toda la fuente y permite estimar coste y recuperación correctamente.
Microbatch

Unidad transaccional de progreso que agrupa un rango de offsets, su cómputo y el intento de escritura al sink.

Es la frontera práctica para reintentos, métricas, idempotencia y diagnóstico de latencia.
Operador stateful

Transformación que mantiene información de lotes anteriores, como una agregación, deduplicación o join stream-stream.

Su estado condiciona memoria, checkpoint, compatibilidad de cambios y necesidad de watermarks.
PySparkConsulta incremental mínima con destino Delta
orders = spark.readStream.table("main.bronze.orders")

query = (
    orders.where("order_id IS NOT NULL")
      .writeStream
      .option("checkpointLocation", "/Volumes/main/ops/checkpoints/orders_silver")
      .trigger(processingTime="1 minute")
      .toTable("main.silver.orders")
)

El objeto `query` permite consultar estado y progreso; el checkpoint debe ser exclusivo de esta consulta.

Puntos clave

  • `readStream` y `writeStream` definen una consulta continua; una acción batch no la inicia.
  • Cada microbatch procesa un rango identificable de entrada y lo confirma en el checkpoint.
  • Una transformación stateful requiere límites temporales y una estrategia de recuperación explícita.

Evita

  • Confundir la latencia del trigger con el tiempo real de proceso: si un lote tarda dos minutos, un trigger de diez segundos no crea capacidad adicional.
  • Usar el mismo checkpoint para dos consultas o destinos distintos, mezclando offsets y commits incompatibles.

Recuerdo activo

¿Qué parte del código inicia realmente la consulta?

Borrador privado · solo en este navegador
02
Implementación

Sources, sinks y output modes

El trigger expresa cuándo intentar procesar datos disponibles; `processingTime` prioriza cadencia y `availableNow` vacía el backlog y termina.

Objetivo
El trigger expresa cuándo intentar procesar datos disponibles; `processingTime` prioriza cadencia y `availableNow` vacía el backlog y termina.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

`processingTime` mantiene una consulta activa y lanza microbatches con la periodicidad solicitada, siempre que el anterior haya acabado. Es apropiado para un dashboard que debe actualizarse durante todo el día. Reducir el intervalo por debajo de la duración real del lote solo aumenta la presión de planificación.

`availableNow=True` procesa todos los datos disponibles en uno o varios microbatches y finaliza. Encaja con trabajos incrementales programados porque conserva checkpoints y límites por lote sin pagar por una consulta ociosa. También simplifica backfills controlados: el trabajo termina cuando alcanza el inicio del trigger.

Modelo mental

Un trigger es una política de servicio para una consulta incremental: decide cuándo pedir al motor que avance, pero no cambia la semántica del plan ni fabrica capacidad. `processingTime` mantiene la consulta activa y busca una cadencia recurrente; `availableNow` captura los datos disponibles para esa ejecución, los procesa en tantos microbatches como requieran los límites de la fuente y termina. La elección se parece a decidir entre un servicio residente y un trabajo incremental finito. Debe partir del SLA de frescura, la forma en que llegan los datos, el tiempo de arranque del compute y el coste de mantener recursos ociosos. Un intervalo de diez segundos no significa diez segundos de latencia si cada lote tarda dos minutos. Del mismo modo, `availableNow` no significa batch completo: conserva offsets, checkpoints y procesamiento incremental entre ejecuciones orquestadas.

Processing time trigger

Política que intenta iniciar microbatches repetidamente con una cadencia temporal mientras la consulta permanece activa.

Permite relacionar latencia continua con capacidad, pero no garantiza que cada lote termine dentro del intervalo.
Available Now

Trigger finito que procesa incrementalmente el conjunto disponible al inicio y termina tras confirmar todos sus microbatches.

Es la opción clave para cargas incrementales orquestadas que deben liberar compute sin perder checkpoints.
Frontera de ejecución

Límite superior de entrada que una ejecución concreta se compromete a alcanzar antes de finalizar.

Distingue los datos pertenecientes al run actual de los que quedarán para el siguiente y hace reproducible un backfill.
PySparkEjecución incremental finita para un Job
(
  spark.readStream.table("main.bronze.order_events")
    .writeStream
    .option("checkpointLocation", "/Volumes/main/ops/checkpoints/order_events")
    .trigger(availableNow=True)
    .toTable("main.silver.order_events")
    .awaitTermination()
)

La tarea del Job termina después de consumir el backlog disponible, pero la siguiente ejecución reanuda desde el mismo checkpoint.

Puntos clave

  • `availableNow` conserva semántica incremental y puede crear varios lotes hasta agotar la entrada.
  • Un trigger no sustituye el dimensionamiento ni controla por sí solo el tamaño del backlog.
  • La elección se basa en SLA, patrón de llegada y modelo operativo, no en preferencia de sintaxis.

Evita

  • Sustituir `availableNow` por una lectura batch y perder el seguimiento incremental de offsets.
  • Suponer que `availableNow` equivale a un único microbatch; puede dividir el backlog según los límites de la fuente.

Recuerdo activo

¿Qué ocurre si llegan archivos mientras una ejecución `availableNow` está activa?

Borrador privado · solo en este navegador
03
Operación

Triggers availableNow y processingTime

El checkpoint conserva offsets, commits y estado; forma parte de la identidad lógica de una consulta, no es una carpeta temporal.

Objetivo
El checkpoint conserva offsets, commits y estado; forma parte de la identidad lógica de una consulta, no es una carpeta temporal.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

Al reiniciar, Spark consulta el checkpoint para conocer qué rangos de la fuente fueron procesados y qué lotes se confirmaron en el sink. Las operaciones con estado añaden metadatos y datos del state store. Borrar la carpeta convierte el siguiente arranque en una consulta nueva y puede reintroducir datos o perder la posición esperada.

No todos los cambios de código son compatibles con un checkpoint existente. Cambiar el número de fuentes, el topic Kafka, el tipo de sink, las claves de una agregación o el esquema del estado suele exigir un checkpoint nuevo y un plan de migración o reproceso. Un simple filtro suele ser compatible, aunque puede cambiar el resultado funcional.

Modelo mental

El checkpoint es la memoria transaccional y operativa de una consulta streaming. No es un simple caché que pueda eliminarse para solucionar errores: contiene la identidad del plan, los rangos de fuente observados, los lotes comprometidos y, cuando hay operadores stateful, el esquema y contenido del estado. Piensa en él como el diario que permite responder qué se leyó, qué se publicó y qué debe restaurarse. La ruta pertenece a una única consulta lógica y a una versión compatible de su topología. Cambiar el topic, el tipo de fuente, el sink, las claves de agregación o el esquema del state store puede invalidar esa continuidad. En ese caso, un checkpoint nuevo no es una reparación neutra; crea una consulta nueva y obliga a decidir explícitamente desde qué punto reponer datos y cómo reconciliar la salida existente.

Offset log

Registro por microbatch de los límites de entrada que la consulta ha planificado y procesado.

Permite retomar desde una posición precisa y diagnosticar si el problema está antes o después de leer la fuente.
Commit log

Registro de lotes cuya salida se considera confirmada para la consulta.

Separa un intento incompleto de un avance durable y participa en la prevención de reprocesados indebidos.
Compatibilidad de estado

Condición por la que la topología, claves, esquema y proveedor del state store pueden restaurarse con un checkpoint existente.

Determina si un despliegue puede reanudar o necesita migración y reconstrucción deliberadas.
PythonRuta de checkpoint estable por entorno y pipeline
catalog = spark.conf.get("app.catalog", "main")
environment = spark.conf.get("app.environment", "dev")
checkpoint = f"/Volumes/{catalog}/ops/checkpoints/{environment}/orders_v1"

assert environment in {"dev", "test", "prod"}
print({"checkpoint": checkpoint, "query_version": 1})

Versiona la consulta cuando un cambio incompatible requiera un checkpoint nuevo; conserva el anterior hasta verificar el cutover.

Puntos clave

  • Un checkpoint pertenece a una consulta y debe residir en almacenamiento durable y gobernado.
  • Offsets y commits permiten reanudar; el state store permite reconstruir operadores stateful.
  • Cambiar la topología stateful sin evaluar compatibilidad puede impedir el reinicio.

Evita

  • Guardar checkpoints críticos en una ubicación efímera o eliminarlos como primera respuesta ante un fallo.
  • Reutilizar el checkpoint de producción durante pruebas, avanzando offsets o alterando estado real.

Recuerdo activo

¿Por qué no debe borrarse un checkpoint para 'forzar' un reintento?

Borrador privado · solo en este navegador
04
Diagnóstico

Offsets, commits y checkpoints

La garantía extremo a extremo depende tanto del seguimiento de offsets como de que el sink confirme lotes de forma idempotente.

Objetivo
La garantía extremo a extremo depende tanto del seguimiento de offsets como de que el sink confirme lotes de forma idempotente.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

Los sinks Delta integrados coordinan los commits de cada microbatch con el checkpoint. Si una tarea falla después de escribir pero antes de confirmar, el reinicio puede volver a presentar el mismo lote y el sink debe reconocerlo. La garantía de la fuente por sí sola no impide duplicados en un servicio externo.

`foreachBatch` abre la puerta a `MERGE`, múltiples destinos o APIs, pero transfiere responsabilidad al código. El `batch_id` y una clave de negocio estable permiten registrar el lote o hacer upsert determinista. Enviar eventos a un endpoint sin clave idempotente conserva, como mucho, semántica at-least-once.

Modelo mental

Exactly-once no es una propiedad que la fuente Kafka, el checkpoint o Delta puedan declarar aisladamente; es un resultado extremo a extremo. Spark puede volver a presentar un microbatch cuando no sabe si el intento anterior quedó publicado. Un sink transaccional integrado puede reconocer el lote y coordinarlo con el progreso, mientras que una API REST, un correo o dos destinos independientes no participan automáticamente en ese protocolo. Por eso el modelo mental correcto separa entrega de procesamiento: at-least-once admite reintentos y posible repetición; idempotencia hace que repetir produzca el mismo estado final. `foreachBatch` ofrece toda la expresividad batch, incluido `MERGE`, pero traslada al autor la responsabilidad de claves, orden y atomicidad. Un `batch_id` identifica un intento lógico de la consulta; una clave de negocio identifica el hecho y suele sobrevivir mejor a reconstrucciones con un checkpoint nuevo.

At-least-once

Garantía de que un registro no se pierde, aunque un fallo pueda provocar uno o más intentos de entrega.

Obliga a diseñar consumidores que toleren repetición en vez de inferir unicidad por la existencia del checkpoint.
Idempotencia

Propiedad por la que aplicar varias veces la misma operación deja el mismo estado que aplicarla una vez.

Convierte reintentos inevitables en recuperación segura para `MERGE`, APIs y efectos externos.
Clave de negocio

Identificador estable del hecho o entidad, independiente del lote y del intento técnico que lo transporta.

Permite deduplicar incluso después de reconstruir una consulta con otro checkpoint o redistribuir eventos.
PySparkUpsert idempotente por microbatch
from delta.tables import DeltaTable

def upsert_orders(batch_df, batch_id):
    target = DeltaTable.forName(spark, "main.silver.orders_current")
    (target.alias("t")
      .merge(batch_df.alias("s"), "t.order_id = s.order_id")
      .whenMatchedUpdateAll()
      .whenNotMatchedInsertAll()
      .execute())

(orders.writeStream
  .foreachBatch(upsert_orders)
  .option("checkpointLocation", "/Volumes/main/ops/checkpoints/orders_upsert")
  .start())

El `MERGE` debe resolver múltiples cambios de la misma clave dentro del lote antes del upsert.

Puntos clave

  • Exactly-once es una propiedad del recorrido fuente–estado–sink, no una etiqueta aislada.
  • `foreachBatch` permite lógica batch por microbatch y debe tolerar la repetición del mismo lote.
  • Una clave de negocio suele ser más útil que confiar únicamente en el número de lote.

Evita

  • Hacer `append` en `foreachBatch` y asumir que el checkpoint evita cualquier repetición posterior al fallo.
  • Ejecutar varias acciones sobre `batch_df` sin persistirlo cuando el coste de recomputación sea relevante.

Recuerdo activo

¿Qué debe garantizar una función `foreachBatch` para soportar reintentos?

Borrador privado · solo en este navegador
05
Decisión de diseño

Recuperación y compatibilidad de cambios

El progreso de una consulta permite separar falta de entrada, capacidad insuficiente, estado creciente y problemas del sink.

Objetivo
El progreso de una consulta permite separar falta de entrada, capacidad insuficiente, estado creciente y problemas del sink.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

`lastProgress` expone tasas de entrada y proceso, duración de fases, offsets y métricas de estado. Si `inputRowsPerSecond` supera sostenidamente `processedRowsPerSecond`, crece el backlog; si no hay filas nuevas, una latencia alta puede proceder del trigger o del propio sink. En Kafka, el desfase de offsets aporta otra señal independiente.

La operación necesita umbrales vinculados al SLA: frescura máxima, duración p95 del lote, filas descartadas y tamaño del estado. Un runbook útil conecta cada alerta con hipótesis y acciones seguras, evitando reiniciar o borrar checkpoints sin diagnóstico.

Modelo mental

Observar streaming consiste en reconstruir una cadena causal entre llegada, procesamiento, estado y publicación. Que una consulta esté `ACTIVE` solo dice que el driver no ha terminado; no demuestra que reciba eventos, avance offsets ni cumpla la frescura de negocio. El progreso de cada microbatch aporta una instantánea: filas y bytes de entrada, tasas, duración por fase, offsets y métricas de operadores stateful. La lectura correcta compara series, no una sola muestra. Si la entrada supera sostenidamente al proceso, crece el backlog; si ambos son bajos pero la duración es alta, puede dominar el sink, el arranque o un operador sin filas. La frescura se mide con event time del último dato válido frente al reloj y debe complementarse con completitud, fallos de calidad y latencia del consumidor. Así, cada alerta conduce a una hipótesis comprobable en vez de a reinicios ciegos.

Backlog

Cantidad de datos disponibles en la fuente que todavía no forman parte de un commit confirmado de la consulta.

Distingue falta de capacidad de ausencia de datos y permite estimar tiempo de recuperación.
Frescura

Diferencia entre el reloj de observación y el event time más reciente que llegó correctamente al producto de datos.

Mide el SLA percibido por el consumidor, que no se deriva del simple estado activo del proceso.
State operator metric

Métrica por operador sobre filas mantenidas, actualizadas, eliminadas y memoria o almacenamiento del estado.

Revela crecimiento no acotado y separa problemas stateful de cuellos en fuente o sink.
PythonInspección defensiva del último progreso
progress = query.lastProgress or {}
summary = {
    "batch_id": progress.get("batchId"),
    "input_rps": progress.get("inputRowsPerSecond", 0),
    "processed_rps": progress.get("processedRowsPerSecond", 0),
    "batch_ms": progress.get("durationMs", {}).get("triggerExecution"),
    "state": progress.get("stateOperators", []),
}
print(summary)

En producción, envía estas métricas a una tabla o plataforma de monitorización en vez de depender de impresiones del driver.

Puntos clave

  • Compara tasa de llegada con tasa de proceso a lo largo de varios lotes, no en una sola muestra.
  • `stateOperators` revela filas y memoria mantenidas por agregaciones, joins o deduplicación.
  • La frescura de negocio requiere comparar el máximo event time publicado con el reloj, no solo observar que el query está activo.

Evita

  • Interpretar una consulta `ACTIVE` como prueba de que cumple el SLA de frescura.
  • Aumentar compute sin comprobar si el cuello está en el sink, el estado, el throttling de la fuente o datos sesgados.

Recuerdo activo

¿Qué señal distingue un backlog creciente de una fuente momentáneamente vacía?

Borrador privado · solo en este navegador
5 lecciones pendientes

Fuente revisada · vista externa

Azure Databricks Hands-on

Azure Databricks Hands-on · commit a91650b

HandsOn.dbc

Archivo importable

Este notebook se abre desde su fuente revisada

El repositorio no permite republicar su contenido dentro de Lakehouse Lab. Conservamos la misma experiencia lateral, la ruta exacta y el commit auditado, y dejamos la lectura en GitHub para respetar la autoría.

Autor
Tsuyoshi Matsuzaki
Licencia
No verificada
Formato
dbc
Abrir / descargar .dbc