Saltar al contenido

Streaming

Menú

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

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

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

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.

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.

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

Profundiza

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

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.

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