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.
- Explicar microbatches y progreso
- Configurar triggers y checkpoints
- Diseñar sinks idempotentes
01Modelo mentalModelo 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.
+
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
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.
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.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.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.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 navegador02Implementació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.
- 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
`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
¿Qué ocurre si llegan archivos mientras una ejecución `availableNow` está activa?
Borrador privado · solo en este navegador03OperaciónTriggers 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.
+
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
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.
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.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.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.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 navegador04DiagnósticoOffsets, 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.
+
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
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.
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.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.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.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 navegador05Decisión de diseñoRecuperación y compatibilidad de cambios
El progreso de una consulta permite separar falta de entrada, capacidad insuficiente, estado creciente y problemas del sink.
+
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
`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.
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.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.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.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