Tema representativo de entrevista

Entrevista de Ingeniería de Datos: ¿Cuándo Deberías Usar Procesamiento por Lotes, Micro-Lotes o Flujo Continuo?

DatosDifícil
Equipo editorial de Offer.ccPublicado Actualizado

Pregunta

Una plataforma de comercio electrónico recibe eventos de pedidos a un pico de 20,000 por segundo, con ráfagas breves de hasta 10 veces esa tasa. Las anomalías de inventario requieren alertas en menos de 10 segundos, un panel de operaciones debe actualizarse en menos de 5 minutos y finanzas cierra los libros al día siguiente con una conciliación ejecutable nuevamente. Alrededor del 2% de los eventos llegan con retraso, hasta por 6 horas. ¿Cómo seleccionarías el procesamiento por lotes, micro-lotes o flujo continuo para cada caso de uso y cómo verificarías la exactitud, reproducción (replay), costo y migración?

Problema y Alcance

Una plataforma de comercio electrónico escribe eventos de pedidos en un registro durable y reproducible (replayable log) y retiene datos sin procesar inmutables. La tasa constante no está especificada; el pico es de 20,000 eventos por segundo, con ráfagas promocionales breves de hasta diez veces esa tasa. Las anomalías de inventario requieren una alerta dentro de los 10 segundos posteriores a la creación del evento. Un panel de operaciones debe actualizarse en menos de 5 minutos. Finanzas cierra los libros al día siguiente y necesita que el mismo día hábil se pueda volver a ejecutar y conciliar. Alrededor del 2% de los eventos llegan tarde, con un retraso de hasta 6 horas.

Dichas tasas, el multiplicador de ráfagas, el porcentaje de retraso y los plazos límite son suposiciones del problema, no afirmaciones de rendimiento sobre un motor. Supón que cada evento tiene un event_id, order_id, event_time y versión de esquema estables. Finanzas utiliza un día hábil en una zona horaria explícita. El equipo de datos de cuatro personas ya opera un almacén de datos (warehouse) y un programador (scheduler) de manera confiable, pero aún no cuenta con operaciones de guardia (on-call) maduras para trabajos de streaming con estado (stateful).

El problema no exige un único modo de procesamiento para los tres consumidores. El objetivo es derivar el diseño mínimo suficiente a partir de los plazos de acción, la completitud de los resultados, la recuperación y el costo de mantenimiento. Esto pertenece a data porque su núcleo radica en la semántica de procesamiento, la frescura y las compensaciones de calidad de datos en lugar de la arquitectura de plataforma entre dominios.

Qué Evalúan los Entrevistadores

La primera señal es si el candidato trabaja en retrospectiva desde la última acción de negocio útil en lugar de optar por streaming por defecto tras ver una cola de mensajes. Una fuente no acotada solo indica que los datos siguen llegando. No exige que cada resultado posterior procese un registro a la vez. El área de finanzas del día siguiente puede leer una instantánea (snapshot) acotada del mismo registro y obtener una reproducibilidad más sencilla mediante el procesamiento por lotes.

La segunda señal es separar la latencia de procesamiento de la latencia de acción. Un motor puede calcular en 500 milisegundos, pero un receptor (sink) que se actualiza cada 5 minutos impide un panel a nivel de segundos. Una respuesta sólida presupuesta la ingesta, el encolamiento, el cómputo, las escrituras, la actualización de caché y la entrega de alertas por separado, midiendo luego cada segmento.

La tercera señal es un límite de exactitud preciso. El procesamiento exactamente una vez (exactly-once) en un motor de streaming no demuestra que los datos tardíos estén completos y no protege automáticamente los efectos secundarios externos. Los trabajos por lotes ejecutables nuevamente tampoco son seguros de forma automática. Sin un rango de entrada definido, clave de negocio, versión de instantánea y publicación atómica, una reejecución puede duplicar o mezclar definiciones.

Finalmente, el entrevistador evalúa el criterio operativo. El streaming continuo necesita capacidad sostenida, recuperación de acumulación (backlog), estado, puntos de control (checkpoints), despliegue seguro y propiedad de guardias. Si los micro-lotes cumplen de manera confiable con un plazo límite de 5 minutos, el estado registro por registro agrega modos de falla sin cambiar una decisión. Por el contrario, el procesamiento por lotes cada hora destruye el valor de una alerta de 10 segundos que efectivamente activa un reabastecimiento o una detención de ventas.

