Planteamiento y alcance
Esta es una pregunta de diseño común para roles de plataformas de datos, analítica en tiempo real e ingeniería de datos sénior. Asume un pico de 50.000 eventos por segundo, un objetivo de frescura de 60 segundos para el panel y un informe del día anterior que debe entregarse a las 07:00. La entrevista evalúa cuándo un grafo de ejecución compartido obliga a cada consumidor a heredar el SLO más estricto y cómo el progreso independiente y los grupos de recursos reducen el acoplamiento.
Qué evalúa el entrevistador
- Si traduces "tiempo real" y "a tiempo" en SLOs de extremo a extremo antes de nombrar herramientas.
- Si distingues entre ramas dentro de un grafo de Beam, suscripciones independientes a una sola fuente y recursos de cómputo totalmente independientes.
- Si puedes explicar las compensaciones (trade-offs) entre lecturas duplicadas, costo, contrapresión (backpressure), replay, eventos tardíos y dominios de fallo.
- Si proporcionas métricas, un despliegue canary, reversión (rollback) y reconciliación que demuestren que la división mejora los resultados para los usuarios.
Preguntas de clarificación que debes hacer primero
Pregunta si 60 segundos significa frescura desde el tiempo del evento hasta que se puede consultar (event-time-to-queryable freshness) o solo latencia del procesador; si el informe financiero puede rellenar (backfill) particiones tardías; si ambos consumidores pueden compartir un almacenamiento sin procesar inmutable; si los eventos necesitan equidad por inquilino (tenant) o prioridad; qué garantías de duplicados, pérdida y ordenamiento se requieren; y si el presupuesto prioriza el costo o la latencia de cola del panel. Si el informe es T+1 pero el panel tiene un objetivo estricto de baja latencia, el aislamiento ya es un punto de partida creíble para evaluar.
Estructura de respuesta en 30 segundos
Anotaría los SLOs de extremo a extremo de ambos consumidores, las garantías de corrección y los objetivos de recuperación. Si un único pipeline compartido debe cumplir con el estricto objetivo de 60 segundos, establecería una línea base compartida y luego usaría suscripciones independientes para separar el progreso en tiempo real del de procesamiento por lotes (batch). Compartiría el parseo común costoso, pero dividiría el cómputo cuando la presión sobre los recursos o la propagación de fallos sea significativa. Cada ruta obtiene sus propias métricas de frescura, acumulación (backlog), retraso, duplicados y finalización de informes, y el replay junto con la inyección de fallos deciden si el aislamiento adicional justifica el costo.
Respuesta a profundidad
1. Deducir la estructura a partir de los SLOs
Define el SLO de tiempo real como una frescura p99 desde el tiempo del evento hasta la consulta inferior a 60 segundos. Define el SLO de batch como la finalización de los eventos válidos del día anterior antes de las 07:00, permitiendo que los datos tardíos entren en una ventana de reparación acotada. La documentación de Google Dataflow señala que un único pipeline que atiende SLOs mixtos debe cumplir con el objetivo más estricto, lo que permite que el trabajo de menor prioridad consuma capacidad de tiempo real. Diferentes presupuestos de error, rutas de alerta o políticas de escalado automático son razones de peso para separar las rutas.
2. Comparar tres topologías
Una rama dentro del grafo es útil para la decodificación compartida y el enrutamiento ligero. Apache Beam documenta que múltiples transformaciones pueden leer la misma PCollection, pero cada transformación procesa la entrada nuevamente; una sola transformación de múltiples salidas puede procesar cada elemento una vez para el trabajo común.
Las suscripciones independientes a un tema permiten que los consumidores de tiempo real y de batch gestionen confirmaciones (acknowledgements), backlog y posiciones de replay por separado. La documentación de Google Dataflow describe múltiples pipelines que utilizan suscripciones independientes para que cada trabajo extraiga y confirme de manera independiente. Los trabajos totalmente independientes aíslan además la CPU, la memoria, la cadencia de lanzamientos y los dominios de fallo, a costa de lecturas duplicadas, serialización y propiedad operativa.
3. Diseñar la ruta recomendada
Mantén una capa inmutable de eventos sin procesar y crea una suscripción de tiempo real más una suscripción batch desde la misma fuente. El trabajo de tiempo real realiza una agregación ligera hacia un almacenamiento de servicio de baja latencia. El trabajo batch lee los datos retenidos por fecha de evento y escribe tablas particionadas. Si el parseo común representa más del 20% de la CPU total, normalízalo una vez en el ingreso y escribe eventos versionados en la capa sin procesar; no consideres el progreso del consumidor de tiempo real como prueba de que la ruta batch se ha confirmado.
raw-events
-> realtime-subscription -> stream-aggregate -> serving-store
-> batch-subscription or retained-raw -> daily-transform -> partitioned-lake4. Gestionar contrapresión, prioridad y costo
Otorga a la ruta de tiempo real su propio límite de concurrencia y alerta de backlog. Permite que el batch reduzca la concurrencia cuando los recursos de tiempo real estén ajustados, pero no coloques ambos detrás de una única cola no acotada. Si dos cómputos completos son demasiado costosos, comparte la decodificación y el aterrizaje sin procesar, y luego aísla las etapas posteriores (downstream). Cuando el backlog de tiempo real supere los 60 segundos, detén la expansión de la capacidad de batch y recupera primero el SLO de tiempo real. Atribuye bytes de entrada, CPU, antigüedad del backlog y costo por millón de eventos a cada ruta en lugar de comparar la cantidad de trabajos.
5. Gestionar datos tardíos, replay y recuperación
Utiliza una marca de agua (watermark) y una latencia permitida acotada para las ventanas de tiempo real. Enruta los eventos que queden fuera de la ventana a una cola de datos tardíos o a la capa sin procesar, y luego permite que el batch repare las particiones afectadas mediante una escritura versionada e idempotente. Realiza el replay desde una posición guardada de la fuente creando una nueva suscripción; nunca rebobines la posición de confirmación del consumidor de producción. Una ruta batch aún debería poder recuperarse a partir de los datos sin procesar tras un fallo en tiempo real. Si el almacenamiento sin procesar no está disponible, ambas rutas necesitan una degradación explícita y alertas.
6. Decidir mediante un experimento de aceptación
Ejecuta primero la línea base compartida y luego habilita recursos aislados para un pequeño segmento de inquilinos. Compara la frescura p50/p95/p99 del panel, la finalización del informe, la antigüedad del backlog, la tasa de duplicados, el tiempo de replay, la CPU, el almacenamiento y el costo por millón de eventos. Inyecta ráfagas de batch, reinicios del procesador, mensajes duplicados, particiones tardías y suscripciones pausadas. Si la división solo mejora la latencia del procesador mientras incrementa los datos duplicados o excede el presupuesto de costos, mantén el plan compartido y optimiza la etapa común.
Ejemplo de una respuesta sólida
Definiría primero dos SLOs de extremo a extremo: frescura p99 del panel desde el tiempo del evento hasta la consulta por debajo de 60 segundos, y un informe financiero del día anterior completado antes de las 07:00 con una ventana acotada de reparación para datos tardíos. Construiría una línea base compartida, pero nunca compartiría la posición de confirmación entre tiempo real y batch. Mi diseño predeterminado es una capa inmutable de eventos sin procesar, dos suscripciones independientes y dos trabajos downstream. El trabajo de tiempo real realiza una agregación ligera; el trabajo batch lee los datos retenidos por fecha de evento. El parseo común se puede versionar una vez en el ingreso, mientras que los trabajos completos se separan solo cuando los recursos o los dominios de fallo lo requieran. Cada ruta gestiona sus propias métricas de frescura, backlog, datos tardíos, duplicados, finalización de informes y costo unitario; el replay utiliza una nueva suscripción y claves de idempotencia. Inyectaría ráfagas, reinicios y eventos tardíos para verificar los SLOs. Si el aislamiento no mejora los resultados para los usuarios, mantendría el cómputo compartido y optimizaría la etapa común.
Errores comunes
- Síntoma: Copiar dos pipelines completos solo porque hay dos consumidores. Por qué falla: Los costos de parseo común y almacenamiento inicial se duplican sin mejorar el aislamiento. Corrección: Comparte primero los datos inmutables sin procesar y luego aísla el cómputo downstream según el SLO.
- Síntoma: Controlar ambos consumidores con un único offset global. Por qué falla: Un consumidor lento bloquea al rápido y el replay no puede realizarse de forma independiente. Corrección: Utiliza suscripciones independientes o un progreso verificable de manera independiente.
- Síntoma: Medir únicamente la latencia del procesador. Por qué falla: El almacenamiento, la consulta y el tiempo de actualización aún pueden violar el objetivo del usuario. Corrección: Mide la frescura de extremo a extremo y los percentiles.
- Síntoma: Escribir eventos tardíos directamente en el resultado en vivo. Por qué falla: Los conteos pueden duplicarse o los informes ya publicados pueden cambiar de forma silenciosa. Corrección: Utiliza un watermark, una ventana de reparación, versionado y una clave de idempotencia.
- Síntoma: Rebobinar el consumidor de producción para reproducir el historial. Por qué falla: Se interrumpe el progreso en línea y el tráfico puede amplificarse. Corrección: Crea una suscripción de replay con límite de tasa a partir de los datos retenidos.
Preguntas de seguimiento
¿Qué sucede si ambas rutas necesitan el mismo cálculo costoso de características (features)?
Haz que la etapa de características sea versionada y reproducible, almacena su salida una vez y permite que ambas rutas la lean. Acepta el cómputo duplicado solo cuando el estado deba permanecer en línea y no pueda compartirse. Compara una capa intermedia compartida frente al cómputo duplicado mediante experimentos de CPU, latencia y consistencia.
¿Las suscripciones independientes duplicarán el costo de entrada?
Añaden sobrecarga de lectura y confirmación, pero el almacenamiento inmutable, la compresión, las ventanas de retención y el replay bajo demanda pueden controlarlo. Evalúa el costo por millón de eventos junto con el SLO de tiempo real y el aislamiento de fallos; el costo de almacenamiento por sí solo es insuficiente.
¿Puede el procesamiento por lotes tomar prestada capacidad de tiempo real cuando se retrasa?
Utiliza capacidad acotada, apropiable (preemptible) y de baja prioridad con un arrendamiento (lease) y recuperación automática. El tiempo real mantiene un límite estricto y una métrica de backlog independiente, de modo que la capacidad prestada no pueda empujar su p99 más allá de los 60 segundos.
¿Cómo demuestras que ambas rutas finalmente coinciden?
Crea conjuntos de reconciliación a partir de la misma versión de evento, clave de negocio y límite temporal. Compara conteos, montos, registros faltantes, duplicados y correcciones tardías. Si la salida en tiempo real es aproximada, define la ventana de convergencia y el margen de diferencia explicable en lugar de asumir que un total idéntico es una prueba suficiente.