publication significative du système de traitement d'événements répartis , marqué par le passage à une nouvelle architecture, réalisée en Java, au lieu de l'ancien langage Clojure.
Le projet permet d'organiser un traitement garanti de divers événements en temps réel. Par exemple, Storm peut être utilisé pour l'analyse de flux de données en temps réel, l'exécution de tâches d'apprentissage automatique, l'organisation de calculs continus, la mise en œuvre de RPC, ETL, etc. Le système prend en charge la mise en cluster, la création de configurations à tolérance de panne, le mode de traitement garanti des données et offre des performances élevées, suffisantes pour traiter plus d'un million de requêtes par seconde sur un seul nœud de cluster.
L'intégration avec divers systèmes de traitement de files d'attente et technologies de bases de données est prise en charge. L'architecture de Storm envisage la réception et le traitement de flux de données non structurées constamment mis à jour en utilisant des gestionnaires complexes avec possibilité de partitionnement entre différentes étapes de calcul. Le projet a été remis à la communauté Apache après l'acquisition de BackType par Twitter, qui a initialement développé le framework. Dans la pratique, Storm a été utilisé chez BackType pour analyser le reflet des événements dans les microblogs, en faisant correspondre en temps réel les nouveaux tweets et les liens qu'ils contiennent (par exemple, l'évaluation de la façon dont les liens externes ou les annonces publiées sur Twitter sont relayés par d'autres participants).
La fonctionnalité de Storm est comparable à celle de la plateforme Hadoop, l'une des principales différences étant que les données ne sont pas stockées, mais proviennent de l'extérieur et sont traitées en temps réel. Dans Storm, il n'y a pas de couche intégrée pour organiser le stockage et la requête analytique commence à s'appliquer aux données entrantes jusqu'à ce qu'elle soit annulée (alors que dans Hadoop, des travaux MapReduce se terminant prennent du temps, Storm utilise l'idée de « topologies » s'exécutant en continu). L'exécution des gestionnaires peut être répartie sur plusieurs serveurs - Storm parallélise automatiquement le travail avec les flux sur différents nœuds du cluster.
Au départ, le système était écrit en Clojure et s'exécutait à l'intérieur d'une machine virtuelle JVM. Une initiative a été lancée au sein de la fondation Apache pour porter Storm sur un nouveau noyau écrit en Java, dont les résultats sont proposés dans la version Apache Storm 2.0. Tous les composants de base de la plateforme ont été réécrits en Java. Le support de l'écriture de gestionnaires en Clojure a été maintenu, mais désormais proposé sous forme de liaisons. Storm 2.0.0 nécessite Java 8. Le modèle de traitement multithread a été complètement repensé, permettant une augmentation significative des performances (pour certaines topologies, les latences ont été réduites de 50 à 80 %).
La nouvelle version propose également un nouvel API typé pour les flux, permettant de définir des gestionnaires utilisant des opérations de style programmation fonctionnelle. La nouvelle API est implémentée au-dessus de l'API de base standard et prend en charge la fusion automatique des opérations pour optimiser leur traitement. Dans l'API Windowing pour les opérations de fenêtres, un support de la sauvegarde et de la restauration d'état a été ajouté en backend.
Le planificateur de lancement des gestionnaires a ajouté un support pour tenir compte de ressources supplémentaires lors de la prise de décisions, ne se limitant pas
au CPU et à la mémoire, mais aussi aux paramètres réseau et GPU. Un grand nombre d'améliorations liées à l'intégration avec la plateforme a été apporté. . Le système de contrôle d'accès a été étendu, permettant la création de groupes d'administrateurs et la délégation de tokens. Des améliorations concernant la prise en charge de SQL et des métriques ont été ajoutées. De nouvelles commandes pour le débogage de l'état du cluster sont désormais disponibles dans l'interface d'administration.
Domaines d'application de Storm :
- Traitement en temps réel des flux de nouvelles données ou des mises à jour de base de données ;
- Calculs continus : Storm peut exécuter des requêtes continues et traiter des flux continus, transmettant les résultats au client en temps réel.
- L'appel de procédure à distance distribué (RPC) : Storm peut être utilisé pour assurer le parallélisme de l'exécution de requêtes intensives en ressources. La tâche (« topologie ») dans Storm représente une fonction distribuée sur les nœuds, qui attend la réception de messages à traiter. Après la réception d'un message, la fonction le traite dans un contexte local et renvoie le résultat. Un exemple d'utilisation du RPC distribué peut être le traitement parallèle de requêtes de recherche ou l'exécution d'opérations sur un grand ensemble d'ensembles.
Caractéristiques de Storm :
- Un modèle de programmation simple, qui simplifie considérablement le traitement des données en temps réel ;
- Prise en charge de tous les langages de programmation. Il existe des modules pour les langages Java, Ruby et Python, et l'adaptation à d'autres langages ne pose pas de problème grâce à un protocole de communication très simple, dont la mise en œuvre nécessite environ 100 lignes de code ;
- Résilience : pour exécuter une tâche de traitement des données, il est nécessaire de créer un fichier jar avec le code. Storm diffusera automatiquement ce fichier jar sur les nœuds du cluster, connectera les gestionnaires associés et assurera la surveillance. À la fin de la tâche, le code sera automatiquement désactivé sur tous les nœuds ;
- Scalabilité horizontale. Tous les calculs sont effectués en parallèle, et en cas d'augmentation de la charge, il suffit d'ajouter de nouveaux nœuds au cluster ;
- Fiabilité. Storm garantit que chaque message reçu sera entièrement traité au moins une fois. Un message ne sera traité qu'une seule fois en cas d'absence d'erreurs lors du passage par tous les gestionnaires. En cas de problèmes, les tentatives de traitement échouées seront répétées.
- Vitesse. Le code de Storm est écrit en tenant compte de hautes performances et utilise pour un échange de messages asynchrone rapide le système .
Source : opennet.ru
