, He tenido a Apache Flink en mi lista de "cosas que realmente necesito entender correctamente". Lo había visto mencionar junto con Kafka, lo escuché en conversaciones sobre canalizaciones en tiempo real y en cierto modo entendí el caso de uso. Pero en realidad nunca me senté y lo aprendí adecuadamente.
Si sientes lo mismo, estás en buena compañía. Hay buenas razones para aprender sobre Flink: es una de las herramientas más populares en ingeniería de software en este momento. Netflix lo utiliza para la detección de anomalías casi en tiempo real en su infraestructura de transmisión. Según se informa, Alibaba ejecuta una de las implementaciones de Flink más grandes del mundo: procesa cientos de miles de millones de eventos por día en decenas de miles de máquinas. Uber construyó su plataforma analítica en torno a esto. Flink se ha convertido en la columna vertebral de cómo algunas de las empresas con mayor uso intensivo de datos del mundo procesan la información a medida que sucede. Entonces, si Flink también ha estado en su lista, este es un buen momento para comprenderlo.
Así que me sumergí. Y, sinceramente, me sorprendió, no sólo lo que es Flink, sino también el por qué existe y cómo está construido. La historia de Flink es en realidad la historia de una idea mucho más profunda: la idea de cómo comprender datos de alta escala y en constante transmisión. El planteamiento del problema es en realidad bastante simple: ¿cómo se construyen respuestas prácticas y del mundo real a partir de una escala masiva de datos continuos? Esta publicación es mi intento de explicar esa idea desde cero y mostrarle dónde encaja Flink en ella.
Vamos a sumergirnos.
Antes de comenzar
En esta publicación surgen constantemente dos conceptos sobre los cuales vale la pena asegurarse de que estamos en la misma página antes de continuar.
¿Qué es una corriente? Una secuencia es una secuencia continua y potencialmente interminable de registros que llegan con el tiempo. Piense en un usuario que navega por un sitio web: cada visita a una página, cada clic, cada desplazamiento es un evento que se produce. Uno tras otro, en tiempo real. No existe un “final” natural para esto: mientras el usuario esté activo, los eventos seguirán llegando. Eso es una corriente.
¿Qué es el procesamiento por lotes? El procesamiento por lotes significa tomar una colección finita y limitada de datos y procesarlos todos a la vez. En lugar de reaccionar a cada evento a medida que llega, se recopilan eventos durante un período de tiempo (digamos, una hora) y luego se ejecuta un cálculo sobre todos ellos juntos. El cálculo tiene un comienzo claro y un final claro.
Ambas son formas legítimas de procesar datos. La tensión entre ellos es para lo que Flink fue creado para resolver, y lo lograremos.
De vuelta al problema: cómo producimos realmente los datos
Permítanme concretar esto con un ejemplo que usaremos a lo largo de esta publicación.
Imagine que está creando un motor de recomendaciones, del tipo que muestra a los usuarios "es posible que también les gusten estos" en función de lo que han estado viendo. Para hacer esto bien, su sistema necesita saber cosas como:
¿En qué ha estado haciendo clic este usuario en los últimos minutos?
¿Qué artículos están de moda en este momento entre todos los usuarios?
¿Qué productos vio este usuario pero no compró en la última sesión?
Ahora bien, ¿de dónde vienen esos datos? Cada vez que un usuario abre la página de un producto, registras un evento. Cada clic, cada compra, cada búsqueda: su aplicación escribe continuamente registros que se ven más o menos así:
{ "user_id": "u-8821", "item_id": "p-443", "event_type": "view", "timestamp": "2024–03–10T14:32:01Z" } { "user_id": "u-1042", "item_id": "p-117", "event_type": "compra", "timestamp": "2024–03–10T14:32:03Z" } { "user_id": "u-8821", "item_id": "p-501", "event_type": "clic", "timestamp": "2024–03–10T14:32:07Z" }
Un registro cada pocos segundos para cada usuario, entre millones de usuarios simultáneos, de forma continua. Esos son tus datos. No es un archivo. No es una mesa que se renueva una vez al día. Una secuencia: una secuencia continua e interminable de eventos que describe lo que sus usuarios están haciendo en este momento.
Nuevamente, así es como se ve esta transmisión:
Y, sin embargo, el paradigma dominante durante años fue tomar esa corriente e… ignorar el hecho de que era una corriente. Vuelca los eventos en archivos cada hora. Espere a que se ejecute el trabajo por lotes. Luego, ofrezca recomendaciones basadas en lo que hicieron los usuarios la última hora.
¿Por qué? Porque el procesamiento por lotes es conceptualmente simple. Sabes exactamente qué datos tienes. Puedes razonar claramente sobre el cálculo: comienza, se ejecuta y finaliza. Sistemas como Hadoop y MapReduce (no es necesario que los conozca en profundidad para esta publicación) se construyeron en torno a este modelo y se ampliaron a tamaños de datos enormes. Trabajaron.
Pero hay un costo fundamental: la latencia. Si su trabajo por lotes se ejecuta cada hora, en el peor de los casos, el comportamiento de un usuario en este momento no influirá en sus recomendaciones hasta dentro de una hora. Para un motor de recomendación, eso significa que a un usuario que acaba de mostrar un gran interés en el equipo de senderismo se le muestran accesorios para computadoras portátiles, porque el sistema aún no se ha puesto al día. El usuario buscó una mochila de senderismo y debes mostrarle las tiendas de campaña y los bastones de senderismo en la siguiente página, no una hora después.
Para la detección de fraude, la latencia horaria significa que las transacciones fraudulentas pasan desapercibidas durante una hora. Para un panel en vivo, significa que sus métricas en “tiempo real” pueden estar obsoletas hasta por 59 minutos. El costo del lote es que los eventos ocurren en tiempo real, pero su sistema solo se entera de ellos según un cronograma.
Entonces, a medida que los volúmenes de datos crecieron y los requisitos de latencia se hicieron más estrictos, los ingenieros comenzaron a construir sistemas de transmisión junto con sus sistemas por lotes: sistemas que podían procesar cada evento a medida que llegaba, en milisegundos. Apache Storm fue uno de los primeros líderes aquí. Kinesis amazónica. Samza de LinkedIn.
Pero construir un nuevo sistema de streaming y al mismo tiempo mantener un sistema por lotes existente no es tan sencillo. Ahora tienes dos sistemas que mantener. Su canal de transmisión calculó resultados aproximados en tiempo real. Su canalización por lotes se ejecutó durante la noche y produjo resultados precisos y completos. Había que escribir la misma lógica de negocios dos veces: una para cada sistema, en diferentes marcos, en diferentes idiomas y mantenida sincronizada manualmente. Cuando el trabajo por lotes y el trabajo de transmisión no coincidían en un número (y eventualmente siempre discrepaban), había que averiguar cuál estaba equivocado.
Su motor de recomendaciones en este nuevo mundo ahora se ve así: un componente de transmisión que actualiza las recomendaciones casi en tiempo real en función de eventos recientes y un componente por lotes que reconstruye el modelo de recomendación completo todas las noches en función de datos históricos.
Dos bases de código. Dos canales de implementación. Dos conjuntos de errores. Una capa que sirve tratando de conciliarlas.
La idea clave: el lote es sólo un caso especial de streaming
Aquí está la idea central de Flink, y es bastante simple:
Un conjunto de datos acotado es sólo un caso especial de un flujo de datos ilimitado que termina.
Su base de datos histórica de 5 años de eventos de usuarios: esa es una secuencia que comenzó hace 5 años y se detuvo hoy. Sus archivos de registro del mes pasado: son una secuencia con un principio y un final. La diferencia entre "datos por lotes" y "datos en streaming" no es una distinción fundamental sobre la naturaleza de los datos. Al final del día, son solo eventos JSON de lo que el usuario buscó y en lo que hizo clic. La pregunta es si el río sigue fluyendo o se ha detenido.
Volviendo a nuestro motor de recomendaciones: los "datos históricos" que procesa en su trabajo por lotes nocturno y los "eventos en tiempo real" que procesa en su canal de transmisión son solo registros en la misma secuencia de eventos de usuario. La única diferencia es cuando los lees. El trabajo por lotes nocturno lee registros de hace 6 meses. La canalización de transmisión lee registros de hace 6 segundos. Mismos datos, diferente ventana de tiempo.
Si construye un sistema que procesa flujos de forma nativa (y maneja tanto flujos infinitos como finitos), no necesita sistemas separados. No es necesario mantener dos bases de código. Tienes un motor, un conjunto de lógica y lo apuntas a cualquier porción de la transmisión que necesites.
Eso es lo que Flink intenta hacer.
Entonces, ¿qué es Apache Flink?
Apache Flink es un marco de procesamiento de flujo distribuido. Toma un flujo de datos potencialmente ilimitado (o un lote de datos limitado, lo mismo), lo procesa en paralelo en un grupo de máquinas y produce resultados continuamente a medida que los datos fluyen.
Internamente, los trabajos de Flink se escriben en código y se convierten en un DAG. Por ejemplo, así es como se vería el código para un trabajo de Flink (no es importante comprender todos los detalles, esto es solo para dar una idea aproximada):
// ── 1. FUENTES ─────────────────────── ─────────────────────── búsquedas = readFromKafka("search-events") clics = readFromKafka("click-events") // ── 2. ACTIVIDAD POR USUARIO (agregación de ventanas) ───────────── // agrupa eventos por usuario, calcula las funciones sucesivas durante los últimos 30 minutos userActivity = (searches + clics) .keyBy(userId) .window(slidingWindow(size=30min, slide=1min)) .aggregate(activityAggregator) // → { userId, RecentQueries, RecentClicks, categorías, … } // ── 3. INTEGRACIÓN DE USUARIO (llamar al modelo de torre de usuarios) ─────────────── // convierte las características de la actividad en un vector userState = userActivity.asyncMap(callUserTowerModel) // → { userId, incrustación[128], características } // ── 4. GENERACIÓN DE CANDIDATOS (2 fuentes, luego fusionar) ───────── annCandidates = userState.asyncMap(vectorAnnLookup) // ~500 elementos trendingCandidates = userState.asyncMap(trendingLookup) // ~200 elementos allCandidates = (annCandidates + trendingCandidates) .keyBy(userId) .window(2sec) .reduce(mergeAndDedupe) // → { userId, candidatos: ~1000 itemIds } // ── 5. OBTENER CARACTERÍSTICAS DEL ARTÍCULO (búsqueda por lotes) ───────────────── scoringInputs = allCandidates .joinWith(userState, on=userId) .asyncMap(fetchItemFeatures) // → { userId, userFeatures, [(itemId, itemFeatures) × ~1000] } // ── 6. RANKING (modelo de ranking de llamadas) ───────────────────────── clasificado = scoringInputs.asyncMap(callRankingModel) // → {userId, 100 pares principales (itemId, puntuación) } // ── 7. SINK ──────────────────────── ───────────────────────── clasificado.writeTo(redis)
Internamente, Flink descompone este código en un gráfico de tareas físicas a realizar y divide estas tareas en un conjunto más pequeño de "subtareas" paralelas:
Flink envía tareas a los nodos trabajadores. Cada trabajador ejecuta sus tareas asignadas continuamente, envía latidos periódicos a Flink e informa si una tarea falla para que Flink pueda reiniciarla.
Analicemos los conceptos centrales de Flink.
Conceptos básicos
Corrientes y operadores
Permítanme comenzar con la imagen más simple posible y desarrollarla.
Cada programa de Flink es un gráfico de flujo de datos: un conjunto de operadores conectados por flujos de datos. No se preocupe si esto suena abstracto en este momento: construiremos la imagen pieza por pieza y hará clic rápidamente.
Las fuentes producen datos (lectura de Kafka, un archivo, una base de datos).
Los operadores lo transforman.
Los receptores consumen la salida (escribir en una base de datos, otro tema de Kafka, un panel).
Un operador es una unidad de lógica de procesamiento. Para nuestro motor de recomendaciones, un operador puede filtrar el tráfico de bots, enriquecer un evento con metadatos de productos o contar cuántas veces se vio cada producto. Cada operador recibe registros de uno o más flujos de entrada, les hace algo y emite registros a uno o más flujos de salida.
Una secuencia es la secuencia de registros que fluyen entre operadores. En nuestro caso, un flujo de eventos de usuario: ver eventos, hacer clic en eventos, comprar eventos, uno tras otro a medida que ocurren.
Esta es la forma básica de cualquier trabajo de Flink.
Paralelismo
Una sola máquina puede procesar eventos rápidamente, pero si maneja millones de usuarios, una sola máquina no es suficiente. Flink resuelve esto ejecutando cada operador en paralelo: cada operador se divide en múltiples subtareas que se ejecutan simultáneamente en diferentes máquinas de su clúster.
Si tiene un operador de filtro con paralelismo 4, hay 4 instancias ejecutándose simultáneamente, cada una de las cuales procesa una porción diferente de la secuencia. Agregue más máquinas, obtenga más subtareas, maneje mayores volúmenes. Así es como Flink escala a miles de millones de eventos por día.
Para nuestro motor de recomendaciones, esto significa que la agregación de ventanas para 10 millones de usuarios no se ejecuta secuencialmente en una máquina, sino que se divide entre docenas de trabajadores.
Estado
Volviendo a nuestro motor de recomendaciones: cuando un usuario ve un producto, ese evento por sí solo no dice casi nada. Necesitas contexto. ¿Qué más ha estado viendo este usuario en los últimos minutos? ¿Han estado mirando productos de la misma categoría? ¿Casi compraron algo similar en la última sesión? Para responder a estas preguntas, su sistema necesita memoria: necesita recordar lo que sucedió antes.
En los primeros días del procesamiento de flujos, la mayoría de los sistemas no tenían estado. Cada evento se procesó de forma aislada: el operador vio el evento, lo transformó y siguió adelante. No hay recuerdos de lo que vino antes. Esto funcionó bien para canalizaciones simples: filtrar el tráfico de bots y enriquecer eventos con metadatos de una tabla de búsqueda. Pero era fundamentalmente demasiado limitado para cualquier cosa que requiriera razonamiento sobre patrones a lo largo del tiempo.
Piense en lo que realmente debe hacer nuestro motor de recomendaciones. Para cada evento entrante, debe preguntar: "¿Qué ha hecho el usuario u-8821 en los últimos 10 minutos?" Para responder a esa pregunta, alguien debe mantener una lista actualizada de los eventos recientes del usuario u-8821. Y los eventos recientes del usuario u-1042. Y todos los demás usuarios. Eso es estado: datos que se acumulan y evolucionan a medida que los registros fluyen a través del operador, en lugar de derivarse de cada registro individual.
Flink hace del estado un concepto de primera clase. Un operador puede declarar el estado explícitamente: un contador, un mapa hash codificado por ID de usuario, una lista ordenada de eventos recientes. Flink le brinda ese estado como un objeto administrado que puede leer y escribir durante el procesamiento. Para nuestro motor de recomendaciones, el estado podría ser un mapa hash desde el ID de usuario hasta la "lista de ID de elementos vistos en los últimos 10 minutos". Cada vez que llega un nuevo evento de vista, busca al usuario en el mapa, agrega el elemento y recorta los eventos de más de 10 minutos.
Pero gestionar el estado en un sistema distribuido es realmente difícil. ¿Qué sucede cuando la máquina que maneja su operador falla? Ese mapa hash en memoria desapareció. Flink se encarga de esto: periódicamente toma instantáneas de todo el estado del operador en un almacenamiento duradero, de modo que en la recuperación pueda restaurar todo a donde estaba antes de la falla. Y garantiza que las actualizaciones de estado se apliquen exactamente una vez; incluso si una máquina falla y se repiten los mismos eventos durante la recuperación, sus recuentos no se duplicarán.
Profundizaremos en cómo Flink logra garantías de exactamente una vez en una futura publicación de arquitectura. Por ahora, solo sepa que Flink le brinda un estado que se siente tan confiable como escribir en una base de datos, con el rendimiento de un mapa hash en memoria.
ventanas
Tenemos un flujo de eventos de usuario, operadores que se ejecutan en paralelo y estados acumulados por usuario. He aquí un problema que surge casi de inmediato en cualquier agregación real.
Supongamos que desea calcular "los 10 productos más vistos en los últimos 5 minutos" para impulsar una sección de "tendencias actuales" de su sitio. Tienes un operador que cuenta las vistas por producto. Pero tu flujo es infinito. ¿Cuándo emites un resultado? No puedes esperar hasta que lleguen “todos los eventos”, nunca dejan de llegar.
Necesita una forma de dividir la corriente infinita en partes finitas y calcular cada parte. Esa es una ventana.
Una ventana es una parte limitada de su flujo. Usted lo define, Flink agrupa los eventos en ese fragmento y, cuando el fragmento está "completo", ejecuta su agregación y emite un resultado. Flink tiene varios tipos de ventanas, ventanas giratorias, ventanas deslizantes, ventanas de sesión, etc. No es muy importante comprender las diferencias entre cada tipo de ventana, pero la esencia de las ventanas es que analiza los datos durante un período de tiempo.
Cositas del artículo original
Pasé algún tiempo leyendo el artículo de Apache Flink de 2015: "Apache Flink: procesamiento de secuencias y lotes en un solo motor" de Carbone, Katsifodimos, Ewen, Markl, Haridi y Tzoumas. Algunas cosas del artículo que añaden color útil a lo que cubrimos anteriormente:
Sobre tolerancia a fallos y garantías exactamente una vez
El documento describe la semántica de exactamente una vez de esta manera: "Flink ofrece estrictas garantías de coherencia de procesamiento de exactamente una vez para operadores con estado a través de una combinación de instantáneas distribuidas y reejecución parcial tras la recuperación". La frase clave es reejecución *parcial*: cuando una máquina falla, Flink no reinicia todo el trabajo desde el principio. Revierte a todos los operadores a su última instantánea exitosa y luego reproduce solo la entrada desde ese punto en adelante. La cantidad máxima de reprocesamiento está limitada por la brecha entre dos puntos de control consecutivos, que es un parámetro ajustable.
El mecanismo que hace que esto funcione sin pausar el cálculo se llama Instantánea de barrera asincrónica (ABS), y es realmente inteligente. Lo cubriremos con todo detalle en la próxima publicación. Pero el titular es: Flink inyecta marcadores de “barrera” especiales en el flujo de datos, que fluyen a través de los operadores como registros normales. Cuando un operador recibe una barrera, captura su estado en un almacenamiento duradero y envía la barrera aguas abajo, todo mientras continúa procesando registros. Sin pausas, sin congelaciones, sin eventos perdidos.
Sobre el procesamiento unificado por lotes y flujos
Una de las afirmaciones más claras del artículo es la siguiente: "Un conjunto de datos acotado es un caso especial de un flujo de datos ilimitado". Los autores hacen una afirmación filosófica, no sólo técnica. Y lo respaldan: "Los cálculos por lotes se ejecutan en el mismo tiempo de ejecución que los cálculos en streaming. El ejecutable en tiempo de ejecución puede parametrizarse con flujos de datos bloqueados para dividir grandes cálculos en etapas aisladas que se programan sucesivamente".
En términos sencillos: no existe un motor por lotes independiente en Flink. Los trabajos por lotes se ejecutan exactamente en el mismo tiempo de ejecución de flujo de datos distribuido que procesa sus transmisiones de Kafka. La única diferencia es que los trabajos por lotes utilizan el intercambio de datos "bloqueado" entre etapas: el operador ascendente finaliza por completo antes de que comience el descendente. Todo lo demás (el modelo del operador, la gestión del estado, la serialización) es idéntico.
Volviendo a nuestro motor de recomendaciones: esto significa que el trabajo que cuenta las tendencias de visualización en tiempo real y el trabajo que procesa 6 meses de eventos históricos para el reentrenamiento del modelo pueden compartir los mismos operadores, el mismo clúster y la misma base de código. La promesa del artículo es que la arquitectura Lambda, con sus dos sistemas y dos bases de código, simplemente ya no es necesaria.
Concluyendo
Hagamos rápidamente un TLDR:
Los datos se producen como flujos continuos, pero históricamente los hemos forzado a hacerlos en lotes, lo que genera latencia y el dolor operativo de mantener dos sistemas.
Flink se basa en la idea de que el lote es solo un caso especial de transmisión y unifica ambos en un solo motor.
Los componentes básicos son: operadores (lógica de procesamiento), flujos (datos en movimiento), estado (memoria que persiste en todos los registros) y ventanas (porciones delimitadas de un flujo para cálculo).
Se incorpora tolerancia a fallos con garantías de una sola vez.
Idealmente, me hubiera gustado profundizar en cada uno de estos temas (y hay mucha profundidad en ellos), pero esta publicación ya se ha vuelto bastante larga, así que lo pospondré para el futuro Sanil por ahora. También puedes seguirme en LinkedIn para ver más publicaciones de tamaño byte y saber qué estoy aprendiendo en este momento.
Hablamos mucho sobre Apache Kafka (dado que es la columna vertebral de la mayoría de las arquitecturas de datos), pero ¿alguna vez te preguntaste cómo funciona Apache Kafka y cómo es su arquitectura? Me sorprendió saber lo simple que es realmente Kafka bajo el capó. Escribí una publicación de blog completa al respecto aquí:
Serie de diseño de sistemas: Apache Kafka desde 10,000 pies
¡Veamos qué es Kafka, cómo funciona y cuándo deberíamos usarlo!medium.com
Si está buscando algo más profundo, le recomiendo que consulte una de mis publicaciones más populares en Temporal, una herramienta de orquestación de flujo de trabajo, con explicaciones detalladas sobre cómo se programan, inician y completan los eventos.
Serie de diseño de sistemas: un desglose paso a paso de la arquitectura interna de Temporal
Una inmersión profunda paso a paso en la arquitectura de Temporal, que cubre flujos de trabajo, tareas, fragmentos, particiones y cómo Temporal…medium.com