Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Lección 3 de 5
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.
- Duración
- 17 min aprox.
- Objetivo
- El checkpoint conserva offsets, commits y estado; forma parte de la identidad lógica de una consulta, no es una carpeta temporal.
- 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
03Operació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.
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