Consigna y contexto
Los pedidos y los pagos se unen por ID de usuario y de pedido, pero un pago puede llegar dos horas después del pedido. Diseñe los límites, las marcas de agua (watermarks), la retención de estado, la ruta para datos tardíos y la estrategia de corrección para el Event Time Interval Join de Flink. Distinga entre event time, processing time e ingestion time, y explique por qué el orden de llegada de los mensajes es insuficiente.
Qué evalúa el entrevistador
- Definir un intervalo de tiempo de negocio en lugar de adivinar el orden a partir del processing time.
- Comprender que una marca de agua es una señal de progreso, no un tiempo globalmente verdadero.
- Explicar el estado de dos flujos, la limpieza, los datos tardíos y el manejo de eventos duplicados.
- Equilibrar latencia, completitud, costo de estado y capacidad de reproducción (replayability).
Preguntas aclaratorias para hacer
- ¿Cuáles son los campos de tiempo de negocio en los pedidos y pagos, y pueden desfasarse sus relojes?
- ¿Qué tan tarde puede llegar un pago, y se requiere una corrección o conciliación manual después del intervalo?
- ¿Es única la clave de unión (join key) y pueden existir duplicados, cancelaciones o múltiples pagos?
- ¿El sistema downstream acepta hechos de solo anexado (append-only), actualizaciones o únicamente una tabla de conciliación final?
Estructura de respuesta de 30 segundos
Defina el intervalo en tiempo de negocio, como un pago de cero a dos horas después del pedido. Particione ambos flujos por clave, avance las marcas de agua de tiempo de evento y permita que el join mantenga los registros de ambos lados hasta que el progreso demuestre que ya no pueden coincidir. Los eventos tardíos dentro del límite permitido pueden unirse; los eventos más allá de este van a una salida lateral (side output) o a un flujo de compensación. Concluya con desduplicación, checkpoints, reproducción e idempotencia downstream.
Análisis detallado paso a paso
1. Elegir el tiempo de evento y los límites del intervalo
Extraiga una marca de tiempo de evento de negocio inmutable de cada registro y utilice la misma clave de usuario y pedido. Si el pago debe ser posterior al pedido, utilice un límite inferior de cero y un límite superior de dos horas; permita un límite inferior negativo solo si el pago anticipado es válido. El intervalo proviene del SLA de negocio, no de una ventana arbitrariamente larga. El processing time es adecuado únicamente cuando el orden histórico no importa.
2. Marcas de agua y estado de dos flujos
Cada entrada crea una marca de agua a partir de su propio límite de desorden (out-of-order bound). El join espera hasta que el progreso en ambos lados sea suficiente para saber que un registro no puede encontrar otra coincidencia; hasta entonces, los registros permanecen en el estado con clave (keyed state). El tamaño del estado depende de la tasa de entrada, la longitud del intervalo, la cardinalidad de la clave y el límite de desorden. Los checkpoints persisten ese estado para que la recuperación continúe desde una posición conocida en lugar de adivinar qué resultados ya se emitieron.
3. Eventos tardíos, duplicados y cancelaciones
Los eventos dentro del rango permitido de desorden y retraso pueden unirse. Los eventos fuera del rango van a una salida lateral o a un tópico de compensación duradero para su conciliación. Desduplique con un ID de evento o clave de negocio, y modele la cancelación o reembolso de pagos como un nuevo evento o una retracción explícita. Si el downstream admite actualizaciones, emita upserts o retracciones; de lo contrario, mantenga una tabla de corrección en lugar de reescribir silenciosamente el historial.
4. Latencia, costo de estado y validación
Límites más cortos y restricciones de desorden reducen el estado y la latencia, pero aumentan las coincidencias perdidas. Límites más amplios mejoran la completitud mientras incrementan el costo de memoria, checkpoints y recuperación. Antes del lanzamiento, reproduzca desorden, duplicados, eventos entre ventanas y fallas de recuperación. Verifique la tasa de aciertos del join, el volumen de salidas laterales, el retraso de marcas de agua, el tamaño del estado, la duración de checkpoints y la tasa de duplicados. Las claves de idempotencia de extremo a extremo evitan cobros duplicados durante el reinicio o la reproducción.
Respuesta modelo
Confirmaría las marcas de tiempo de negocio de pedidos y pagos y el SLA de retraso permitido, luego particionaría por ID de usuario y de pedido. Si el pago solo puede llegar dentro de las dos horas posteriores a un pedido, expresaría un intervalo de tiempo de evento de cero a dos horas en lugar de una ventana de processing time. Cada flujo emite marcas de agua; el join mantiene los registros en keyed state y los limpia cuando el progreso demuestra que no quedan coincidencias.
Los eventos tardíos dentro del límite participan; los eventos fuera de él van a una salida lateral o flujo de compensación. Desduplicaría por ID de evento y representaría reembolsos y cancelaciones como nuevos eventos. Emitiría upserts o retracciones cuando sea compatible; de lo contrario, mantendría una tabla de corrección. Reproduciría desorden, duplicados y recuperación antes del lanzamiento, observaría la tasa de aciertos, la salida lateral, el retraso de marcas de agua, el estado y los checkpoints, y usaría claves de idempotencia para la seguridad en la reproducción.
Errores comunes
- Reemplazar el tiempo de evento de negocio con el tiempo de llegada del mensaje.
- Tratar una marca de agua como prueba de que todos los eventos upstream han llegado.
- Elegir una ventana grande sin discutir la limpieza del estado, los checkpoints y la recuperación.
- Descartar eventos tardíos fuera de la ventana sin una salida lateral o ruta de conciliación.
- Omitir la desduplicación y la idempotencia downstream, produciendo resultados de pago duplicados después del reinicio.
Preguntas de seguimiento y respuestas
Pregunta de seguimiento 1: ¿Por qué no usar dos ventanas independientes y un join normal?
Las ventanas independientes pierden el progreso del tiempo de evento y el límite de limpieza de los dos flujos, lo que dificulta expresar un intervalo de tiempo relativo. Interval Join combina la coincidencia de claves con límites inferiores y superiores explícitos para esta relación.
Pregunta de seguimiento 2: ¿Qué hace cuando una marca de agua se estanca?
Verifique la inactividad de las particiones (partition idleness), las marcas de tiempo de origen, la contrapresión (backpressure) y la configuración de desorden. Configure la inactividad para particiones genuinamente inactivas de modo que una partición vacía no bloquee el progreso, pero nunca avance las marcas de agua arbitrariamente para ocultar una falla de origen.
Pregunta de seguimiento 3: ¿Qué pasa si el negocio acepta pagos que llegan después de dos horas?
Separe la unión en tiempo real de la conciliación. Emita un resultado provisional, persista los eventos tardíos en un flujo de compensación y permita que un trabajo por lotes o un segundo trabajo de streaming genere una corrección con un upsert idempotente.