Saltar al contenido
Lakehouse LabLakehouse LabPreparación Databricks Data Engineer
Módulo 14 · Professional

Estado

Contenido abierto

Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.

Professional

Estado, ventanas, watermarks y datos tardíos

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

Lectura pública
Al terminar podrás
  • Definir event time y processing time
  • Aplicar watermark con intención
  • Deduplicar y agregar ventanas acotando estado
Ver fuentes y revisión

Metadatos editoriales

Última revisión
21 jul 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
01
Modelo mental

Event time frente a processing time

El event time pertenece al hecho de negocio; el processing time describe cuándo lo observa la plataforma y no corrige el desorden de llegada.

Objetivo
El event time pertenece al hecho de negocio; el processing time describe cuándo lo observa la plataforma y no corrige el desorden de llegada.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m13
Reportar un error en esta lección

Un pago generado a las 10:03 puede llegar a las 10:17 por una desconexión móvil. Agrupar por la hora de ingesta lo asignaría a una ventana distinta y haría que un reproceso produzca otro resultado. La columna de event time debe proceder del evento, convertirse a `timestamp` y validarse antes de cualquier operación temporal.

El retraso `processing_time - event_time` es una distribución, no una constante. Para elegir tolerancia se estudian percentiles y casos extremos por fuente. Un reloj del productor defectuoso debe enviarse a cuarentena; ampliar indefinidamente el watermark para ocultarlo traslada el problema al state store.

Modelo mental

Event time y processing time responden preguntas diferentes. Event time pertenece al hecho: cuándo ocurrió la compra, lectura o clic según el productor. Processing time pertenece a la plataforma: cuándo el evento fue observado y transformado. En un sistema distribuido, reintentos, desconexiones móviles, buffers y particiones hacen que el orden de llegada difiera del orden de negocio. Structured Streaming no puede corregir ese desorden mirando el reloj del cluster; necesita una columna de event time válida y una política explícita de tardanza. El modelo mental útil es una línea temporal que avanza con evidencia imperfecta: cada fuente revela máximos observados, el watermark deriva una frontera conservadora y los operadores deciden cuándo dejar de esperar. Antes de agregar, hay que normalizar zona horaria, precisión, valores imposibles y semántica del productor, porque un timestamp incorrecto puede adelantar la frontera y expulsar datos legítimos.

Event time

Instante en que ocurrió el hecho según el dominio productor, transportado como parte del evento.

Es la base correcta para ventanas, orden de negocio y análisis reproducible pese a retrasos de red.
Processing time

Instante en que el motor recibe o procesa el evento en una ejecución concreta.

Sirve para operación y triggers, pero produce resultados dependientes de retrasos y reejecuciones si se usa como tiempo de negocio.
Ingestion time

Marca añadida al entrar en una frontera controlada de la plataforma.

Permite medir retraso y detectar relojes anómalos sin sustituir la semántica del event time.
PySparkNormalización del tiempo de evento
from pyspark.sql import functions as F

events = (
    raw.withColumn("event_ts", F.to_timestamp("event_time_iso"))
       .withColumn("ingested_at", F.current_timestamp())
       .withColumn(
           "lateness_seconds",
           F.col("ingested_at").cast("long") - F.col("event_ts").cast("long")
       )
)

valid = events.where("event_ts IS NOT NULL AND lateness_seconds >= 0")

Conserva ambos tiempos: uno gobierna la semántica y el otro permite medir el comportamiento de la fuente.

Puntos clave

  • Event time determina ventanas reproducibles; processing time mide la observación del sistema.
  • La calidad y zona horaria del timestamp son parte del contrato del evento.
  • La distribución de tardanza informa el watermark y el SLA de correcciones.

Evita

  • Usar `current_timestamp()` como event time porque siempre está presente, haciendo que un replay cambie los resultados.
  • Aceptar timestamps sin zona o muy futuros y permitir que adelanten prematuramente el watermark.

Recuerdo activo

¿Por qué un reproceso debe conservar el event time original?

Borrador privado · solo en este navegador
02
Implementación

Ventanas tumbling y sliding

Una ventana agrupa por event time y el watermark limita cuánto estado se conserva antes de considerar una ventana finalizable.

Objetivo
Una ventana agrupa por event time y el watermark limita cuánto estado se conserva antes de considerar una ventana finalizable.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m13
Reportar un error en esta lección

Las ventanas tumbling no se solapan; las sliding pueden asignar el mismo evento a varias ventanas. En una agregación, Spark mantiene acumuladores para ventanas todavía abiertas. El watermark avanza a partir del máximo event time observado menos el retraso configurado y permite retirar estado antiguo.

Una tolerancia corta reduce memoria y acelera resultados finales, pero descarta más eventos tardíos. Una tolerancia larga mejora completitud a costa de latencia y estado. La decisión debe expresar un compromiso medible, por ejemplo aceptar el 99,5% de eventos en quince minutos y corregir el resto mediante un proceso separado.

Modelo mental