Preguntas a Aclarar Antes de Responder

  • ¿Dónde comienza y termina el plazo límite? Aquí comienza cuando el sistema de origen confirma el evento y termina cuando llega la alerta o el resultado de la consulta es visible. Una promesa que termina en la ingesta del almacén ignora el retraso del receptor y de la caché.
  • ¿La alerta de 10 segundos activa una acción automática? Los cambios automáticos de inventario, las detenciones de ventas o las notificaciones requieren deduplicación, idempotencia y auditoría. Una alerta observacional puede tolerar más falsos positivos o duplicados.
  • ¿El panel de 5 minutos es una estimación o una cifra completa? El micro-lote es suficiente para una estimación revisable con una hora "a la fecha" (as-of). Si cada visualización debe incluir una cola tardía de 6 horas, los requisitos de frescura de 5 minutos y completitud entran en conflicto y el contrato de producto debe cambiar.
  • ¿Cuándo congela finanzas un cierre y puede ajustarlo más adelante? Producir un borrador al día siguiente seguido de asientos de ajuste difiere de esperar seis horas antes de congelar. La respuesta determina el corte del lote y el protocolo de revisión.
  • ¿Pueden los tres usos compartir una capa de datos sin procesar y definiciones de transformación? Pueden compartir eventos normalizados, claves de negocio y conjuntos de pruebas (test fixtures). Copiar de forma independiente filtros, zonas horarias y reglas de montos eventualmente creará una desviación (drift) entre el lote y el flujo.
  • ¿Puede el equipo operar un flujo con estado las 24 horas? Sin alertas de acumulación, simulacros de recuperación de puntos de control y despliegues seguros, el streaming continuo debe restringirse a la ruta más pequeña que realmente necesite una acción a nivel de segundos.

Respuesta de 30 Segundos

“Dividiría a los consumidores según el plazo de acción de negocio en lugar de aplicar un solo modo a la fuente. La ruta de anomalías de inventario de 10 segundos utiliza streaming continuo. El panel de operaciones de 5 minutos utiliza micro-lotes de uno o dos minutos, dejando presupuesto para escrituras y actualización de caché. Finanzas al día siguiente utiliza una instantánea del día hábil por lotes y sigue siendo la fuente de conciliación. Los tres comparten eventos sin procesar inmutables, reglas de normalización y claves de negocio.

Los resultados de streaming son vistas revisables, no afirmaciones de que los datos tardíos estén completos. Las acciones externas son idempotentes por evento y versión de regla. Antes del corte definitivo (cutover), reproduzco el mismo historial a través de las rutas de lote, micro-lote y flujo, comparo en la sombra los recuentos, montos y latencia de cola, y habilito acciones solo para la ruta a nivel de segundos. La regla es elegir el modo más simple que cumpla con el SLO de acción y cuya recuperación converja.”

Análisis Detallado Paso a Paso

Paso 1: Escribir tres contratos de resultados

Define la clave, el plazo límite, la completitud y el comportamiento de revisión de cada salida antes de seleccionar Spark, Flink u otro producto.

UsoClave de resultadoPlazo límite de visibilidadCompletitud y revisiónModo recomendado
Anomalía de inventarioitem, rule version, time window10 segundosrápidamente revisable; acciones idempotentescontinuous stream
Panel de operacionesmetric, dimensions, window5 minutosmuestra hora as-of; datos tardíos reemplazan una versión1–2 minute micro-batch
Cierre financierobusiness day, account, currencyal día siguienterecalcula una instantánea congelada; ajusta discrepanciasbatch

Un solo evento puede servir a tres contratos distintos. La cola de mensajes es un transporte de entrada y no puede forzar a las tres filas al mismo modo de ejecución.

Paso 2: Asignar un presupuesto de acción de extremo a extremo

Un primer presupuesto verificable para la alerta de 10 segundos podría asignar 2 segundos a cada uno de los siguientes: confirmación y transporte de la fuente, encolamiento, cómputo, receptor y acción de regla, y entrega de la alerta. Las pruebas de carga deben reemplazar esas asignaciones provisionales, pero su total no puede exceder el plazo límite del negocio. Registra p95, p99 y la antigüedad máxima de la acumulación para cada segmento. El tiempo del operador por sí solo oculta un receptor lento.

