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.
- Definir event time y processing time
- Aplicar watermark con intención
- Deduplicar y agregar ventanas acotando estado
03OperaciónWatermarks y finalización
La deduplicación con watermark mantiene identificadores durante un horizonte finito y elimina repeticiones sin hacer crecer el estado para siempre.
+
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.
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.
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.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.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.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