Tema representativo de entrevista

Entrevista de ingeniería de datos: manejo de eventos fuera de orden y tardíos con marcas de agua (watermarks)

DatosDifícil
Equipo editorial de Offer.ccPublicado Actualizado

Pregunta

Diseñe una canalización (pipeline) de streaming que calcule métricas basadas en tiempo de evento cuando los eventos pueden llegar fuera de orden, duplicados o después de que se cierre una ventana. Explique las marcas de agua (watermarks), el retraso (lateness), las correcciones y el monitoreo.

Planteamiento y contexto

Usted calcula métricas de monto y recuento de pedidos en intervalos de cinco minutos. Los clientes fuera de línea, los reintentos y la entrega multirregión hacen que los eventos lleguen tarde, y un mismo pedido puede enviarse más de una vez. El negocio requiere un valor inicial dentro de un minuto y correcciones disponibles para el reporte en un lapso de 24 horas.

Distinga el tiempo de evento (event time) del tiempo de procesamiento (processing time), defina cuándo emite resultados una ventana y muestre a dónde van los eventos tardíos y cómo los consumidores identifican las correcciones.

Qué está evaluando el entrevistador

Semántica de tiempo

Una respuesta sólida define las marcas de tiempo de los eventos, el tiempo de procesamiento, los límites de las ventanas y las zonas horarias antes de explicar que una marca de agua (watermark) es una estimación de que la mayor parte de los datos para una ventana ya ha llegado.

Ciclo de vida del resultado

Separe los resultados preliminares, el resultado a tiempo, las correcciones tardías y los datos posteriores al tiempo límite (cutoff). Una sola emisión no es automáticamente la verdad permanente.

Estado y costo

Analice la retención de estado, el retraso permitido (allowed lateness), el alcance del recálculo, las claves calientes (hot keys) y los puntos de control (checkpoints). Esperar indefinidamente hace que el estado y los costos sean ilimitados.

Observabilidad

Monitoree el retraso de la marca de agua (watermark lag), la distribución del retraso de eventos, la tasa de corrección, los eventos descartados, la proporción de duplicados y la frescura de los resultados.

Preguntas de clarificación que debe hacer

  • ¿Las métricas se agrupan por tiempo de evento o por tiempo de llegada?
  • ¿Cuánto error es aceptable para el primer resultado y cuándo es el corte final?
  • ¿Los eventos tardíos deben corregir reportes que ya fueron publicados?
  • ¿Cada evento incluye un event_id estable para la deduplicación?
  • ¿Cuáles son la tasa pico por clave y el presupuesto máximo de estado?
  • ¿Los eventos que superan el límite de corte deben descartarse, enviarse a cuarentena o recalcularse fuera de línea?

Estructura para una respuesta de 30 segundos

“Utilizaré ventanas de tiempo de evento y requeriré un event_id junto con una marca de tiempo de evento. El procesador de streaming usa una marca de agua (watermark) para emitir un resultado inicial y un período acotado de retraso permitido para correcciones; las correcciones utilizan la misma clave de ventana con una revisión. Los eventos que superan el límite de corte van a cuarentena y a recálculo por lotes en lugar de reescribir silenciosamente el historial. Monitorearé el retraso de la marca de agua, los percentiles de retraso, la tasa de correcciones, los descartes y el tamaño del estado”.

Análisis detallado paso a paso

Paso 1: Definir ventanas y tiempo

Utilice el tiempo de evento en UTC y ventanas fijas como [10:00, 10:05). El tiempo de procesamiento es para alertas operativas y disparadores tempranos, no para agrupaciones de negocio. Incluya inquilino (tenant), producto o región en la clave.

Paso 2: Avanzar la marca de agua (watermark)

Cada partición avanza a partir de las marcas de tiempo de eventos observadas, mientras que una política global toma un límite inferior seguro. Detecte particiones inactivas; una partición silenciosa no debe bloquear toda la canalización.

Paso 3: Disparar y acumular

Emita una aproximación con un disparador basado en tiempo de procesamiento, luego emita el resultado a tiempo cuando la marca de agua supere el final de la ventana. Elija paneles (panes) acumulativos o de descarte e incluya window_end, revision y is_final en cada salida.

Paso 4: Manejar eventos tardíos y duplicados

Dentro del retraso permitido, deduplique por event_id, actualice el estado y emita una nueva revisión. Una retransmisión (replay) no debe sumar el monto dos veces. Retenga el estado hasta el corte comercial y luego límpielo.

Paso 5: Definir la alternativa para datos tardíos

Escriba los eventos que superen el retraso permitido en cuarentena con el motivo y la carga útil original. Un trabajo por lotes recalcula las últimas 24 horas y aplica una actualización/inserción idempotente (upsert) o una revisión más alta en el almacén de correcciones.

Paso 6: Probar y publicar