El panel de 5 minutos puede iniciar un micro-lote cada 1 o 2 minutos. Con una fracción de 2 minutos, la espera de programación es de aproximadamente 2 minutos en el peor de los casos, el cómputo y la escritura reciben 1 minuto cada uno, y la actualización de la caché recibe el último minuto. Si una ráfaga de diez veces empuja el cómputo más allá del presupuesto, primero aumenta el paralelismo, acorta la fracción o pospone dimensiones de bajo valor. Un tiempo de ejecución promedio no demuestra el SLO de cola.

Finanzas está limitada por la reproducibilidad y una definición congelada, no por la latencia mínima. El lote lee una instantánea explícita o un rango de desplazamientos (offsets), escribe los resultados en un área de preparación (staging) etiquetada con run_id, los valida y publica de forma atómica. Repetir la misma entrada y versión de regla debería producir el mismo resultado.

Paso 3: Compartir hechos sin crear deuda de implementación dual

Los eventos sin procesar ingresan primero a un registro reproducible o almacén de objetos con la carga útil original, la versión del esquema, el tiempo de ingesta y la posición de origen. Una capa de normalización gestiona la evolución del esquema, las zonas horarias, las unidades de monto, el estado de cancelación y las claves de negocio de manera uniforme. Cada modo de ejecución lee después de ese límite.

Los hechos compartidos no requieren compartir código línea por línea entre tres motores. Un límite más útil es un contrato de datos compartido, fixtures de referencia (golden fixtures) y reglas deterministas. Si el lote y el flujo deben implementar la lógica de ventanas por separado, ejecuta pruebas de equivalencia sobre las mismas fixtures. Finanzas puede conservar una verificación independiente más estricta; la reutilización de código no debe eliminar un control.

Paso 4: Diseñar duplicados, datos tardíos y efectos secundarios por separado

El flujo de inventario deduplica mediante event_id, calcula ventanas cortas de tiempo de evento y emite las versiones de regla y resultado. Una detención de venta, reabastecimiento o notificación utiliza una clave de idempotencia de acción estable. Un punto de control exitoso no demuestra que una llamada HTTP externa ocurrió una sola vez, por lo que los reintentos deben reconocer una acción ya realizada.

El micro-lote de operaciones lee rangos de desplazamiento de origen semiabiertos fijos y reemplaza una ventana de destino versionada en lugar de anexar otro total. Los eventos tardíos ingresan en lotes posteriores y elevan la versión del resultado. El panel muestra una hora "as of", dejando en claro que la frescura de 5 minutos no significa que la cola de 6 horas esté completa.

El lote de finanzas lee una entrada congelada tras el corte acordado. Los eventos posteriores ingresan a una corrida de ajuste que retiene la corrida original, el rango de entrada, la versión de la regla y la aprobación de la discrepancia. El procesamiento de flujo exactamente una vez puede evitar resultados de procesamiento persistidos duplicados en ciertos límites, pero no demuestra la completitud ante la presencia de datos tardíos ni se extiende automáticamente a efectos secundarios externos.

Paso 5: Hacer reproducibles la capacidad y el costo

El pico establecido es de 20,000 eventos por segundo, por lo que una ráfaga de diez veces es de 200,000 por segundo. Si un evento normalizado es ilustrativamente de 1 KiB, el ingreso máximo es de aproximadamente 195 MiB por segundo. Esta es solo una estimación de capacidad; el dimensionamiento para producción debe recalcularse a partir de distribuciones de tamaño comprimidas y descomprimidas.

El streaming continuo necesita capacidad sostenida tanto para el tráfico pico como para la recuperación de acumulación. Diez minutos de falla a 20,000 eventos por segundo generan 12 millones de eventos en cola. Si la capacidad recuperada simplemente es igual a la entrada actual, la acumulación nunca se despeja. Un micro-lote debe terminar una fracción antes de que llegue la siguiente. El procesamiento por lotes puede concentrar recursos en períodos más económicos, pero lotes muy pequeños y frecuentes hacen que la sobrecarga de inicio, confirmación y archivos pequeños sea dominante.

La comparación incluye cómputo continuo, almacenamiento de estado y puntos de control, escrituras en receptores, escaneos, trabajo de guardia y complejidad de despliegue. Los cargos en la nube por sí solos omiten el costo de un equipo de cuatro personas manteniendo tres implementaciones similares. El diseño recomendado limita el streaming continuo a las alertas de inventario, evitando el estado continuo las 24 horas para finanzas y el panel.

