Problema y alcance
Los clientes móviles reportan eventos de compra con al menos event_id, event_time, user_id y un amount ya convertido a una única moneda. Los reintentos de red pueden generar duplicados, los dispositivos fuera de línea pueden subir datos varias horas después y las diferentes particiones no preservan un orden global. Calcule los ingresos para cada hora del reloj UTC por event_time, con una entrada máxima de 20,000 eventos por segundo.
El negocio requiere una estimación para la hora actual en un plazo de un minuto, un resultado puntual una vez que la marca de agua supere el final de la ventana, y correcciones automáticas durante 24 horas después de que finalice dicha ventana. Los datos que lleguen después de las 24 horas no deben desaparecer silenciosamente; se envían a auditoría y conciliación fuera de línea. El rendimiento, la puntualidad y el período de corrección son entradas del problema, no afirmaciones de rendimiento sobre un motor de procesamiento de flujos.
Asuma que el productor asigna un event_id estable y que los reintentos válidos con ese ID contienen el mismo contenido de negocio. Los eventos sin procesar permanecen en un almacenamiento reproducible. El alcance cubre semántica temporal, marcas de agua, disparadores de ventanas, desduplicación, actualizaciones de resultados, capacidad del estado, recuperación y conciliación. La selección del broker y la conversión de divisas quedan fuera del alcance. Esta es una pregunta de ingeniería de datos: el objetivo es un contrato de datos verificable que haga explícito el balance entre completitud, latencia y costo de estado.
Qué evalúan los entrevistadores
La primera señal es si el candidato separa tres relojes. event_time determina la hora comercial que contiene un evento. processing_time indica cuándo lo ve el motor y puede impulsar una actualización anticipada de un minuto. Una marca de agua es la estimación del motor sobre el progreso del tiempo del evento. Si el tiempo de procesamiento determina la ventana, reproducir el mismo historial en un momento diferente puede producir resultados por hora distintos.
La segunda señal es tratar una marca de agua como una estimación de progreso en lugar de una promesa absoluta. Avanzarla rápidamente proporciona salidas oportunas, pero clasifica más registros como tardíos. Retenerla mejora la completitud, pero retrasa el cierre de la ventana y retiene más estado. Una respuesta sólida deriva la política a partir de datos medidos de retraso de llegada y el SLO de corrección, en lugar de memorizar un retraso fijo de cinco minutos o una hora.
La tercera señal es separar la marca de agua, el retraso permitido y la retención de desduplicación. La marca de agua controla la emisión puntual. El retraso permitido controla cuánto tiempo permanece corregible el estado de la ventana. El estado de desduplicación debe abarcar el período en el que puede llegar otra copia de un evento. Los valores pueden estar relacionados, pero no son una única configuración.
Por último, la salida debe converger. Los eventos tardíos provocan múltiples resultados para una ventana. Agregar cada instantánea hace que una suma aguas abajo cuente la ventana repetidamente. Los resultados necesitan un upsert idempotente indexado por la ventana, con un marcador de revisión o finalidad. Un checkpoint protege el estado del operador; si el sink no puede confirmar con el checkpoint, el diseño aún necesita escrituras transaccionales o reemplazo idempotente de versiones.
Preguntas para aclarar antes de responder
- ¿La métrica de negocio se basa en la ocurrencia o en la llegada? Este problema utiliza el
event_timede compra. El tiempo de procesamiento es apropiado únicamente cuando la métrica es específicamente "solicitudes procesadas por el sistema ahora". - ¿El resultado de un minuto es una estimación o un número final? Aquí es una estimación, por lo que un disparador temprano basado en tiempo de procesamiento es válido. Si cada resultado mostrado debe ser completo, el sistema tiene que esperar más o utilizar un cálculo por lotes.
- ¿Pueden cambiar los reportes después de 24 horas? Las 24 horas son el período de corrección automática del trabajo de streaming. Los datos posteriores se envían a conciliación. Un requisito regulatorio o de liquidación para corregir cada registro tardío significa que la limpieza de estado no puede definir la verdad definitiva de los datos.
- ¿Se pueden identificar los duplicados por ID? Aquí existe un
event_idestable. Adivinar a partir del usuario, monto y hora eliminará compras legítimas y omitirá duplicados; corrija primero el contrato del productor. - ¿Qué pasa si un mismo ID contiene contenido diferente? No elija que gane la primera escritura ni la última. Registre un conflicto en la huella digital del payload y póngalo en cuarentena porque el productor violó el contrato de idempotencia.
- ¿El sink acepta instantáneas, deltas o retracciones? Este diseño emite instantáneas de ventana completas y realiza upserts por
(window_start, dimensions). Un sink de solo anexado necesita un changelog versionado cuya capa de lectura seleccione la última revisión. - ¿Se puede confiar en la marca de tiempo del evento? Ponga en cuarentena las marcas de tiempo muy en el futuro, más antiguas que la retención de datos sin procesar o inválidas para la zona horaria acordada. Una marca de tiempo futura incorrecta que participe en el tiempo máximo de evento puede avanzar la marca de agua de forma demasiado agresiva.
Respuesta de 30 segundos
“Asignaría ventanas horarias UTC por event_time, generaría una marca de agua por partición de origen activa y avanzaría con el progreso mínimo seguro. Un disparador por tiempo de procesamiento emite estimaciones cada minuto; la marca de agua emite el resultado puntual, reteniendo el estado durante 24 horas para correcciones. event_id y una huella digital del payload desduplican los eventos, mientras que el sink realiza upserts por ventana y revisión. Los datos excesivamente tardíos van a una salida lateral y a conciliación diaria. La recuperación combina una fuente reproducible, checkpoints y un sink transaccional o idempotente”.
Análisis detallado paso a paso
Comience con el contrato de resultados. Utilice ventanas semiabiertas como [10:00, 11:00) e incluya al menos window_start y las dimensiones de reporte en la clave. Cada salida es la instantánea de ingresos completa para esa ventana en su versión actual:
HourlyRevenue {
window_start
window_end
revenue
revision
result_state // EARLY | ON_TIME | FINAL
}revision aumenta monótonamente para una ventana, y el sink acepta únicamente una revisión mayor. EARLY es la actualización de un minuto y no reclama completitud. ON_TIME significa que la marca de agua ha superado el final de la ventana. FINAL significa que el período de corrección automática de 24 horas ha finalizado. FINAL es un contrato operativo, no una afirmación de que no existan más datos; la conciliación aún maneja los registros más allá del límite de corte.
Derive la política de tiempo a partir de las restricciones. Cada partición de origen activa extrae un event_time validado. Una estrategia común de desorden acotado es:
partition_watermark = max_valid_event_time_seen - out_of_order_bound
operator_watermark = min(active_partition_watermarks)Elija out_of_order_bound a partir de una distribución observada de retraso de llegada, la tasa aceptable de datos tardíos y el objetivo de latencia ON_TIME. Las marcas de agua por partición evitan que una partición rápida declare completada a una lenta. Un operador de múltiples entradas toma el progreso mínimo para no superar una entrada que aún puede producir registros más antiguos. Marque una partición como inactiva solo después de que no haya producido datos durante un intervalo acordado; de lo contrario, puede detener la marca de agua global para siempre. Un umbral de inactividad excesivamente corto hace que los registros antiguos de una partición reanudada sean tardíos, por lo que la salida lateral sigue siendo necesaria.
La ventana tiene dos familias de disparadores. Un disparador temprano basado en tiempo de procesamiento emite la instantánea completa actual una vez por minuto. Cuando la marca de agua supera window_end, emite la revisión ON_TIME. Retenga el estado de la ventana hasta que la marca de agua supere window_end + 24h; cada evento tardío válido actualiza el agregado y produce una revisión superior. En la limpieza, emita FINAL y libere el estado. El modelo de acumulación y disparadores de Beam ilustra por qué "cuándo emitir" y "si un pane es un delta o una instantánea acumulada" son elecciones independientes. Este diseño utiliza instantáneas acumuladas y upserts, de modo que los sistemas aguas abajo no tengan que ensamblar panes.
Desduplique por event_id antes de la agregación. Almacene event_id → (event_time, payload_fingerprint). Permita el paso del primer evento, descarte otra copia con la misma huella digital y envíe el mismo ID con una huella digital diferente a un flujo de conflictos. La retención depende del intervalo más largo en el que el sistema aguas arriba puede reproducir otra copia, no simplemente del tamaño de la ventana. Las semánticas de desduplicación basadas en marcas de agua de Spark exigen de manera similar un umbral de retraso mayor que la brecha temporal entre el duplicado más temprano y el más tardío. Una limpieza de estado demasiado temprana permite que un duplicado tardío se vuelva a contar.
Estime el estado explícitamente. Si el pico de 20,000 eventos por segundo dura 24 horas y cada evento es único, el trabajo recuerda 20,000 × 86,400 = 1,728,000,000 IDs. Incluso una carga útil lógica ilustrativa de solo 40 bytes por entrada hace que este límite superior sea de aproximadamente 69.12 GB. Los objetos del motor, los índices, el backend de estado, los checkpoints y las réplicas aumentan la huella física. Por lo tanto, el estado necesita particionamiento basado en claves y checkpoints incrementales, mientras que la retención sigue el contrato de reproducción. Si el 99% de los duplicados llega en un período más corto, se puede combinar un horizonte de streaming más corto con una clave de negocio en el sink y conciliación fuera de línea, pero el riesgo residual de duplicados debe cuantificarse en lugar de ocultarse tras un TTL.
La ruta de salida debe tolerar la reproducción. Idealmente, el sink participa en transacciones de checkpoint para que la posición de origen, el estado del operador y la salida de la ventana se confirmen juntos. De lo contrario, haga de (window_key, revision) un upsert idempotente: reproducir la misma revisión o una más antigua después de una caída no puede sobrescribir el nuevo valor. La etiqueta de "exactamente una vez" de un broker es insuficiente cuando la escritura en una base de datos externa queda fuera de la confirmación del checkpoint. La semántica de extremo a extremo depende del límite de confirmación compartido entre la fuente, el estado y el sink.
Los eventos posteriores a las 24 horas van a una salida lateral de datos demasiado tardíos y permanecen en el registro sin procesar. Un trabajo por lotes diario recalcula las ventanas afectadas por event_time con las mismas reglas de ID y monto, y luego las compara con las instantáneas FINAL de streaming. Una diferencia material crea una revisión superior o entra en aprobación financiera. El trabajo por lotes reemplaza atómicamente una ventana o realiza un upsert versionado. Sumar el resultado del lote al total existente contaría los mismos datos dos veces cuando el lote se reproduzca.
Utilice una secuencia mínima para verificar la semántica. La ventana [10:00, 11:00) recibe A=100, B=50 y luego una copia idéntica de A. El resultado ON_TIME desduplicado es 150. Después de que la marca de agua supera las 11:00, llega C=20 durante el retraso permitido, y una revisión superior cambia la ventana a 170. D llega después de que la marca de agua supera el punto de limpieza y va a conciliación en lugar de tocar directamente el estado borrado. Si el segundo payload de A tiene un monto de 120, va al flujo de conflictos; el resultado no debe convertirse en 170 ni en 190.
La matriz de pruebas también cubre el desorden entre particiones, una partición inactiva que detiene la marca de agua, la cuarentena de marcas de tiempo futuras, los límites inmediatamente anteriores y posteriores a la marca de agua, caídas antes y después de la escritura en el sink, recuperación de checkpoints y reproducción, revisiones fuera de orden que llegan al sink y ejecuciones repetidas de conciliación. Las métricas de producción incluyen retraso de llegada p50/p95/p99 y de cola, marca de agua actual y retraso por operador, conteos y montos tempranos/puntuales/tardíos/demasiado tardíos, aciertos de desduplicación y conflictos de huella digital, bytes de estado, duración del checkpoint, revisiones rechazadas por el sink y diferencias entre lotes y streaming.
Respuesta sólida de ejemplo
“Definiría el número de un minuto como una estimación y las 24 horas como el período de corrección automática. event_time asigna una compra a una hora UTC; el tiempo de procesamiento solo impulsa el disparador temprano de un minuto. Cada partición de entrada activa genera una marca de agua a partir de su mayor tiempo de evento validado menos el margen de desorden, y las operaciones aguas abajo utilizan la entrada activa más lenta. Las particiones inactivas se marcan explícitamente y las marcas de tiempo futuras inválidas se ponen en cuarentena.
Antes de la agregación, desduplico por event_id, almacenando el tiempo del evento y una huella digital del payload. Un reintento idéntico se cuenta una sola vez; el mismo ID con diferente contenido entra en un flujo de conflictos. La ventana de hora natural emite una instantánea EARLY completa cada minuto, una instantánea ON_TIME una vez que su final queda detrás de la marca de agua, y retiene el estado durante 24 horas. Los eventos tardíos en ese período actualizan el total e incrementan la revisión; la limpieza emite FINAL.
El sink realiza upserts por clave de ventana y revisión en lugar de agregar cada instantánea como un delta. La fuente es reproducible y el estado del operador se restaura desde checkpoints. Si el sink no puede unirse a una transacción de checkpoint, las escrituras versionadas evitan que la reproducción de recuperación reemplace un resultado más nuevo. Los eventos posteriores a las 24 horas van a una salida lateral y un trabajo por lotes diario recalcula el registro sin procesar y lo compara con FINAL.
Si el pico de 20,000 eventos por segundo dura 24 horas y todos los IDs son únicos, el límite superior es de 1,728 millones de entradas de desduplicación. Incluso 40 bytes lógicos por entrada representan 69.12 GB, por lo que las distribuciones medidas de retraso deben justificar la retención, y el costo del estado y de los checkpoints requiere monitoreo. Terminaría con inyección de fallos para duplicados, desorden, tardanza, marcas de tiempo futuras, particiones inactivas y recuperación, asegurando que la ventana final equivalga al recálculo fuera de línea”.
Errores comunes
- Asignar ventanas de negocio por tiempo de procesamiento → el mismo evento puede caer en una hora diferente durante una reproducción histórica → use tiempo de evento para la pertenencia y tiempo de procesamiento solo para salidas tempranas.
- Describir la marca de agua como "no pueden llegar datos más antiguos" → suele ser una estimación heurística de progreso, por lo que aún puede aparecer un evento más antiguo → defina retraso permitido, una salida lateral para datos demasiado tardíos y conciliación.
- Utilizar un único valor sin justificar para el retraso de la marca de agua, el retraso permitido y el TTL de desduplicación → gobiernan la emisión, el estado de la ventana y el reconocimiento de duplicados respectivamente → derívelos del SLO de latencia, el período de corrección y el contrato de reproducción aguas arriba.
- Anexar un total en cada emisión → las emisiones EARLY, ON_TIME y tardías se suman repetidamente → emita una instantánea completa y realice upsert por ventana y revisión, o defina un protocolo de deltas retractables.
- Avanzar a partir de un tiempo de evento máximo global → una partición rápida o una marca de tiempo futura incorrecta hace que los datos de particiones lentas lleguen tarde demasiado pronto → genere progreso por partición, tome el mínimo de las entradas activas y valide marcas de tiempo.
- Mantener una partición permanentemente inactiva en el mínimo → la marca de agua se detiene, impidiendo que se limpien las ventanas y el estado de desduplicación → utilice un manejo observable de inactividad y enrute correctamente los datos tardíos reanudados.
- Almacenar solo el ID del evento → la reutilización de un ID por parte del productor se absorbe silenciosamente → almacene también una huella digital del payload y ponga en cuarentena los conflictos.
- Decir "exactamente una vez está habilitado" → un sink externo puede no compartir el límite del checkpoint → rastree la ruta de confirmación de extremo a extremo y use un sink transaccional o escrituras idempotentes versionadas.
- Descartar datos tras la limpieza del estado → el flujo parece estable mientras los resultados financieros se vuelven inexplicables → retenga los eventos sin procesar e implemente salida lateral, recálculo por lotes y revisión de discrepancias.
Preguntas de seguimiento y respuestas
Pregunta de seguimiento 1: ¿Qué tan grande debe ser el margen de desorden de la marca de agua?
Mida processing_time - event_time en producción y segméntelo por fuente, versión de cliente y región. Primero elija el resultado ON_TIME más tardío aceptable y el porcentaje de eventos permitidos en la ruta de corrección; luego seleccione un percentil que satisfaga ambos. Continúe monitoreando p50, p95, p99 y la cola. Un cambio en la distribución debe gestionarse a través de un lanzamiento de configuración en lugar de permitir que un valor atípico mueva la política automáticamente. La marca de agua gobierna la emisión puntual, mientras que el período de corrección de 24 horas aún cubre una cola más larga.
Pregunta de seguimiento 2: ¿Por qué una partición inactiva de Kafka puede detener la ventana?
El progreso seguro de un operador de múltiples entradas es la marca de agua de entrada mínima. Una partición que no avanza mantiene ese mínimo sin cambios. Tras confirmar que no ha producido eventos durante el umbral de inactividad, márquela como inactiva para que abandone temporalmente el mínimo. No configure el umbral demasiado corto: los registros de una partición reanudada pueden ser más antiguos que la nueva marca de agua. Deben entrar en el procesamiento de retraso permitido o en la salida lateral, y las transiciones de estado de inactividad necesitan su propia métrica.
Pregunta de seguimiento 3: ¿Por qué el estado de desduplicación de 24 horas no garantiza que la duplicación sea imposible para siempre?
Veinticuatro horas solo cubren el período de corrección automática de este problema. Si el sistema aguas arriba reproduce un ID en la hora 25, el estado del streaming ha desaparecido y el evento puede volver a entrar en el agregado. La desduplicación permanente necesita un índice único de negocio de mayor duración, un registro de eventos consultable o una desduplicación completa fuera de línea; cada opción añade costo de almacenamiento y escritura. Declare la garantía como "desduplicado dentro del horizonte de reproducción declarado" y luego concilie los registros posteriores.
Pregunta de seguimiento 4: ¿Qué pasa si la base de datos aguas abajo admite INSERT pero no upsert?
Escriba cada resultado en un changelog inmutable indexado por ventana y revisión. La capa de lectura construye la vista actual a partir de la revisión máxima para cada ventana. Los consumidores deben saber que cada registro es una instantánea completa, no un delta, y no deben sumar todas las revisiones. Si el costo de consulta es demasiado alto, compacte asincrónicamente en una tabla de servicio que admita reemplazo atómico mientras retiene el registro de versiones para auditoría y recuperación.
Pregunta de seguimiento 5: ¿Cómo demuestra que la recuperación no cuenta de más ni de menos?
Para una entrada fija, registre el conjunto de desduplicación y las revisiones de ventana esperados. Termine el trabajo en tres límites: después de la lectura de origen pero antes de un checkpoint de estado, después del checkpoint de estado pero antes de la confirmación del sink, y después de la escritura en el sink pero antes de la confirmación del checkpoint. Reproduzca después de la recuperación y compruebe que cada ID de evento contribuya una sola vez, que el sink conserve solo la revisión más grande y que la ventana final sea igual al recálculo fuera de línea. Un estado RUNNING sin inyección de fallos en estos límites de confirmación no es prueba de corrección de extremo a extremo.