7.6 Datos en tiempo real y en flujo
Presentación y motivación
La mayor parte de lo que sabes sobre los canales de datos asume que los datos se quedan quietos. Recopilas un día de registros, ejecutas un trabajo durante la noche, y lees los resultados por la mañana. Los datos en tiempo real y en flujo invierten esa suposición. En lugar de procesar un montón terminado de datos, procesas un flujo interminable de eventos a medida que llegan, y produces respuestas continuamente. Esta es la diferencia entre el procesamiento por lotes, que opera sobre un conjunto de datos acotado y completo, y el procesamiento de flujo, que opera sobre un flujo no acotado y nunca terminado.
Para los equipos grandes, el flujo aparece en el momento en que la latencia empieza a importarle al negocio. Una decisión de fraude que llega una hora tarde no vale nada. Una señal de personalización que aterriza mañana no personaliza nada. Un tablero operacional que se retrasa un turno respecto a la realidad engaña a quienes lo observan. El capítulo 7.2 (ingeniería de datos) argumenta que deberías elegir el lote por defecto y recurrir al flujo solo donde la latencia genuinamente lo paga, y este capítulo te lleva el resto del camino: cuándo el tiempo real se gana su costo, y cómo construirlo sin incendiar tu presupuesto de operaciones. El flujo se sitúa cerca de los patrones de mensajería impulsada por eventos del capítulo 3.12 (arquitectura y mensajería impulsada por eventos), las elecciones de almacenamiento del capítulo 3.4 (arquitectura de datos y almacenamiento), y las prácticas de telemetría del capítulo 9.2 (observabilidad y telemetría).
Los entornos empresariales y gubernamentales elevan las apuestas. Un banco puntúa las transacciones en busca de fraude en el tiempo que tarda un lector de tarjetas en parpadear. Una agencia de tránsito rastrea vehículos y predice llegadas para millones de pasajeros. Una agencia de beneficios vigila anomalías en las reclamaciones mientras mantiene un registro auditable de cada decisión. En todos estos casos, el valor viene de actuar sobre los datos mientras todavía están frescos, y el riesgo viene de actuar sobre datos que están equivocados, incompletos, o son imposibles de reconstruir después. Este capítulo es contundente sobre ambos.
Principios fundamentales
- Recurre al flujo solo cuando la latencia tenga un valor de negocio claro; el lote es más barato y simple.
- Distingue los datos acotados (finitos) de los datos no acotados (nunca terminados), y diseña en consecuencia.
- Trata el tiempo del evento, no el tiempo de llegada, como la fuente de verdad, y planifica para datos tardíos y fuera de orden.
- Las ventanas y las marcas de agua son cómo obtienes respuestas finitas de flujos infinitos.
- Prefiere resultados efectivamente una vez a través de sumideros idempotentes en lugar de promesas frágiles de exactamente una vez.
- El procesamiento con estado necesita puntos de control para poder recuperarse sin perder ni contar doble.
- Diseña para la contrapresión y el reprocesamiento desde el primer día, no como una ocurrencia tardía.
- Mantén la lógica de flujo observable y auditable; un flujo silencioso es peor que un lote fallido.
Recomendaciones
Justifica el tiempo real antes de construirlo
La decisión de flujo más importante es si transmitir en absoluto. El tiempo real aproximadamente duplica tu complejidad y costo operacional, porque cambias un trabajo que se ejecuta y se detiene por un sistema que debe mantenerse saludable cada segundo. Antes de comprometerte, nombra la decisión que los datos frescos habilitan y el costo de que esa decisión llegue tarde. La puntuación de fraude, las alertas operacionales, y la personalización en vivo usualmente superan el umbral. Un tablero que un humano mira dos veces al día casi nunca lo hace, sin importar cuán satisfactorio suene «tiempo real» en una reunión de planificación. Escribe el requisito de latencia como un número, en segundos o minutos, y contrástalo con la realidad. Mucho de lo que la gente llama tiempo real está bien servido por micro-lotes que se ejecutan cada pocos minutos a una fracción del costo.
Diseña en torno al tiempo del evento, no al tiempo de procesamiento
La idea más difícil en el flujo es que los eventos ocurren en un momento y se procesan en otro. El tiempo del evento es cuándo la cosa realmente ocurrió, por ejemplo cuando un pasajero pasó su tarjeta. El tiempo de procesamiento es cuándo tu sistema se ocupó de manejarlo. Estos se separan constantemente: un teléfono pierde señal en un túnel y sube tres minutos de toques a la vez, un tropiezo de red reordena mensajes, una partición se retrasa. Si calculas sobre el tiempo de procesamiento, tus números se tambalean con tu infraestructura en lugar de reflejar el mundo. Este problema de datos tardíos y fuera de orden es el corazón de la disciplina, y se conecta directamente con el modelado de eventos en la arquitectura impulsada por eventos. Marca cada evento con su tiempo de evento en la fuente, lleva esa marca de tiempo a través de todo el canal, y calcula tus resultados contra ella.
Usa ventanas y marcas de agua para obtener respuestas finitas
Un flujo no acotado nunca termina, así que «cuenta los eventos» no tiene respuesta hasta que lo acotes. Las ventanas hacen ese acotamiento. Las ventanas volcadas cortan el tiempo en cubos fijos y no superpuestos, por ejemplo cada minuto. Las ventanas deslizantes se superponen, así que una ventana de cinco minutos que avanza cada minuto te da una cifra móvil suave. Las ventanas de sesión agrupan ráfagas de actividad separadas por brechas de inactividad, lo cual se ajusta bien a las sesiones de usuario. Una vez que tienes ventanas, necesitas decidir cuándo una ventana está terminada, porque los datos tardíos podrían todavía llegar. Una marca de agua es la estimación del sistema de que probablemente ha visto todos los eventos hasta un tiempo de evento dado. Cuando la marca de agua pasa el final de una ventana, emites el resultado. Ajusta cuánto tiempo esperas: mantén las ventanas abiertas más tiempo y toleras más retraso al costo de latencia y memoria, ciérralas más rápido y arriesgas descartar rezagados. Decide explícitamente qué pasa con los datos que llegan después de que una ventana cierra, si los descartas, los registras, o emites una corrección.
Haz idempotentes los sumideros y prefiere efectivamente una vez
Las garantías de entrega suenan simples y no lo son. La entrega al menos una vez significa que cada evento se procesa, pero algunos pueden procesarse más de una vez después de un reintento, así que los conteos pueden inflarse. Exactamente una vez suena ideal pero es costoso y, tomado literalmente a través de sistemas externos arbitrarios, a menudo imposible. El objetivo práctico es efectivamente una vez: el resultado observable es como si cada evento se hubiera procesado una vez, incluso si la maquinaria reintentó por debajo. Llegas allí haciendo tus sumideros idempotentes, seguros para escribir repetidamente, usando claves deterministas y upserts para que un evento repetido sobrescriba en lugar de duplicar. Combina la entrega al menos una vez con escrituras idempotentes y obtienes resultados correctos sin pagar por una coordinación transaccional pesada en todas partes. Reserva la maquinaria de exactamente una vez verdadera para los lugares estrechos que genuinamente la necesitan.
Pon puntos de control al procesamiento con estado para que pueda recuperarse
Muchos cálculos de flujo útiles tienen estado: conteos corrientes, uniones entre flujos, deduplicación, modelos de fraude que recuerdan el comportamiento reciente. Ese estado vive en memoria y desaparecería cuando un proceso se reinicia. Los puntos de control toman periódicamente una instantánea del estado y la posición del flujo juntos, así que después de un fallo el sistema reanuda desde un punto consistente en lugar de reproducir todo o perder su memoria. Dimensiona tu estado deliberadamente, porque el estado no acotado es una manera común de que un trabajo de flujo se quede sin memoria en producción. Usa expiración y tiempo de vida en el estado que ya no necesitas, y monitorea el tamaño del estado como una métrica de primera clase. El tiempo de recuperación después de un fallo es una preocupación real de nivel de servicio, así que pruébalo antes de que lo hagan tus usuarios.
Transmite desde bases de datos operacionales con captura de datos de cambio
A menudo quieres reaccionar a los cambios en una base de datos que nunca fue diseñada para emitir eventos. La captura de datos de cambio (CDC) resuelve esto leyendo el registro de transacciones de la base de datos y convirtiendo cada inserción, actualización y eliminación en un flujo de eventos de cambio. Esto es mucho mejor que sondear la tabla con un temporizador, lo cual es lento, se pierde estados intermedios, y machaca la fuente. La CDC te permite mantener un índice de búsqueda, una caché, un almacén de analítica, o un servicio posterior continuamente sincronizado con un sistema de registro, y lo hace sin cambios invasivos en la aplicación. Trata el flujo de cambios como un producto de datos de primera clase: versiona su esquema, documenta su significado, y vigila su retraso, porque todo lo posterior hereda ese retraso.
Prefiere una arquitectura de flujo primero sobre mantener dos bases de código
La clásica arquitectura Lambda ejecuta una capa de lotes para un historial preciso y completo junto a una capa de velocidad para resultados frescos y aproximados, luego las fusiona. Funciona, pero te hace escribir y mantener la misma lógica de negocio dos veces, en dos sistemas, y reconciliar las diferencias para siempre. La arquitectura Kappa colapsa esto: mantén un registro duradero y reproducible de eventos y ejecuta todo el procesamiento como procesamiento de flujo, reprocesando el historial al reproducir el registro cuando la lógica cambia. La industria ha derivado hacia esta forma de flujo primero porque una única base de código es dramáticamente más barata de mantener y razonar. Si puedes expresar tus necesidades de lotes como reproducciones sobre un registro de eventos retenido, evitas por completo el impuesto de las dos bases de código. Usa intermediarios basados en registro que retengan el historial para que el reprocesamiento sea cuestión de rebobinar, no de reconstruir.
Expón los flujos como SQL, vistas materializadas y OLAP en tiempo real
No todos los que necesitan flujo deberían tener que escribir código de procesamiento de flujo de bajo nivel. El SQL de flujo permite a los analistas e ingenieros expresar ventanas, uniones y agregaciones en un lenguaje que ya conocen, y mantiene los resultados continuamente actualizados como vistas materializadas. Para consultas analíticas de baja latencia sobre datos frescos, un almacén de procesamiento analítico en línea (OLAP) en tiempo real ingiere el flujo y responde consultas de corte y análisis en milisegundos, que es lo que impulsa un tablero operacional genuinamente en vivo. Empareja estos con las prácticas de analítica de producto del capítulo 7.4 (analítica de producto y experimentación) cuando la meta es retroalimentación rápida sobre funciones y experimentos. Elige estas herramientas de más alto nivel donde encajen, y reserva los procesadores de flujo escritos a mano para la lógica que no pueden expresar.
Planifica la contrapresión y el reprocesamiento desde el inicio
Un flujo puede llegar más rápido de lo que puedes procesarlo. La contrapresión es el mecanismo que permite a un consumidor lento señalar aguas arriba que se reduzca la velocidad en lugar de caerse o descartar datos silenciosamente. Asegúrate de que cada etapa de tu canal la honre, y monitorea el retraso del consumidor como una métrica destacada, porque un retraso creciente es la advertencia más temprana de que estás perdiendo la carrera. El reprocesamiento es la otra capacidad que la gente desearía haber construido de entrada. Cuando encuentras un error o cambias una regla, quieres reproducir el historial a través de la lógica corregida. Eso solo es posible si tu registro de eventos retiene suficiente historial y tus sumideros son lo bastante idempotentes para absorber la reproducción. Diseña ambos desde el primer día; adaptarlos bajo presión de incidente es miserable.
Ventajas y desventajas
| Elección | Ventajas | Desventajas | Mejor ajuste |
|---|---|---|---|
| Lote | Simple, barato, fácil de probar y rellenar retroactivamente | Alta latencia, obsoleto entre ejecuciones | Reporte, la mayoría de la analítica |
| Micro-lote (minutos) | Casi en tiempo real, mucho más simple que el flujo | No verdaderamente instantáneo | Tableros «en tiempo real» |
| Flujo verdadero (menos de un segundo) | Reacción instantánea, resultados continuos | Complejo, costoso, difícil de probar | Fraude, alertas, personalización en vivo |
| Al menos una vez + sumidero idempotente | Resultados correctos, asequible, resiliente | Requiere diseño de claves disciplinado | La mayoría de los canales de flujo |
| Maquinaria de exactamente una vez | Garantía fuerte de extremo a extremo | Costosa, limitada entre sistemas | Caminos estrechos de alto riesgo |
| Lambda (lote + velocidad) | Historial preciso más vista fresca | Dos bases de código que mantener | Migraciones heredadas |
| Kappa (flujo primero) | Una base de código, reproducible | Necesita un registro duradero y retenido | Nuevas plataformas de flujo |
La tensión central es latencia contra complejidad. Cada paso hacia el tiempo real te cuesta en carga operacional, dificultad de prueba, y dinero, y los retornos no son lineales: pasar de diario a cada pocos minutos es barato y a menudo suficiente, mientras que pasar de minutos a menos de un segundo es donde se concentra el gasto. Resuelve la tensión poniendo precio a la decisión, no a la tecnología. Pregunta qué acción habilita la frescura y qué cuesta el retraso, luego compra solo tanta reducción de latencia como esa acción justifique. Cuando sí necesites flujo, apóyate en la entrega al menos una vez con sumideros idempotentes y un registro de flujo primero, porque esa combinación te da corrección y reproducibilidad sin las garantías más pesadas.
Preguntas para discutir con tu equipo
¿Qué decisión habilitan realmente los datos en tiempo real para nosotros, y qué cuesta cuando esos datos llegan un minuto tarde en lugar de instantáneamente? Esta es la pregunta que debería filtrar cada proyecto de flujo, porque el flujo aproximadamente duplica tu costo y complejidad operacional comparado con el lote. Un equipo grande puede quemar trimestres construyendo una plataforma en tiempo real que sirve tableros que un humano revisa dos veces al día, lo cual es dinero prendido en fuego. Trae la acción concreta que impulsan los datos, ya sea bloquear una transacción fraudulenta, avisar a un operador, o cambiar lo que ve un usuario, y pon un número al costo de la latencia para cada una. Si la respuesta honesta es que un micro-lote de cinco minutos serviría la necesidad, ese es un hallazgo que vale la pena celebrar, no ocultar. La respuesta debería cambiar directamente si construyes flujo verdadero, te conformas con micro-lotes, o te quedas en lote.
¿Cómo manejamos los eventos tardíos y fuera de orden, y qué pasa con los datos que llegan después de que una ventana cierra? Los datos tardíos y fuera de orden son la parte difícil del flujo, y los equipos que se saltan esta pregunta la descubren en producción cuando sus números se niegan a reconciliarse. Las presiones en competencia son la latencia y la corrección: mantén las ventanas abiertas más tiempo para atrapar rezagados y retrasas cada resultado y consumes más memoria, ciérralas más rápido y descartas silenciosamente datos reales. Trae evidencia sobre cuán tarde llegan realmente tus datos, medida como la brecha entre el tiempo de evento y el tiempo de procesamiento a través de tus fuentes, ya que una fuente móvil en túneles se comporta muy distinto a un evento del lado del servidor. Decide explícitamente si los datos tardíos se descartan, se registran, o disparan una corrección, y asegúrate de que todos aguas abajo sepan cuál. En un contexto gubernamental donde las cifras deben ser defendibles, descartar silenciosamente eventos tardíos puede ser un problema de cumplimiento, así que la política necesita ser deliberada y documentada.
¿Son nuestros sumideros lo bastante idempotentes para que podamos reproducir el historial con seguridad, y nuestro registro de eventos retiene suficiente para hacer posible la reproducción? El reprocesamiento es la capacidad que los equipos con más frecuencia desearían haber construido y con más frecuencia no construyeron, y depende de dos cosas trabajando juntas: sumideros idempotentes que absorben eventos reproducidos sin duplicar, y un registro duradero que retiene suficiente historial para reproducir desde ahí. Sin ambos, arreglar un error de lógica significa que no puedes recalcular limpiamente el período afectado, y te quedas parchando números a mano bajo presión. Trae tu ventana de retención actual y una prueba concreta: elige un error real del último trimestre y pregunta si podrías haber reproducido la lógica corregida sobre los datos afectados. La resistencia contra esto es el costo, ya que retener el historial y diseñar escrituras idempotentes toma almacenamiento y disciplina por adelantado. Pero la alternativa aparece en el peor momento posible, durante un incidente, así que la respuesta moldea cuánto inviertes en reproducibilidad antes de necesitarla.
Cuando un trabajo de flujo falla, ¿qué tan rápido debe recuperarse, cuánto estado se le permite mantener, y realmente hemos cronometrado una recuperación bajo carga de producción? Un trabajo por lotes que muere puede volver a ejecutarse mañana, pero un flujo siempre activo que muere es una interrupción en progreso, y los trabajos con estado que mantienen conteos corrientes, uniones, o modelos de fraude pueden perder minutos de memoria o tardar mucho en recargar el estado después de un reinicio. Para un equipo grande, aquí es donde un detalle poco glamuroso fija silenciosamente tu disponibilidad real: el estado no acotado crece hasta que un trabajo se queda sin memoria, y una restauración de punto de control lenta convierte un parpadeo de diez segundos en uno de diez minutos. Las presiones en competencia son la frescura contra la seguridad, porque los puntos de control más frecuentes acortan la recuperación pero añaden sobrecarga, y una retención de estado generosa mejora la exactitud pero arriesga el agotamiento de memoria. Trae un objetivo de tiempo de recuperación concreto, tu tamaño de estado actual y su curva de crecimiento, tu intervalo de punto de control, y los resultados de un simulacro de conmutación real en lugar de una estimación esperanzada. En entornos empresariales y gubernamentales donde el flujo respalda la puntuación de fraude o un feed de seguridad pública, un camino de recuperación no probado es un riesgo operacional que has aceptado sin medir, así que trata el simulacro como un requisito, no un extra deseable.
¿Ejecutamos una única base de código de flujo primero o una capa de lote separada y una capa de velocidad separada, y qué nos cuesta realmente mantener las dos reconciliadas? El patrón Lambda de una capa de lotes para historial preciso más una capa de velocidad para resultados frescos te obliga a escribir la misma lógica de negocio dos veces, en dos sistemas, y reconciliar sus respuestas para siempre, mientras que una forma de flujo primero (Kappa) mantiene un registro duradero y reproducible y ejecuta todo el procesamiento como procesamiento de flujo. Para una organización grande, la lógica duplicada es donde se cría la deriva y los números disputados, porque una regla cambia en una capa y no en la otra, y los ingenieros pasan tiempo real explicando por qué las dos discrepan. La atracción hacia mantener ambas es la inercia y la comodidad de una capa de lotes probada, así que pesa eso honestamente contra el impuesto de mantenimiento. Trae la lista de cálculos que actualmente ejecutas en ambos lugares, los incidentes causados por el desacuerdo de las dos capas, y una evaluación de si tu registro de eventos retiene suficiente historial para expresar las necesidades de lotes como reproducciones. En contextos gubernamentales y empresariales auditados, tener dos capas que pueden reportar cifras distintas para el mismo período es en sí mismo un pasivo de cumplimiento, ya que debes poder decir qué número es autoritativo y por qué.
¿Quién opera este sistema siempre activo cuando se rompe a las tres de la madrugada, y hemos presupuestado la carga de guardia y las habilidades especializadas que exige, o estamos asumiendo una dotación de personal con forma de lote? El flujo desplaza el costo de construir a operar: el sistema debe mantenerse saludable cada segundo, lo cual significa cobertura de guardia real, ingenieros fluidos en tiempo de evento, marcas de agua, estado, y semántica de entrega, y pruebas más difíciles que para un trabajo que se ejecuta y se detiene. Los equipos rutinariamente aprueban una plataforma de flujo por la fuerza de sus capacidades y nunca financian a la gente que la mantiene viva, así que la plataforma se degrada y la confianza se erosiona. La contrapartida es alcance contra sostenibilidad: cada canal en tiempo real adicional es otra cosa que puede avisar a alguien, así que la pregunta es si la latencia que compra justifica un compromiso operacional permanente. Trae un inventario honesto de quién es dueño de cada flujo en producción, tu rotación de guardia actual y su margen, y dónde realmente reside la experiencia en tiempo de evento, ya sea una contratación, un socio, o un servicio gestionado. Para un organismo público o una gran empresa, añade los tiempos de espera de contratación y adquisición y cualquier opción de servicio gestionado, porque una plataforma en tiempo real que depende de talento escaso que no puedes reclutar o retener es un plan para operar un sistema propenso a interrupciones con poco personal.
Perspectiva sectorial
Startup. El flujo rara vez es tu primer movimiento, y levantar una plataforma pesada puede hundir a un equipo diminuto. Elige la única señal que toca tu valor central, pon los eventos en un único intermediario basado en registro retenido, y ejecuta un procesador ligero con sumideros idempotentes con clave para que un reintento al menos una vez nunca cuente doble. Mantén unos pocos días de historial para poder reproducir a través de la lógica fija, y prefiere un servicio de flujo gestionado sobre operar tu propio clúster, porque tu recurso más escaso es la atención de ingeniería.
Pequeña empresa. Probablemente no tienes un especialista en flujo y no tienes apetito por operar infraestructura siempre activa, así que trata el tiempo real como algo que compras dentro de herramientas que ya usas en lugar de un sistema que dotas de personal. Enmarca la necesidad como una pregunta de latencia con un número adjunto, y en la mayoría de los casos un micro-lote que se actualiza cada pocos minutos la satisfará a una fracción del costo y riesgo. Elige proveedores cuyas funciones en tiempo real sean transparentes sobre el retraso y fáciles de degradar, y reserva el flujo personalizado para el raro caso donde los datos frescos impulsan directamente los ingresos o la seguridad.
Empresa. El problema es la consistencia y el costo entre muchos equipos: una plataforma compartida basada en registro, una política estándar de tiempo de evento y datos tardíos, y sumideros idempotentes para que los grupos dejen de reinventar canales frágiles. Presupuesta explícitamente las operaciones siempre activas y la carga de guardia, estandariza en un registro de flujo primero para evitar una base de código de lotes duplicada, y gestiona los flujos como productos de datos gobernados con dueños, versionado de esquema, y retraso monitoreado en lugar de una dispersión de trabajos a medida. Rastrea la latencia, el tiempo de recuperación, y el costo por flujo como métricas de cartera.
Gobierno. La auditabilidad y la rendición de cuentas pública moldean cada elección. Retén cada evento procesado en un registro duradero para que las cifras reportadas a los organismos de supervisión, el uso de pasajeros, las anomalías de beneficios, las decisiones de fraude, puedan reconstruirse exactamente, y haz explícita y documentada la política de datos tardíos en lugar de descartar eventos silenciosamente. La contratación pública debería exigir portabilidad de datos y divulgación de las garantías de entrega y retención de un servicio gestionado, y cualquier reformulación después de un cambio de regla debería ser una reproducción defendible a través de la lógica corregida, no un parche manual que nadie puede rastrear.
Ejemplos
Startup. Una aplicación de consumo quiere mostrar a los usuarios un feed de actividad en vivo y marcar inicios de sesión sospechosos a medida que ocurren. El equipo se resiste a levantar una plataforma de flujo pesada. Ponen los eventos en un único intermediario basado en registro retenido, ejecutan un procesador de flujo ligero para la lógica de riesgo de inicio de sesión, y alimentan un almacén OLAP en tiempo real que impulsa el feed de actividad. Cada sumidero tiene clave y es idempotente, así que un reintento al menos una vez nunca cuenta doble. Cuando más tarde encuentran un error en la regla de riesgo, simplemente reproducen el registro a través de la lógica fija durante la noche, porque mantuvieron una semana de historial y nunca necesitaron una segunda base de código de lotes.
Empresa. Un banco minorista puntúa cada transacción de tarjeta en busca de fraude dentro de la ventana de autorización, uniendo el flujo de transacciones en vivo contra un modelo con estado del comportamiento reciente de la cuenta. Los puntos de control permiten que el servicio de puntuación se recupere de un fallo de nodo en segundos sin perder su memoria de los últimos minutos. Por separado, la captura de datos de cambio transmite actualizaciones de la base de datos bancaria central hacia un índice de búsqueda y un servicio de personalización, manteniendo ambos frescos sin sondeo. Los tableros operacionales leen de un almacén OLAP en tiempo real para que los equipos de riesgo y operaciones observen el negocio a medida que se mueve, y todo el canal emite la telemetría de retraso y rendimiento descrita en el capítulo 9.2.
Gobierno. Una autoridad de tránsito metropolitana ingiere posiciones de vehículos y toques de tarifa para predecir llegadas y monitorear la aglomeración en tiempo real, alimentando tanto aplicaciones públicas como un centro de operaciones. Como los pasajeros en túneles suben toques en ráfagas retrasadas, el equipo calcula el uso de pasajeros sobre el tiempo de evento con marcas de agua ajustadas al retraso observado, y registra cualquier evento que llegue después de que su ventana cierra en lugar de descartarlo silenciosamente. Cada evento procesado se retiene en un registro auditable para que las cifras de uso de pasajeros reportadas a los organismos de supervisión puedan reconstruirse exactamente. Cuando cambia una regla de tarifa, reproducen el período afectado a través de la lógica corregida y producen una reformulación defendible.
Caso de negocio: motivaciones, ROI y TCO
El retorno de los datos en tiempo real viene de actuar mientras la acción todavía importa. El fraude atrapado durante la autorización previene una pérdida que un lote nocturno solo reportaría. La personalización que responde dentro de una sesión eleva la conversión de una manera que la recomendación de mañana no puede. El monitoreo operacional que refleja el presente te permite intervenir antes de que un problema pequeño se convierta en una interrupción o un incidente público. En cada caso, el valor es la diferencia entre actuar ahora y actuar después, y esa diferencia es lo que deberías cuantificar cuando presentes el caso.
El costo total de propiedad es más alto que el lote, y la honestidad sobre eso protege tu credibilidad. Pagas por infraestructura siempre activa, por ingenieros que entienden el tiempo de evento, las marcas de agua, el estado, y la semántica de entrega, y por las pruebas más difíciles y la carga de guardia de un sistema que debe mantenerse saludable continuamente en lugar de ejecutarse y detenerse. Una arquitectura de flujo primero sobre un registro retenido reduce el costo continuo al ahorrarte una base de código de lotes duplicada, y elegir al menos una vez con sumideros idempotentes evita el gasto de la maquinaria de exactamente una vez de extremo a extremo. El error más costoso es construir tiempo real donde el micro-lote o el lote bastarían, así que el argumento de costo más fuerte a menudo es una decisión de no transmitir. Enmarca la propuesta al liderazgo en torno a decisiones específicas sensibles a la latencia y su retorno medible, y sé igualmente claro sobre dónde quedarse en lote ahorra dinero sin pérdida de valor.
Antipatrones y trampas
- Construir flujo por prestigio cuando un micro-lote cada pocos minutos satisfaría la necesidad.
- Calcular sobre el tiempo de procesamiento, así que tus números se tambalean con tu infraestructura en lugar del mundo.
- Ignorar los datos tardíos y fuera de orden hasta que la reconciliación falla en producción.
- Perseguir exactamente una vez literal en todas partes en lugar de al menos una vez con sumideros idempotentes.
- Estado no acotado sin expiración, creciendo silenciosamente hasta que un trabajo se queda sin memoria.
- Sin puntos de control, así que un reinicio pierde el estado o fuerza una reproducción completa.
- Sondear bases de datos operacionales con un temporizador en lugar de usar la captura de datos de cambio.
- Mantener una capa de lotes Lambda y una capa de velocidad con lógica duplicada y en deriva.
- Una ventana de retención demasiado corta para reproducir el historial cuando encuentras un error.
- Flujos sin métricas de retraso, rendimiento, o frescura, fallando silenciosamente.
Modelo de madurez
- Nivel 1, Iniciar: Todo es por lotes, o unos pocos trabajos de flujo hechos a mano se ejecutan reactivamente sin monitoreo. Los números se calculan sobre el tiempo de procesamiento, los datos tardíos se ignoran, y un reinicio pierde el estado. Nadie puede reproducir el historial para arreglar un error, y los problemas se descubren cuando las cifras posteriores se niegan a reconciliarse.
- Nivel 2, Desarrollar: Algunos equipos ejecutan canales de flujo centrales en un intermediario basado en registro con puntos de control, y distinguen el tiempo de evento del tiempo de procesamiento y usan ventanas básicas. La práctica es inconsistente de equipo a equipo: la entrega es al menos una vez pero no todos los sumideros son idempotentes, el manejo de datos tardíos es improvisado, y el retraso se vigila informalmente en lugar de alertarse.
- Nivel 3, Estandarizar: El tiempo de evento, las marcas de agua, y una política explícita de datos tardíos están documentados y se aplican en toda la organización. Los sumideros son idempotentes para resultados efectivamente una vez, el estado tiene expiración, y la captura de datos de cambio alimenta sistemas posteriores por convención. Un registro retenido soporta la reproducción, y el retraso, el rendimiento, y la frescura se monitorean con alertas como un estándar en toda la organización en lugar de un hábito por equipo.
- Nivel 4, Gestionar: El patrimonio de flujo se mide y controla contra líneas base. Cada canal lleva objetivos de nivel de servicio para la latencia de extremo a extremo, el retraso del consumidor, el tiempo de recuperación, el sesgo de tiempo de evento, la tasa de eventos tardíos, el tamaño del estado, y el costo por millón de eventos, todos rastreados contra objetivos acordados y alertando sobre la regresión. La recuperación se simula y cronometra en lugar de asumirse, el margen de contrapresión y el crecimiento del estado se vigilan como señales de capacidad, y un nuevo flujo debe superar estas métricas antes de pasar a producción.
- Nivel 5, Orquestar: Una arquitectura de flujo primero sirve tanto las necesidades frescas como las históricas desde un único registro reproducible, y el SQL de flujo, las vistas materializadas, y el OLAP en tiempo real hacen los datos frescos ampliamente accesibles. El reprocesamiento es rutinario y probado, la plataforma se autoescala y reequilibra contra la carga y el costo medidos, y los flujos se retiran, reajustan de alcance, o reemplazan con base en evidencia. El flujo está integrado con la planificación de negocio y riesgo, y cada flujo es observable y auditable de extremo a extremo a medida que cambia la carga y el panorama de costos.
Ideas para el debate
- ¿Dónde en tu pila el «tiempo real» realmente se gana su costo, y dónde es un deseo no examinado?
- ¿Qué tan grande es la brecha entre el tiempo de evento y el tiempo de procesamiento a través de tus fuentes, y la mides?
- ¿Podrías colapsar una configuración Lambda de lote y velocidad en una única base de código de flujo primero, y qué lo bloquearía?
- ¿Cuáles de tus sumideros son verdaderamente idempotentes, y podrías reproducir con seguridad los datos del trimestre pasado a través de la lógica corregida hoy?
- ¿Cuál es tu política para los datos que llegan después de que una ventana cierra, y todos aguas abajo la conocen?
- ¿Cómo cambiaría la captura de datos de cambio la forma en que mantienes sincronizados la búsqueda, las cachés, y la analítica?
Puntos clave
- Recurre al flujo solo cuando una decisión sensible a la latencia lo paga; el lote y el micro-lote son valores predeterminados más baratos.
- Calcula sobre el tiempo de evento, y trata los datos tardíos y fuera de orden como el problema central, manejado con ventanas y marcas de agua.
- Prefiere la entrega al menos una vez con sumideros idempotentes para resultados efectivamente una vez sobre el exactamente una vez literal en todas partes.
- Pon puntos de control al procesamiento con estado, acota tu estado, y monitorea el retraso del consumidor como una métrica destacada.
- Usa la captura de datos de cambio para transmitir desde bases de datos operacionales en lugar de sondear.
- Favorece una arquitectura de flujo primero sobre un registro retenido y reproducible en lugar de mantener dos bases de código.
- Expón los flujos a través de SQL de flujo, vistas materializadas, y OLAP en tiempo real, y mantén cada flujo observable y auditable.
Referencias y lecturas adicionales
- Tyler Akidau, Slava Chernyak, y Reuven Lax, «Streaming Systems».
- Martin Kleppmann, «Designing Data-Intensive Applications».
- Nathan Marz y James Warren, «Big Data» (arquitectura Lambda).
- Jay Kreps, «Questioning the Lambda Architecture» (O’Reilly Radar).
- Fabian Hueske y Vasiliki Kalavri, «Stream Processing with Apache Flink».
- Ben Stopford, «Designing Event-Driven Systems».
- Tyler Akidau y colegas, «The Dataflow Model» (artículo de VLDB sobre ventanas y marcas de agua).