Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Lección 4 de 5
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.
- Duración
- 17 min aprox.
- Objetivo
- Un join entre dos streams necesita watermarks y una condición temporal para que Spark pueda descartar pares que ya no pueden coincidir.
- 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
04DiagnósticoDeduplicació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.
+
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.
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.
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.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.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.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