Una ventana transforma un flujo infinito en grupos temporales finitos; un watermark proporciona la evidencia para dejar de mantener algunos de esos grupos. No es un temporizador que espera exactamente N minutos después de cada evento ni una promesa de que todo lo anterior se descartará en un instante preciso. Spark observa el máximo event time visto y resta el retraso configurado para obtener una frontera. Una ventana cuyo final queda suficientemente atrás puede cerrarse y su estado eliminarse según el modo de salida. El compromiso es explícito: ampliar la tardanza protege más eventos desordenados, pero conserva más claves y ventanas, aumenta checkpoint y retrasa resultados finales. Reducirla mejora coste y latencia a cambio de una ruta de excepciones. La columna marcada debe ser la misma que alimenta `window`; perder su metadato mediante expresiones mal ubicadas puede impedir el comportamiento esperado.

Ventana tumbling

Intervalos contiguos de tamaño fijo sin solapamiento, donde cada evento pertenece a una sola ventana.

Simplifica totales por minuto u hora y limita el número de acumuladores activos por clave.
Ventana sliding

Intervalos de tamaño fijo iniciados con una cadencia menor, de modo que un evento puede pertenecer a varias ventanas.

Permite métricas móviles, pero multiplica actualizaciones y estado respecto a una ventana tumbling.
Watermark

Frontera derivada del máximo event time observado menos una tolerancia de tardanza.

Da al motor una condición para limpiar estado y al negocio una política cuantificable sobre datos tardíos.
PySparkVentas por ventanas de cinco minutos
from pyspark.sql import functions as F

sales_5m = (
    events.withWatermark("event_ts", "15 minutes")
      .groupBy(F.window("event_ts", "5 minutes"), "store_id")
      .agg(
          F.countDistinct("order_id").alias("orders"),
          F.sum("amount").alias("revenue")
      )
)

El retraso de quince minutos debe justificarse con la distribución real de tardanza y el SLA de publicación.

Puntos clave

  • Watermark no significa esperar exactamente ese tiempo desde la llegada de cada fila.
  • La columna marcada debe participar en la ventana o condición temporal del operador stateful.
  • Output mode y watermark determinan cuándo se emiten y retiran resultados.

Evita

  • Aplicar `withWatermark` a una columna y agrupar por otra, impidiendo que el operador use la marca temporal esperada.
  • Elegir un watermark igual al intervalo de trigger; resuelven problemas diferentes.

Recuerdo activo

¿Qué se sacrifica al reducir de una hora a diez minutos el watermark?

Borrador privado · solo en este navegador
03
Operación

Watermarks y finalización

La deduplicación con watermark mantiene identificadores durante un horizonte finito y elimina repeticiones sin hacer crecer el estado para siempre.

Objetivo
La deduplicación con watermark mantiene identificadores durante un horizonte finito y elimina repeticiones sin hacer crecer el estado para siempre.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m13
Reportar un error en esta lección

`dropDuplicatesWithinWatermark(["event_id"])` recuerda las claves observadas mientras el evento pueda seguir teniendo duplicados dentro de la tolerancia. La marca temporal debe definirse antes. Si dos copias pueden diferir en timestamp, la tolerancia debe superar la máxima distancia temporal entre ellas para garantizar la deduplicación.

La clave debe representar identidad de negocio, no una combinación accidental de todas las columnas. Antes de deduplicar conviene validar nulos y normalizar tipos. Los eventos que llegan más tarde que el umbral pueden descartarse; la métrica de descarte y una ruta de reconciliación son parte del diseño.

Modelo mental

Deduplicar un stream significa recordar identidades ya observadas durante el periodo en que todavía podrían repetirse. Sin un límite, cualquier identificador histórico podría reaparecer y el state store tendría que crecer para siempre. El watermark acota esa obligación temporal, pero la clave y el timestamp deben representar el contrato real. `dropDuplicatesWithinWatermark` está pensado para mantener la clave dentro del horizonte aun cuando los duplicados lleven timestamps ligeramente distintos; combinar `dropDuplicates` con una columna temporal en la clave puede tratar cada timestamp como un registro diferente. La deduplicación técnica tampoco reemplaza la idempotencia del destino: un evento puede salir una sola vez de este operador y aun así repetirse por un fallo en una llamada externa. Primero se decide qué constituye el mismo hecho, cuánto tiempo puede repetirse y qué hacer con una repetición más antigua que la retención.

Clave de deduplicación

Conjunto mínimo de atributos que identifica de forma estable un único hecho de negocio.

Una clave incorrecta elimina eventos legítimos o deja pasar reintentos como si fueran hechos distintos.
Horizonte de deduplicación

Periodo durante el cual el motor conserva una clave porque todavía acepta la llegada de una repetición.

Controla directamente el equilibrio entre corrección ante reintentos tardíos y tamaño del state store.
dropDuplicatesWithinWatermark

Operador que deduplica por claves dentro de la tolerancia temporal marcada sin exigir el timestamp como parte de la identidad.