Utilice marcas de tiempo controladas para probar desorden, duplicados, particiones inactivas, recuperación tras reinicios y retraso en los límites. Los consumidores deduplican (metric, window_end, revision) y utilizan revisiones finales o aprobadas por el corte para los reportes.

Ejemplo de respuesta de alta calidad

“Primero defino el tiempo de evento, las ventanas y los límites de corte como un contrato explícito. Cada evento incluye un eventid estable, eventtime y schema_version. El procesador agrupa en ventanas por inquilino y métrica durante cinco minutos. Las marcas de agua por partición consideran particiones inactivas, y la marca de agua global es una estimación conservadora de la finalización de la ventana.

El sistema emite una revisión preliminar de inmediato, emite la revisión a tiempo después de que la marca de agua supera el final y acepta eventos tardíos durante 30 minutos. Cada salida incluye límites de ventana, una revisión y un indicador de finalización (final flag), por lo que las escrituras posteriores son idempotentes. Los eventos con más de 30 minutos de retraso ingresan a cuarentena; un trabajo de recálculo de 24 horas genera una revisión superior. Después del corte del reporte, auditamos el evento pero no reescribimos silenciosamente el libro contable del negocio.

Monitoreo el retraso de la marca de agua, el retraso p50/p95/p99, las tasas de corrección y descarte, los duplicados, los bytes de estado, el trabajo pendiente de recálculo y la demora del resultado final. La capacidad equivale a las ventanas activas multiplicadas por el estado por ventana, delimitada por puntos de control y TTL de estado”.

Errores comunes

  • Reemplazar el tiempo de evento por el tiempo de procesamiento → los eventos fuera de línea caen en la ventana incorrecta → mantenga event_time y reserve el tiempo de procesamiento para operaciones.
  • Tratar la marca de agua como una verdad absoluta → los eventos tardíos desaparecen silenciosamente → documéntela como una estimación y configure retraso permitido más cuarentena.
  • Emitir un solo resultado inmutable → las correcciones no pueden propagarse → use revisiones y un indicador de finalización para actualizaciones idempotentes.
  • Esperar indefinidamente → el estado y el costo no tienen límite → establezca un corte comercial y recalcule fuera de línea después de este.
  • Deduplicar solo por carga útil → el orden de los reintentos altera el conteo doble → utilice un event_id estable y un estado de deduplicación duradero.
  • Ignorar particiones inactivas → la marca de agua se detiene y las alertas mienten → detecte particiones inactivas y exclúyalas temporalmente del límite inferior.
  • Probar únicamente con entradas ordenadas → los fallos de límites aparecen en producción → inyecte desorden, duplicados, datos tardíos, reinicios y recuperación.

Preguntas de seguimiento y respuestas

Pregunta de seguimiento 1: ¿Por qué puede estancarse una marca de agua?

Una partición puede estar en silencio, desconectada o estimando el progreso de manera conservadora. Combine tiempos de espera por inactividad (idle timeouts), señales de latido (heartbeats) de partición y alertas de retraso de marca de agua para distinguir el silencio de una falla.

Pregunta de seguimiento 2: ¿Cómo se elige el retraso permitido (allowed lateness)?

Utilice el retraso histórico, el corte comercial y el presupuesto de estado. Comience con una estimación del retraso p99 y valídela mediante retransmisión (replay). Un valor más largo no es automáticamente más correcto; aumenta el estado y el costo de corrección.

Pregunta de seguimiento 3: ¿Cómo se previene una tormenta de correcciones?

Procese eventos tardíos en micro-lotes, limite las correcciones por ventana y mantenga solo la revisión más reciente aguas abajo. Durante picos de carga, degrade métricas no críticas a corrección por lotes.

Pregunta de seguimiento 4: ¿Se pueden descartar los eventos posteriores al corte?

Nunca los descarte silenciosamente. Registre datos de cuarentena y auditoría y mida el impacto en el negocio. Los requisitos contables o de cumplimiento normativo pueden exigir un recálculo o manejo manual.

Pregunta de seguimiento 5: ¿Cómo se evita que los resultados retrocedan tras un reinicio?

Guarde en puntos de control el estado de la ventana, el estado de deduplicación y la marca de agua. Utilice revisiones monotónicas, rechace revisiones más antiguas aguas abajo, retransmita el registro (log) y ejecute verificaciones de consistencia tras la recuperación.

Fuente 1: Apache Beam Programming Guide

Beam define marcas de agua, disparadores, retraso permitido y modos de acumulación, incluyendo cómo los datos tardíos pueden producir nuevos paneles (panes).

Fuente 2: Apache Kafka Streams Core Concepts

Kafka Streams documenta el período de gracia para registros fuera de orden y la semántica de descarte después del final de la ventana más la gracia.

Fuente 3: Dataford streaming interview question

La pregunta de entrevista pública considera las marcas de agua, el enrutamiento de eventos tardíos, el recálculo y el monitoreo como los puntos de evaluación profunda para este escenario.

Fuentes públicas

Preguntas relacionadas