1. Pregunta
Un flujo de eventos de órdenes agrega el monto y la cantidad de órdenes por customerId en ventanas de cinco minutos. Los relojes de los dispositivos pueden desincronizarse, los reintentos de red generan entregas fuera de orden y algunas particiones de Kafka pueden no tener mensajes nuevos temporalmente. Los resultados deben aparecer rápidamente mientras se permite que los datos tardíos corrijan una ventana durante un período acotado. Diseña el manejo del tiempo de evento.
2. Restricciones y aclaraciones
- Confirma si las ventanas utilizan tiempo de evento, tiempo de escritura o tiempo de procesamiento; este problema utiliza tiempo de evento.
- Establece un límite máximo de desorden, como 30 segundos, y define qué sucede más allá de este.
- Decide si los resultados son eventos de solo adición (append-only) o instantáneas actualizables; los consumidores downstream deben identificar las revisiones.
- Analiza particiones inactivas, marcas de tiempo inválidas, eventos duplicados y trabajos de reproducción (replay) para que una entrada sin mensajes no se confunda con el fin del flujo.
3. Conceptos clave
El tiempo de evento proviene del propio registro. Una marca de agua indica que el sistema cree que el tiempo de evento ha avanzado hasta una posición determinada. Comúnmente, una ventana se activa cuando la marca de agua supera su final; un registro que llega después de esa marca de agua cuyo timestamp aún pertenece a la ventana se considera tardío. Cuando se combinan entradas paralelas, un operador comúnmente espera a la marca de agua de entrada mínima, por lo que una partición inactiva puede retrasar el progreso global. Por lo tanto, se requiere detección de inactividad o una política de tiempo de espera por partición.
4. Flujo de referencia
onRecord(event):
ts = extractEventTimestamp(event)
key = canonicalKey(event.customerId)
updateWatermarkGenerator(ts)
window = floorToFiveMinutes(ts)
if ts <= currentWatermark + allowedLateness:
state[window, key] = aggregate(state[window, key], event)
emitUpsert(window, key, state[window, key], revision + 1)
else:
routeToLateData(window, key, event)
onWatermark(wm):
finalizeWindowsBefore(wm)
expireStateAfterRetention()La fuente extrae las marcas de tiempo de los eventos y genera una marca de agua con desorden acotado (bounded-out-of-orderness). Marca una partición como inactiva tras un período prolongado sin datos para que no detenga la marca de agua fusionada. Retén el estado de la ventana durante el período de retraso permitido; actualiza los resultados con un evento de upsert o de corrección. Enruta los datos que superen el límite a una salida lateral (side output) para su revisión o backfill offline.
5. Compensaciones entre precisión y latencia
Un retraso permitido más largo hace que la vista de tiempo de evento sea más completa, pero incrementa el estado y el volumen de correcciones. Emitir un único resultado final minimiza la latencia downstream pero no puede representar revisiones tardías. Los eventos de actualización que transportan una clave de ventana, revisión y motivo requieren fusiones idempotentes downstream. Para métricas no tolerantes a pérdidas como el dinero, retén los eventos sin procesar y programa un recálculo offline; para una tabla de clasificación en tiempo real, puede ser aceptable actualizar solo en el siguiente lote tras el límite.
6. Verificación y observabilidad
- Genera eventos ordenados, fuera de orden, con retraso menor a 30 segundos y posteriores al límite; compara cada ventana con una línea base offline.
- Inyecta particiones inactivas, saltos de reloj, eventos duplicados y reinicios de tareas; verifica que las marcas de agua no se detengan ni retrocedan.
- Registra la marca de agua actual, el desfase entre tiempo de procesamiento y tiempo de evento, el tamaño del estado de la ventana, la tasa de eventos tardíos, el volumen de salidas laterales y el recuento de correcciones.
- Incluye la clave de ventana, la revisión y el ID del evento de entrada en cada actualización; reproduce el mismo registro y compara las instantáneas finales.
7. Errores comunes
- Reemplazar ventanas de tiempo de evento por ventanas de tiempo de procesamiento mientras se afirma tener resistencia al desorden de red.
- Tratar una marca de agua como un marcador de finalización absoluto; en la práctica puede ser una heurística basada en una suposición de retraso.
- Ignorar particiones inactivas, permitiendo que una partición silenciosa bloquee el cierre de todas las ventanas.
- Descartar datos tardíos sin una salida lateral, versión o ruta de backfill offline.
8. Puntos de evaluación en entrevistas
Distingue las tres nociones de tiempo
El candidato debe definir el tiempo de evento, de ingestión y de procesamiento, y explicar cómo la elección de la ventana altera la precisión y la latencia.
Genera y fusiona marcas de agua
La respuesta debe cubrir el desorden acotado, la marca de agua mínima a través de entradas paralelas y el manejo de inactividad, señalando que una marca de agua es una estimación de progreso en lugar de una promesa absoluta.
Diseña rutas de datos tardíos y de corrección
El candidato debe definir un límite de retraso permitido, salida lateral, revisión y fusión idempotente downstream tanto para datos dentro del límite como posteriores a este.
Verifica el resultado mediante reproducción (replay)
El candidato debe probar desorden, inactividad, reinicios y duplicados contra una línea base offline y métricas en vivo, en lugar de verificar únicamente si el trabajo sigue ejecutándose.