Paso 6: Migrar con reproducción en la sombra en lugar de un único corte

Congela un segmento de historial que incluya tráfico normal, ráfagas de diez veces, duplicados, desorden y la cola tardía de 6 horas. Usa el resultado por lotes anterior como línea base mientras los nuevos trabajos de micro-lote y flujo se ejecutan en modo sombra (shadow mode) sin activar acciones de inventario. Compara recuentos, montos, versiones e historiales de revisiones tardías para cada clave de negocio, y explica cada discrepancia.

Aumenta la exposición por etapas: escribe solo en tablas sombra, abre un panel interno y luego permite que la alerta del flujo active acciones reversibles. La inyección de fallas cubre caídas antes y después de los puntos de control, tiempos de espera en receptores, particiones estancadas, cambios de esquema y recuperación de acumulación. Las métricas de aceptación incluyen p99 de extremo a extremo, antigüedad máxima de la acumulación, tiempo de finalización del micro-lote, tasa de revisiones tardías, recuento de acciones duplicadas, diferencias de montos entre lote y flujo, y tiempo de recuperación.

Define también criterios de salida. Si los micro-lotes no cumplen repetidamente con los 5 minutos, inspecciona la programación, el sesgo (skew), los receptores y la capacidad para ráfagas antes de considerar el streaming continuo. Si una alerta de 10 segundos ya no impulsa una acción, degrádala a micro-lote para reducir los costos de guardia. Un modo de procesamiento es una elección de negocio verificable, no una identidad permanente.

Respuesta de Ejemplo Sólida

“Separo una entrada no acotada de sus modos de procesamiento. Estos consumidores tienen diferentes plazos de acción, por lo que no los forzaría a todos a streaming por uniformidad técnica. La anomalía de inventario realmente debe actuar dentro de 10 segundos, por lo que uso un flujo continuo, asigno presupuestos entre origen, cola, cómputo, escritura y entrega, y hago que cada acción sea idempotente por event_id y versión de regla. El panel de operaciones solo necesita 5 minutos, así que comienzo con micro-lotes de uno o dos minutos que reemplazan resultados por ventana y versión y exponen una hora as-of. El cierre financiero lee una instantánea congelada del día hábil por lotes, la prepara por run_id, la valida y luego la publica atómicamente; los datos tardíos ingresan en corridas de ajuste.

Las tres rutas comparten eventos sin procesar reproducibles, un contrato normalizado y fixtures de referencia, mientras que finanzas mantiene un control independiente. A 20,000 eventos por segundo con una breve ráfaga de diez veces, pruebo un ingreso de 200,000 por segundo y la tasa de recuperación tras fallas, no solo el estado constante. Los puntos de control de flujo no reemplazan la idempotencia de acciones externas, y el procesamiento exactamente una vez no demuestra que la cola tardía de 6 horas esté completa.

Para la migración, reproduzco el mismo historial y comparo los resultados a nivel de clave y las revisiones de las rutas de lote, micro-lote y flujo en tablas sombra. El streaming continuo comienza con alertas observacionales. Activa acciones automáticas solo después de que el p99 sea inferior a 10 segundos, las acciones duplicadas sean cero y la acumulación se despeje dentro de su objetivo. Mi regla de decisión es elegir el modo más simple que cumpla con el plazo de acción, la completitud y el contrato de recuperación. Una menor latencia justifica su costo de estado y de guardia únicamente cuando cambia una decisión de negocio.”

Errores Comunes

  • Mover todo a streaming porque Kafka está presente → la entrada continua no hace que cada resultado sea registro por registro, por lo que finanzas y los paneles heredan costos de estado y operativos sin beneficio → elegir según el plazo límite de acción de cada consumidor.
  • Decir únicamente que el streaming es de baja latencia → los receptores, las cachés y la entrega de notificaciones pueden consumir todo el presupuesto → medir la ruta de extremo a extremo y cada cola.
  • Llamar definitiva a una vista de 5 minutos → los eventos aún pueden llegar hasta con 6 horas de retraso → definir semánticas de estimación, revisión, congelación y ajuste.
  • Tratar los puntos de control como acciones externas exactamente una vez → los reintentos pueden repetir llamadas de inventario o notificación → usar claves de idempotencia de acción, un registro de auditoría y pruebas de reproducción.
  • Anexar un nuevo total desde cada micro-lote → la misma ventana se cuenta repetidamente → reemplazar atómicamente por ventana y versión o definir un protocolo de deltas retractables.
  • Copiar reglas de negocio en lote y flujo de forma independiente → las definiciones de zona horaria, cancelación y montos se desvían → compartir contratos y fixtures, comparando luego a nivel de clave.
  • Dimensionar solo para una entrada constante → las ráfagas de diez veces y las acumulaciones por fallas no pueden despejarse dentro del plazo límite → probar juntos el pico, la tasa de recuperación neta y la capacidad del receptor.
  • Habilitar acciones reales en el primer despliegue → una discrepancia semántica altera directamente el inventario → escribir en la sombra y observar alertas antes de habilitar gradualmente acciones reversibles.

