Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu 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.
- Duración
- 17 min aprox.
- Objetivo
- El trigger expresa cuándo intentar procesar datos disponibles; `processingTime` prioriza cadencia y `availableNow` vacía el backlog y termina.
- Siguiente paso
- Continuar con la siguiente lección
Ver detalles del módulo
Structured Streaming, triggers y checkpoints
Construye consultas incrementales con estado durable y semántica de recuperación clara.
- Explicar microbatches y progreso
- Configurar triggers y checkpoints
- Diseñar sinks idempotentes
02ImplementaciónSources, sinks y output modes
El trigger expresa cuándo intentar procesar datos disponibles; `processingTime` prioriza cadencia y `availableNow` vacía el backlog y termina.
+
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.
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.
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.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.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.(
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