Saltar al contenido

Estado

Menú

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

Guardar progreso

Lección 5 de 5

State store, métricas y presión de memoria

Las métricas del state store revelan si el compromiso entre tardanza, cardinalidad y capacidad sigue siendo sostenible.

17 min aprox.

Detalles

Estado, ventanas, watermarks y datos tardíos

Controla el crecimiento de estado y la corrección temporal con eventos fuera de orden.

Reto observable

Configura ventanas y watermark para un SLA temporal, e identifica con datos de frontera qué eventos se actualizan, descartan y mantienen en estado.

Al terminar podrás
  • Definir event time y processing time
  • Aplicar watermark con intención
  • Deduplicar y agregar ventanas acotando estado
Prerrequisitos
m13
Última revisión
25 ago 2026
Nivel
Professional
Ruta relacionada
streaming
Dominios blueprint
Streaming state · Reliability
Estado
Revisión editorial interna
Fuentes principales
Apply watermarks to control data processing thresholds · Databricks · dropDuplicatesWithinWatermark · Databricks PySpark reference
Reportar un error
05
Decisión de diseño

State store, métricas y presión de memoria

Las métricas del state store revelan si el compromiso entre tardanza, cardinalidad y capacidad sigue siendo sostenible.

`stateOperators` informa filas totales y actualizadas, memoria usada y filas retiradas por watermark. Un crecimiento continuo después de varios horizontes puede indicar claves de cardinalidad extrema, timestamps futuros, ausencia de watermark efectivo o una condición de join no acotada.

El diagnóstico debe reproducir datos puntuales: a tiempo, duplicados, tardíos dentro de tolerancia, demasiado tardíos y timestamps inválidos. Una prueba de recuperación reinicia desde el mismo checkpoint; cambiar access mode o esquema de estado durante el experimento puede introducir otra variable.

PythonResumen de operadores con estado
progress = query.lastProgress or {}
for operator in progress.get("stateOperators", []):
    print({
        "rows_total": operator.get("numRowsTotal"),
        "rows_updated": operator.get("numRowsUpdated"),
        "rows_removed": operator.get("numRowsRemoved"),
        "memory_bytes": operator.get("memoryUsedBytes"),
    })

Guarda la serie temporal; una única observación no permite distinguir una ventana grande de una fuga lógica de estado.

¿Qué patrón indica que un watermark no está retirando estado como se esperaba?

Profundiza

El state store es una base de datos incremental local y checkpointada que materializa la memoria de agregaciones, deduplicaciones y joins. Su tamaño no depende solo de filas por segundo: depende del número de claves vivas, ventanas simultáneas, versiones por clave, retraso aceptado y velocidad con que el watermark permite evacuar. Por eso un stream de bajo volumen pero millones de dispositivos únicos puede ser más difícil que uno denso con pocas claves. Las métricas deben leerse como un balance: filas actualizadas entran, filas eliminadas salen y el total retenido refleja la deuda de estado. Memoria, latencia de commit y tamaño de checkpoint crecen de manera distinta según el proveedor. RocksDB con changelog checkpointing puede reducir presión y duración de checkpoints en cargas stateful, pero cambiar el mecanismo de estado de una consulta ya iniciada puede exigir un checkpoint nuevo y reconstrucción.

State store

Almacén versionado por operador que mantiene claves y valores necesarios entre microbatches y participa en la recuperación.

Es el principal determinante de estabilidad para operaciones stateful y debe dimensionarse como parte del diseño, no al final.
Cardinalidad viva

Número de claves y ventanas que todavía no pueden eliminarse según la política temporal del operador.

Predice mejor el tamaño del estado que el throughput bruto de entrada.
Changelog checkpointing

Estrategia que persiste cambios del estado entre snapshots periódicos en lugar de cargar un snapshot completo en cada checkpoint.

Puede reducir latencia de checkpoint para grandes estados RocksDB, pero requiere una decisión de arquitectura y recuperación compatible.
Resumen

Puntos clave

  • El número de filas de estado debe estabilizarse para una carga estacionaria y un horizonte finito.
  • Los timestamps futuros pueden empujar el watermark y provocar pérdidas silenciosas de eventos normales.
  • La política global `max` favorece latencia y puede descartar datos del stream más lento; `min` prioriza seguridad.

Evita

  • Medir solo uso de CPU y omitir filas/memoria del estado.
  • Cambiar a la política `max` para ocultar un stream lento sin aceptar explícitamente los descartes resultantes.

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 14

Contenido del módulo