Resuelve productores que reemiten el mismo identificador con pequeñas variaciones temporales y mantiene estado acotado.
PySparkDeduplicación acotada por identificador
deduplicated = (
    events.where("event_id IS NOT NULL AND event_ts IS NOT NULL")
      .withWatermark("event_ts", "30 minutes")
      .dropDuplicatesWithinWatermark(["event_id"])
)

Mide duplicados y eventos tardíos con datos de prueba que crucen varios microbatches.

Puntos clave

  • `distinct` sin watermark puede conservar todas las filas únicas indefinidamente.
  • `dropDuplicatesWithinWatermark` requiere una marca temporal en el DataFrame streaming.
  • El horizonte debe cubrir la separación máxima esperada entre copias del mismo evento.

Evita

  • Deduplicar por todas las columnas cuando dos copias llevan distinto `ingested_at` y, por tanto, nunca coinciden.
  • Escoger un umbral menor que el retraso entre duplicados y prometer una garantía que el estado ya no puede cumplir.

Recuerdo activo

¿Qué condición permite retirar del estado un `event_id` antiguo?

Borrador privado · solo en este navegador
04
Diagnóstico

Deduplicación con estado

Un join entre dos streams necesita watermarks y una condición temporal para que Spark pueda descartar pares que ya no pueden coincidir.

Objetivo
Un join entre dos streams necesita watermarks y una condición temporal para que Spark pueda descartar pares que ya no pueden coincidir.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m13
Reportar un error en esta lección

Relacionar clics con pagos solo por `session_id` deja abierta para siempre la posibilidad de una coincidencia futura. Añadir que el pago ocurra entre el clic y treinta minutos después proporciona un intervalo acotado. Cada entrada debe declarar su tolerancia de tardanza para que el motor calcule cuándo retirar estado.

Los joins outer requieren esperar a que transcurra el horizonte antes de emitir una fila no emparejada. Eso aumenta la latencia de los `NULL` esperados. En múltiples streams, la política global por defecto avanza con el watermark más lento; cambiarla a `max` reduce latencia, pero puede descartar datos del stream rezagado.

Modelo mental

Un join stream-stream intenta emparejar dos conjuntos que nunca dejan de crecer. Una igualdad de clave no basta para saber cuándo eliminar una fila no emparejada: en teoría, su pareja podría llegar años después. Para acotar estado se necesitan watermarks en las entradas y una condición temporal que limite qué pares son válidos, por ejemplo una impresión ocurrida entre cero y treinta minutos antes de una compra. El motor combina esas restricciones para decidir cuándo un evento ya no puede encontrar pareja futura. Los inner joins pueden ejecutarse sin watermark, pero retendrían estado sin límite; los outer joins necesitan límites para determinar cuándo emitir el lado nulo. El modelo mental es doble: la clave reduce candidatos y el intervalo temporal cierra la búsqueda. El watermark global y la fuente más lenta condicionan cuándo se limpia y cuándo aparecen resultados no emparejados.

Condición temporal

Predicado que limita la diferencia permitida entre los event times de dos filas candidatas al join.

Hace posible demostrar que una fila antigua ya no podrá emparejarse y permite limpiar su estado.
Watermark global

Frontera derivada de los watermarks de varias entradas que gobierna operadores stateful conjuntos.

Explica por qué una fuente lenta puede retener estado y retrasar salidas aunque la otra avance.
Stream-static join

Join entre un flujo incremental y una relación batch leída como referencia durante la consulta.

Suele requerir menos estado que stream-stream y es apropiado cuando un lado no necesita emparejamiento temporal continuo.
PySparkJoin temporal de clics y pagos
clicks_wm = clicks.withWatermark("click_ts", "20 minutes")
payments_wm = payments.withWatermark("payment_ts", "10 minutes")

matched = clicks_wm.join(
    payments_wm,
    (clicks_wm.session_id == payments_wm.session_id) &
    (payments_wm.payment_ts >= clicks_wm.click_ts) &
    (payments_wm.payment_ts <= clicks_wm.click_ts + F.expr("INTERVAL 30 MINUTES")),
    "leftOuter",
)

Documenta por separado la tardanza de cada fuente y el máximo intervalo de negocio entre los dos eventos.

Puntos clave

  • Las claves de igualdad no acotan por sí solas el estado de un stream-stream join.
  • La condición de rango temporal y los watermarks trabajan conjuntamente.
  • Un outer join no puede declarar una fila sin pareja hasta que expire la posibilidad de coincidencia.

Evita

  • Añadir watermarks pero omitir el rango temporal del join, por lo que el estado sigue sin una frontera útil.
  • Esperar que un outer join emita inmediatamente los no emparejados, antes de que expire el watermark.

Recuerdo activo

¿Por qué un left outer join retrasa las filas sin pago?

Borrador privado · solo en este navegador
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.

Objetivo
Las métricas del state store revelan si el compromiso entre tardanza, cardinalidad y capacidad sigue siendo sostenible.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m13
Reportar un error en esta lección

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

Modelo mental

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

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.

Recuerdo activo

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

Borrador privado · solo en este navegador
5 lecciones pendientes

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