Saltar al contenido

Streaming

Menú

Puedes leer sin crear un espacio. Créalo solo cuando quieras guardar.

Guardar progreso

Lección 2 de 5

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.

17 min aprox.

Detalles

Structured Streaming, triggers y checkpoints

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

Reto observable

Reinicia una consulta streaming desde su checkpoint y demuestra mediante offsets, conteos y claves que no pierde ni duplica eventos.

Al terminar podrás
  • Explicar microbatches y progreso
  • Configurar triggers y checkpoints
  • Diseñar sinks idempotentes
Prerrequisitos
m12
Última revisión
25 ago 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
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.

`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.

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.

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

Profundiza

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.
Resumen

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.

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

Módulo 13

Contenido del módulo