Apache Flink es un marco de trabajo unificado de código abierto para el procesamiento de flujos y procesamiento por lotes , desarrollado por la Apache Software Foundation . El núcleo de Apache Flink es un motor de flujo de datos distribuido escrito en Java y Scala . [ 3 ] [ 4 ] Flink ejecuta programas de flujo de datos arbitrarios de forma paralela a los datos y en paralelo (por lo tanto , en paralelo de tareas ). [ 5 ] El sistema de tiempo de ejecución en paralelo de Flink permite la ejecución de programas de procesamiento masivo/por lotes y de flujos. [ 6 ] [ 7 ] Además, el tiempo de ejecución de Flink admite la ejecución de algoritmos iterativos de forma nativa. [ 8 ]
Flink proporciona un motor de transmisión de alto rendimiento y baja latencia [ 9 ] , así como soporte para el procesamiento en tiempo de evento y la gestión de estado. Las aplicaciones de Flink son tolerantes a fallos en caso de fallo de la máquina y admiten semántica de ejecución exactamente una vez [10]. Los programas se pueden escribir en Java, Python [ 11 ] y SQL [ 12 ] y se compilan y optimizan automáticamente [ 13 ] en programas de flujo de datos que se ejecutan en un entorno de clúster o nube [ 14 ] .
Flink no proporciona su propio sistema de almacenamiento de datos, pero proporciona conectores de origen y destino de datos a sistemas como Apache Doris, Amazon Kinesis , Apache Kafka , HDFS , Apache Cassandra y ElasticSearch . [ 15 ]
Desarrollo
Apache Flink se desarrolla bajo la Licencia Apache 2.0 [ 16 ] por la Comunidad Apache Flink dentro de la Fundación de Software Apache . El proyecto está impulsado por 127 [ 17 ] committers y más de 1354 contribuyentes.
Descripción general
El modelo de programación de flujo de datos de Apache Flink proporciona procesamiento evento por evento en conjuntos de datos finitos e infinitos. En un nivel básico, los programas de Flink constan de flujos y transformaciones. “Conceptualmente, un flujo es un flujo (potencialmente infinito) de registros de datos, y una transformación es una operación que toma uno o más flujos como entrada y produce uno o más flujos de salida como resultado.” [ 18 ]
Apache Flink incluye dos API principales: una API DataStream para flujos de datos con o sin límite de tamaño y una API DataSet para conjuntos de datos con límite de tamaño. Flink también ofrece una API Table, un lenguaje de expresiones similar a SQL para el procesamiento relacional de flujos y lotes, que se puede integrar fácilmente en las API DataStream y DataSet de Flink. El lenguaje de nivel superior compatible con Flink es SQL, que es semánticamente similar a la API Table y representa los programas como expresiones de consulta SQL.
Modelo de programación y entorno de ejecución distribuido
Al ejecutarse, los programas Flink se asignan a flujos de datos en tiempo real . [ 18 ] Cada flujo de datos de Flink comienza con una o más fuentes (una entrada de datos, por ejemplo, una cola de mensajes o un sistema de archivos) y termina con uno o más sumideros (una salida de datos, por ejemplo, una cola de mensajes, un sistema de archivos o una base de datos). Se puede realizar un número arbitrario de transformaciones en el flujo. Estos flujos se pueden organizar como un grafo de flujo de datos dirigido y acíclico, lo que permite a una aplicación ramificar y fusionar flujos de datos.
Flink ofrece conectores de origen y destino listos para usar con Apache Kafka , Amazon Kinesis, [ 19 ] HDFS , Apache Cassandra y más. [ 15 ]
Los programas de Flink se ejecutan como un sistema distribuido dentro de un clúster y pueden implementarse en modo independiente, así como en YARN, Mesos, configuraciones basadas en Docker y otros marcos de gestión de recursos. [ 20 ]
Estado: Puntos de control, puntos de guardado y tolerancia a fallos
Apache Flink incluye un mecanismo ligero de tolerancia a fallos basado en puntos de control distribuidos. [ 10 ] Un punto de control es una instantánea automática y asíncrona del estado de una aplicación y su posición en un flujo de origen. En caso de fallo, un programa Flink con puntos de control habilitados reanudará el procesamiento desde el último punto de control completado tras la recuperación, garantizando así que Flink mantenga una semántica de estado de "única vez" dentro de la aplicación. El mecanismo de puntos de control también ofrece puntos de acceso para que el código de la aplicación incluya sistemas externos (como abrir y confirmar transacciones con un sistema de base de datos).
Flink también incluye un mecanismo llamado puntos de guardado, que son puntos de control activados manualmente. [ 21 ] Un usuario puede generar un punto de guardado, detener un programa Flink en ejecución y luego reanudarlo desde el mismo estado de la aplicación y posición en el flujo. Los puntos de guardado permiten actualizar un programa Flink o un clúster Flink sin perder el estado de la aplicación. A partir de Flink 1.2, los puntos de guardado también permiten reiniciar una aplicación con un paralelismo diferente, lo que permite a los usuarios adaptarse a cargas de trabajo cambiantes.
API DataStream
La API DataStream de Flink permite realizar transformaciones (por ejemplo, filtros, agregaciones, funciones de ventana) en flujos de datos limitados o ilimitados. La API DataStream incluye más de 20 tipos diferentes de transformaciones y está disponible en Java y Scala. [ 22 ]
Un ejemplo sencillo de un programa de procesamiento de flujo con estado es una aplicación que emite un recuento de palabras a partir de un flujo de entrada continuo y agrupa los datos en ventanas de 5 segundos:
import org.apache.flink.streaming.api.scala._ import org.apache.flink.streaming.api.windowing.time.Timecase class WordCount ( palabra : String , count : Int )objeto WindowWordCount { def main ( args : Array [ String ]) {val env = StreamExecutionEnvironment . getExecutionEnvironment val text = env . socketTextStream ( "localhost" , 9999 )val counts = text . flatMap { _ . toLowerCase . split ( "\\W+" ) filter { _ . nonEmpty } } . map { WordCount ( _ , 1 ) } . keyBy ( "word" ) . timeWindow ( Time . seconds ( 5 )) . sum ( "count" )recuentos.imprimirenv.execute ( " Window Stream WordCount" ) } }Apache Beam - Flink Runner
Apache Beam “proporciona un modelo de programación unificado avanzado, que permite (a un desarrollador) implementar trabajos de procesamiento de datos por lotes y en tiempo real que pueden ejecutarse en cualquier motor de ejecución”. [ 23 ] El ejecutor Apache Flink-on-Beam es el más completo en cuanto a funcionalidades, según una matriz de capacidades mantenida por la comunidad Beam. [ 24 ]
data Artisans, en conjunto con la comunidad Apache Flink, trabajó estrechamente con la comunidad Beam para desarrollar un ejecutor de Flink. [ 25 ]
API de DataSet
La API DataSet de Flink permite realizar transformaciones (por ejemplo, filtros, mapeo, unión, agrupación) en conjuntos de datos acotados. La API DataSet incluye más de 20 tipos diferentes de transformaciones. [ 26 ] La API está disponible en Java, Scala y una API experimental en Python. La API DataSet de Flink es conceptualmente similar a la API DataStream. Esta API está obsoleta a partir de la versión 2.0 de Flink. [ 27 ]
API de tablas y SQL
La API de tablas de Flink es un lenguaje de expresiones similar a SQL para el procesamiento relacional de flujos y lotes, que se puede integrar en las API de DataSet y DataStream de Flink para Java y Scala. La API de tablas y la interfaz SQL operan sobre una abstracción relacional de tablas. Estas se pueden crear a partir de fuentes de datos externas o de DataStreams y DataSets existentes. La API de tablas admite operadores relacionales como selección, agregación y uniones en tablas.
Las tablas también se pueden consultar con SQL estándar. La API de tablas y SQL ofrecen una funcionalidad equivalente y se pueden combinar en el mismo programa. Cuando una tabla se convierte de nuevo en un DataSet o DataStream, el plan lógico, que se definió mediante operadores relacionales y consultas SQL, se optimiza con Apache Calcite y se transforma en un programa DataSet o DataStream. [ 28 ]
Flink hacia adelante
Flink Forward es una conferencia anual sobre Apache Flink. La primera edición tuvo lugar en 2015 en Berlín. La conferencia, de dos días de duración, contó con más de 250 asistentes de 16 países. Las sesiones se organizaron en dos ejes temáticos: uno con más de 30 presentaciones técnicas de desarrolladores de Flink y otro con formación práctica sobre Flink.
En 2016, 350 participantes asistieron a la conferencia y más de 40 ponentes presentaron charlas técnicas en tres sesiones paralelas. El tercer día, los asistentes fueron invitados a participar en sesiones de formación práctica.
En 2017, el evento también se extendió a San Francisco. La jornada de la conferencia estuvo dedicada a charlas técnicas sobre el uso de Flink en el ámbito empresarial, el funcionamiento interno del sistema, las integraciones con el ecosistema y el futuro de la plataforma. Se ofrecieron ponencias magistrales, charlas de usuarios de Flink del sector empresarial y académico, y sesiones prácticas de formación sobre Apache Flink.
En 2020, a raíz de la pandemia de COVID-19, la edición de primavera de Flink Forward, que debía celebrarse en San Francisco, fue cancelada. En su lugar, la conferencia se celebró virtualmente, del 22 al 24 de abril, e incluyó ponencias magistrales en directo, casos de uso de Flink, aspectos internos de Apache Flink y otros temas sobre procesamiento de flujos de datos y análisis en tiempo real. [ 29 ]
En 2024, Flink Forward [ 30 ] regresó a Berlín, su lugar de origen, para celebrar su décimo aniversario. La conferencia destacó los nuevos planes de Flink 2.0, donde se abandona Java 8 y se introduce un nuevo backend de estado. También se presentó Flink CDC [ 31 ] , que permite crear flujos sin código en YAML. Asimismo, hubo una sesión sobre la adopción de OpenLineage por parte de Flink. [ 32 ]
Historia
En 2010, el proyecto de investigación "Stratosphere: Information Management on the Cloud" [ 33 ] , liderado por Volker Markl (financiado por la Fundación Alemana de Investigación (DFG) ) [ 34 ] , se inició como una colaboración entre la Universidad Técnica de Berlín , la Universidad Humboldt de Berlín y el Instituto Hasso-Plattner de Potsdam. Flink surgió de una bifurcación del motor de ejecución distribuida de Stratosphere y se convirtió en un proyecto de Apache Incubator en marzo de 2014. [ 35 ] En diciembre de 2014, Flink fue aceptado como un proyecto de nivel superior de Apache. [ 36 ] [ 37 ] [ 38 ] [ 39 ]
Fechas de lanzamiento
- 12/2025: Apache Flink 2.2 (12/2025: v2.2.0)
- 07/2025: Apache Flink 2.1 (10/2025: v2.1.1)
- 03/2025: Apache Flink 2.0 (10/2025: v2.0.1)
- 08/2024: Apache Flink 1.20 (09/2025: v1.20.3)
- 03/2024: Apache Flink 1.19 (06/2024: v1.19.1, 02/2025: v1.19.2)
- 10/2023: Apache Flink 1.18 (01/2024: v1.18.1)
- 03/2023: Apache Flink 1.17 (05/2023: v1.17.1; 11/2023: v1.17.2)
- 10/2022: Apache Flink 1.16 (01/2023: v1.16.1; 05/2023: v1.16.2; 11/2023: v1.16.3)
- 05/2022: Apache Flink 1.15 (07/2022: v1.15.1; 08/2022: v1.15.2; 11/2022: v1.15.3; 03/2023: v1.15.4)
- 09/2021: Apache Flink 1.14 (12/2021: v1.14.2; 01/2022: v1.14.3; 03/2022: v1.14.4; 06/2022: v1.14.5; 09/2022: v1.14.6)
- 05/2021: Apache Flink 1.13 (05/2021: v1.13.1; 08/2021: v1.13.2; 10/2021: v1.13.3; 12/2021: v1.13.5; 02/2022: v1.13.6)
- 12/2020: Apache Flink 1.12 (01/2021: v1.12.1; 03/2021: v1.12.2; 04/2021: v1.12.3; 05/2021: v1.12.4; 08/2021: v1.12.5; 12/2021: v1.12.7)
- 07/2020: Apache Flink 1.11 (07/2020: v1.11.1; 09/2020: v1.11.2; 12/2020: v1.11.3; 08/2021: v1.11.4; 12/2021: v1.11.6)
- 02/2020: Apache Flink 1.10 (05/2020: v1.10.1; 08/2020: v1.10.2; 01/2021: v1.10.3)
- 08/2019: Apache Flink 1.9 (10/2019: v1.9.1; 01/2020: v1.9.2)
- 04/2019: Apache Flink 1.8 (07/2019: v1.8.1; 09/2019: v1.8.2; 12/2019: v1.8.3)
- 11/2018: Apache Flink 1.7 (12/2018: v1.7.1; 02/2019: v1.7.2)
- 08/2018: Apache Flink 1.6 (09/2018: v1.6.1; 10/2018: v1.6.2; 12/2018: v1.6.3; 02/2019: v1.6.4)
- 05/2018: Apache Flink 1.5 (07/2018: v1.5.1; 07/2018: v1.5.2; 08/2018: v1.5.3; 09/2018: v1.5.4; 10/2018: v1.5.5; 12/2018: v1.5.6)
- 12/2017: Apache Flink 1.4 (02/2018: v1.4.1; 03/2018: v1.4.2)
- 06/2017: Apache Flink 1.3 (06/2017: v1.3.1; 08/2017: v1.3.2; 03/2018: v1.3.3)
- 02/2017: Apache Flink 1.2 (04/2017: v1.2.1)
- 08/2016: Apache Flink 1.1 (08/2016: v1.1.1; 09/2016: v1.1.2; 10/2016: v1.1.3; 12/2016: v1.1.4; 03/2017: v1.1.5)
- 03/2016: Apache Flink 1.0 (04/2016: v1.0.1; 04/2016: v1.0.2; 05/2016: v1.0.3)
- 11/2015: Apache Flink 0.10 (11/2015: v0.10.1; 02/2016: v0.10.2)
- 06/2015: Apache Flink 0.9 (09/2015: v0.9.1)
- 04/2015: Apache Flink 0.9-milestone-1
Fechas de lanzamiento de Apache Incubator
- 01/2015: Apache Flink 0.8-incubación
- 11/2014: Apache Flink 0.7-incubación
- 08/2014: Apache Flink 0.6-incubating (09/2014: v0.6.1-incubating)
- 05/2014: Stratosphere 0.5 (06/2014: v0.5.1; 07/2014: v0.5.2)
Fechas de lanzamiento de Pre-Apache Stratosphere
- 01/2014: Stratosphere 0.4 (se omitió la versión 0.3)
- 08/2012: Estratosfera 0.2
- 05/2011: Stratosphere 0.1 (08/2011: v0.1.1)
Las versiones 1.14.1, 1.13.4, 1.12.6 y 1.11.5, que supuestamente solo debían contener una actualización de Log4j a la versión 2.15.0, se omitieron porque se descubrió la vulnerabilidad CVE- 2021-45046 durante la publicación de la versión. [ 40 ]
Véase también
Referencias
- ↑ "Versión 2.3.0" . 22 de junio de 2026. Consultado el 23 de junio de 2026 .
- ↑ "Todas las versiones estables de Flink" . flink.apache.org . Apache Software Foundation . Consultado el 20 de diciembre de 2021 .
- ↑ "Apache Flink: Procesamiento de datos por lotes y en tiempo real escalable" . apache.org .
- ↑ "apache/flink" . GitHub . 29 de enero de 2022.
- ↑ Alexander Alexandrov, Rico Bergmann, Stephan Ewen, Johann-Christoph Freytag, Fabian Hueske, Arvid Heise, Odej Kao, Marcus Leich, Ulf Leser, Volker Markl , Felix Naumann, Mathias Peters, Astrid Rheinländer, Matthias J. Sax, Sebastian Schelter, Mareike Höger, Kostas Tzoumas y Daniel Warneke. 2014. La plataforma Stratosphere para análisis de big data . The VLDB Journal 23, 6 (diciembre de 2014), 939-964. DOI
- ↑ Ian Pointer (7 de mayo de 2015). "Apache Flink: El nuevo competidor de Hadoop se enfrenta a Spark" . InfoWorld .
- ↑ "Sobre Apache Flink. Entrevista con Volker Markl" . odbms.org .
- ↑ Stephan Ewen, Kostas Tzoumas, Moritz Kaufmann y Volker Markl . 2012. Generación de flujos de datos iterativos rápidos . Proc. VLDB Endow. 5, 11 (julio de 2012), 1268-1279. DOI
- ↑ "Evaluación comparativa de motores de computación en tiempo real en Yahoo!" . Yahoo Engineering . Consultado el 23 de febrero de 2017 .
- 1 2 Carbone, Paris; Fóra, Gyula; Ewen, Stephan; Haridi, Seif; Tzoumas, Kostas (2015-06-29). "Instantáneas asíncronas ligeras para flujos de datos distribuidos". arXiv : 1506.08603 [ cs.DC ].
- ↑ "Documentación de Apache Flink 1.2.0: Guía de programación en Python" . ci.apache.org . Consultado el 23 de febrero de 2017 .
- ↑ "Documentación de Apache Flink 1.2.0: Tabla y SQL" . ci.apache.org . Consultado el 23 de febrero de 2017 .
- ↑ Fabian Hueske, Mathias Peters, Matthias J. Sax, Astrid Rheinländer, Rico Bergmann, Aljoscha Krettek y Kostas Tzoumas. 2012. Abriendo las cajas negras en la optimización del flujo de datos . Proc. VLDB Endow. 5, 11 (julio de 2012), 1256-1267. DOI
- ↑ Daniel Warneke y Odej Kao. 2009. Nephele: procesamiento de datos paralelo eficiente en la nube . En Actas del 2.º Taller sobre Computación Multitarea en Redes y Supercomputadoras (MTAGS '09). ACM, Nueva York, NY, EE. UU., Artículo 8, 10 páginas. DOI
- 1 2 "Documentación de Apache Flink 1.2.0: Conectores de transmisión" . ci.apache.org . Consultado el 23 de febrero de 2017 .
- ↑ "ASF Git Repos - flink.git/blob - LICENSE" . apache.org . Archivado del original el 23/10/2017 . Consultado el 12/04/2015 .
- ↑ "Información del Comité Apache Flink" . Consultado el 10 de abril de 2025 .
- 1 2 "Documentación de Apache Flink 1.2.0: Modelo de programación de flujo de datos" . ci.apache.org . Consultado el 23 de febrero de 2017 .
- ↑ "Kinesis Data Streams: procesamiento de datos en tiempo real" . 5 de enero de 2022.
- ↑ "Documentación de Apache Flink 1.2.0: Entorno de ejecución distribuido" . ci.apache.org . Consultado el 24 de febrero de 2017 .
- ↑ "Documentación de Apache Flink 1.2.0: Entorno de ejecución distribuido - Puntos de guardado" . ci.apache.org . Consultado el 24 de febrero de 2017 .
- ^ "Documentación de Apache Flink 1.2.0: Guía de programación de la API de Flink DataStream" . ci.apache.org . Consultado el 24 de febrero de 2017 .
- ↑ "Apache Beam" . beam.apache.org . Consultado el 24 de febrero de 2017 .
- ↑ "Matriz de capacidades de Apache Beam" . beam.apache.org . Consultado el 24 de febrero de 2017 .
- ↑ "¿Por qué Apache Beam? Una perspectiva de Google | Blog de Google Cloud sobre Big Data y aprendizaje automático | Google Cloud Platform" . Google Cloud Platform . Archivado del original el 25 de febrero de 2017. Consultado el 24 de febrero de 2017 .
- ^ "Documentación de Apache Flink 1.2.0: Guía de programación de la API de Flink DataSet" . ci.apache.org . Consultado el 24 de febrero de 2017 .
- ↑ "API de Flink versión 2.0 obsoletas" . Consultado el 10 de abril de 2025 .
- ↑ "Procesamiento de flujos para todos con SQL y Apache Flink" . flink.apache.org . 24 de mayo de 2016. Consultado el 8 de enero de 2020 .
- ↑ "Conferencia virtual Flink Forward 2020" .
- ↑ "Flink Forward" . Consultado el 9 de abril de 2025 .
- ↑ "Captura de datos de cambios" . Consultado el 9 de abril de 2025 .
- ↑ «OpenLineage» . Consultado el 9 de abril de 2025 .
- ↑ "Estratosfera" . estratosfera.eu .
- ↑ «Stratosphere - Gestión de la información en la Nube» . Deutsche Forschungsgemeinschaft (DFG) . Consultado el 1 de diciembre de 2023 .
- ↑ "Estratosfera" . apache.org .
- ↑ "Detalles del proyecto para Apache Flink" . apache.org .
- ↑ "La Fundación de Software Apache anuncia Apache™ Flink™ como un proyecto de primer nivel : Blog de la Fundación de Software Apache" . apache.org . 12 de enero de 2015.
- ↑ "¿Encontrará el misterioso Apache Flink su nicho de mercado ideal?" . siliconangle.com . 9 de febrero de 2015.
- ↑ (en alemán)
- ↑ "Versiones de emergencia de Apache Flink Log4j" . flink.apache.org . Apache Software Foundation. 16 de diciembre de 2021. Consultado el 22 de diciembre de 2021 .
Enlaces externos
- Proyectos de la Apache Software Foundation
- Procesamiento distribuido de flujos
- Software libre programado en Java.
- Software de sistema gratuito
- Software que utiliza la licencia Apache.