Saltar al contenido

Streaming

Menú

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

Guardar progreso

Lección 5 de 5

Recuperación y compatibilidad de cambios

El progreso de una consulta permite separar falta de entrada, capacidad insuficiente, estado creciente y problemas del sink.

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
05
Decisión de diseño

Recuperación y compatibilidad de cambios

El progreso de una consulta permite separar falta de entrada, capacidad insuficiente, estado creciente y problemas del sink.

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

PythonInspección defensiva del último progreso
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.

¿Qué señal distingue un backlog creciente de una fuente momentáneamente vacía?

Profundiza

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.

Backlog

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

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.
State operator metric

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

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.

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