Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu 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.
- Duración
- 17 min aprox.
- Objetivo
- El progreso de una consulta permite separar falta de entrada, capacidad insuficiente, estado creciente y problemas del sink.
- Siguiente paso
- Continuar con el laboratorio
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
05Decisió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.
`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