Saltar al contenido

Estado

Menú

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

Guardar progreso

Lección 3 de 5

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.

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

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

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.

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

Profundiza

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

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.

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