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

Watermarks y finalización

Contenido abierto

Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu 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.

Duración
17 min aprox.
Objetivo
La deduplicación con watermark mantiene identificadores durante un horizonte finito y elimina repeticiones sin hacer crecer el estado para siempre.
Siguiente paso
Continuar con la siguiente lección
Ver detalles del módulo

Estado, ventanas, watermarks y datos tardíos

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

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

Mostrar prerrequisitos
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
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