una versión significativa del sistema de procesamiento distribuido de eventos , destacándose por la transición a una nueva arquitectura implementada en Java, en lugar del lenguaje Clojure usado anteriormente.
El proyecto permite organizar el procesamiento garantizado de diversos eventos en tiempo real. Por ejemplo, Storm se puede utilizar para analizar flujos de datos en tiempo real, realizar tareas de aprendizaje automático, organizar cálculos continuos, implementar RPC, ETL, entre otros. El sistema soporta clustering, creación de configuraciones redundantes, modo de procesamiento de datos garantizado y tiene un alto rendimiento, suficiente para procesar más de un millón de solicitudes por segundo en un nodo de clúster.
Se soporta la integración con varios sistemas de procesamiento de colas y tecnologías de bases de datos. La arquitectura de Storm implica la recepción y procesamiento de flujos de datos no estructurados que se actualizan constantemente mediante el uso de procesadores complejos arbitrarios con la posibilidad de particionar entre diferentes etapas de los cálculos. El proyecto fue transferido a la comunidad Apache tras la adquisición de Twitter de la empresa BackType, que originalmente desarrolló el framework. En la práctica, Storm se utilizó en BackType para analizar la representación de eventos en microblogs, emparejando en tiempo real nuevos tweets con los enlaces que contenían (por ejemplo, se evaluó cómo los enlaces externos o los anuncios publicados en Twitter son retransmitidos por otros participantes).
La funcionalidad de Storm se compara con la plataforma Hadoop, siendo la principal diferencia que los datos no están almacenados en un repositorio, sino que llegan desde el exterior y se procesan en tiempo real. En Storm no hay una capa integrada para organizar el almacenamiento y la consulta analítica comienza a aplicarse a los datos entrantes hasta que se cancela (mientras que en Hadoop se utilizan trabajos MapReduce que ocupan tiempo finito, en Storm se aplica la idea de ‘topologías’ que se ejecutan continuamente). La ejecución de los procesadores puede distribuirse entre varios servidores: Storm paraleliza automáticamente el trabajo con flujos en diferentes nodos del clúster.
El sistema fue inicialmente escrito en el lenguaje Clojure y se ejecuta dentro de la máquina virtual JVM. La fundación Apache inició la iniciativa de traducir Storm a un nuevo núcleo, escrito en Java, cuyos resultados se presentaron en la versión Apache Storm 2.0. Todos los componentes básicos de la plataforma han sido reescritos en Java. Se ha mantenido el soporte para escribir controladores en Clojure, pero ahora se ofrece en forma de enlaces. Para utilizar Storm 2.0.0, se requiere Java 8. Se ha rediseñado completamente el modelo de procesamiento multihilo, lo que ha permitido un aumento significativo en el rendimiento (para algunas topologías, las latencias se redujeron entre un 50% y un 80%).
En la nueva versión, también se ha introducido una nueva API tipada de Streams, que permite definir controladores utilizando operaciones en estilo de programación funcional. La nueva API se implementa sobre la API base estándar y soporta la combinación automática de operaciones para optimizar su procesamiento. En la API de Windowing para operaciones con ventanas, se ha añadido soporte para el almacenamiento y recuperación del estado en el backend.
Se ha añadido al planificador de ejecución de controladores soporte para tener en cuenta recursos adicionales al tomar decisiones, que no se limitan a
CPU y memoria, tales como parámetros de red y GPU. Se han realizado numerosas mejoras relacionadas con la integración con la plataforma . Se ha ampliado el sistema de control de acceso, que ahora permite la creación de grupos de administradores y la delegación de tokens. Se han añadido mejoras relacionadas con el soporte de SQL y métricas. En la interfaz de administrador, se han presentado nuevos comandos para la depuración del estado del clúster.
Áreas de aplicación de Storm:
- Procesamiento de flujos de nuevos datos o actualizaciones de bases de datos en tiempo real;
- Cálculos continuos: Storm puede ejecutar consultas continuas y procesar flujos continuos, enviando los resultados del procesamiento al cliente en tiempo real.
- Llamadas a procedimientos remotos distribuidos (RPC): Storm se puede utilizar para garantizar la paralelización de la ejecución de solicitudes que consumen recursos. La tarea (“topología”) en Storm consiste en una función distribuida entre nodos que espera la llegada de mensajes que deben ser procesados. Tras recibir un mensaje, la función lo procesa en un contexto local y devuelve el resultado. Un ejemplo del uso de RPC distribuido puede ser el procesamiento paralelo de consultas de búsqueda o la realización de operaciones sobre un gran conjunto de conjuntos.
Características de Storm:
- Modelo de programación simple que simplifica significativamente el procesamiento de datos en tiempo real;
- Soporte para cualquier lenguaje de programación. Existen módulos para los lenguajes Java, Ruby y Python; la adaptación para otros lenguajes no presenta dificultades gracias a un protocolo de comunicación muy simple, cuyo soporte se puede implementar con aproximadamente 100 líneas de código;
- Tolerancia a fallos: para ejecutar una tarea de procesamiento de datos, es necesario generar un archivo jar con el código. Storm distribuirá automáticamente este archivo jar entre los nodos del clúster, conectará los controladores relacionados y organizará la supervisión. Al finalizar la tarea, el código se desconectará automáticamente en todos los nodos;
- Escalabilidad horizontal. Todos los cálculos se realizan en paralelo; ante un aumento de la carga, basta con conectar nuevos nodos al clúster;
- Fiabilidad. Storm garantiza que cada mensaje recibido será procesado completamente al menos una vez. Un mensaje solo se procesará una vez si no hay errores durante el paso por todos los controladores; si surgen problemas, los intentos fallidos de procesamiento serán repetidos.
- Velocidad. El código de Storm está escrito con un enfoque en un alto rendimiento y utiliza para el rápido intercambio de mensajes un sistema .
Fuente: opennet.ru