Preguntas de Seguimiento y Respuestas

Pregunta de seguimiento 1: ¿Por qué no usar streaming continuo para el panel de 5 minutos?

Si los micro-lotes de uno o dos minutos finalizan de manera confiable bajo condiciones de pico y recuperación, con las escrituras finales y la actualización de caché aún dentro de los 5 minutos, el streaming continuo no cambia una decisión operativa. Agrega estado de larga duración, puntos de control, recuperación de acumulación y costo de despliegue. Si el plazo límite se convierte posteriormente en 30 segundos, o si la programación y el inicio consumen la mayor parte del presupuesto del micro-lote, reproduce el mismo historial para comparar con streaming continuo. La evidencia contra el SLO, no la etiqueta de "tiempo real", es lo que activa la actualización.

Pregunta de seguimiento 2: ¿Podría un solo flujo producir alertas, paneles y resultados financieros?

Podría hacerlo, pero un solo punto de control, cambio de esquema o ventana incorrecta no debe bloquear los tres usos. Una alerta puede leer el flujo normalizado y mantener un estado corto, un panel puede leer agregaciones versionadas y finanzas aún puede recalcular a partir del historial inmutable en una instantánea congelada. Comparte entradas y definiciones mientras aíslas despliegues y dominios de falla. Combina unidades de ejecución únicamente después de demostrar que el impacto de fallas, el reprocesamiento (backfill) y la auditoría siguen siendo aceptables.

Pregunta de seguimiento 3: ¿El procesamiento exactamente una vez hace que la salida de flujo sea adecuada para el cierre financiero?

Esa etiqueta por sí sola es insuficiente. La semántica de procesamiento puede evitar que la salida confirmada se duplique en los reintentos, pero no garantiza que los registros tardíos hayan llegado y no cubre automáticamente los efectos secundarios. Finanzas todavía necesita un rango de entrada congelado, versión de reglas, ejecución repetible, restricciones contables y aprobación de discrepancias. La salida de flujo puede ser una estimación temprana. Solo se convierte en una entrada de cierre si esos controles de auditoría también se cumplen y la conciliación a largo plazo demuestra su equivalencia.

Pregunta de seguimiento 4: ¿Cómo demuestras la recuperación de una ráfaga de diez veces?

Registra la duración de la ráfaga y la duración de la falla, calcula los eventos en cola y mide la tasa neta de drenaje. Si la entrada actual es de 20,000 por segundo y los consumidores procesan 30,000 por segundo, el drenaje neto es de solo 10,000 por segundo; 12 millones de eventos en cola toman alrededor de 20 minutos. La aceptación también observa la latencia de alerta de extremo a extremo, los límites del receptor, el crecimiento del estado y la demora de escalado automático (autoscaling). Una cifra de rendimiento pico sin el cálculo de recuperación neta oculta un incumplimiento en la cola larga.

Pregunta de seguimiento 5: ¿Cuándo debería degradarse el streaming continuo a micro-lotes?

Degrada cuando el negocio ya no actúe sobre la salida a nivel de segundos, el micro-lote cumpla con el nuevo plazo límite, o el costo de guardias y de estado del streaming supere persistentemente la pérdida que previene. Primero ejecuta micro-lotes en la sombra y compara resultados y latencia. Luego desactiva las acciones reales mientras mantienes una entrada reproducible y una ventana de reversión (rollback). Continúa monitoreando las revisiones tardías y el tiempo de finalización en picos después del cambio para que el ahorro de costos no oculte datos obsoletos.

Fuentes públicas

Preguntas relacionadas