Automatisation de la livraison des flux dans Apache NiFi

Bonjour à tous !

Automatisation de la livraison des flux dans Apache NiFi

La tâche consiste à déployer le flux, illustré dans l'image ci-dessus, sur N serveurs avec Apache NiFi. Le flux est un test - il génère un fichier et l'envoie vers une autre instance de NiFi. Le transfert de données se fait via le protocole NiFi Site to Site.

NiFi Site to Site (S2S) est un moyen sûr et facilement configurable de transférer des données entre les instances de NiFi. Pour savoir comment fonctionne S2S, consultez documentation et il est important de ne pas oublier de configurer l'instance NiFi pour autoriser S2S, voir ici.

Dans les cas où il s'agit de transfert de données par S2S, une instance est appelée client, l'autre serveur. Le client envoie des données, le serveur les reçoit. Il y a deux façons de configurer le transfert de données entre eux :

  1. Push. Les données sont envoyées depuis l'instance cliente via un Remote Process Group (RPG). Sur l'instance serveur, les données sont reçues via un Input Port
  2. Tirer. Le serveur reçoit les données via l'RPG, le client les envoie via un Output port.


Le flux à déployer est stocké dans Apache Registry.

Apache NiFi Registry est un sous-projet d'Apache NiFi, représentant un outil de stockage des flux et de gestion des versions. Une sorte de GIT. Des informations sur l'installation, la configuration et l'utilisation du registre peuvent être trouvées dans la documentation officielle. Les flux à stocker sont regroupés dans un process group et sont ainsi conservés dans le registre. Nous y reviendrons plus tard dans cet article.

Au départ, lorsque N est un petit nombre, les flux sont livrés et actualisés manuellement dans un délai acceptable.

Mais avec l'augmentation de N, les problèmes se multiplient :

  1. le temps nécessaire pour actualiser les flux augmente. Il faut se connecter à tous les serveurs
  2. des erreurs d'actualisation des modèles surviennent. Ici, cela a été mis à jour, mais là, cela a été oublié
  3. les erreurs humaines lors de l’exécution d’un grand nombre d’opérations similaires

Tout cela nous mène à la nécessité d'automatiser le processus. J'ai essayé les méthodes suivantes pour résoudre ce problème :

  1. Utiliser MiNiFi au lieu de NiFi
  2. NiFi CLI
  3. NiPyAPI

Utilisation de MiNiFi

Apache MiNiFi — sous-projet Apache NiFi. MiNiFy — un agent compact utilisant les mêmes processeurs que NiFi, permettant de créer les mêmes flux que dans NiFi. La légèreté de l'agent est également due au fait que MiNiFy n'a pas d'interface graphique pour la configuration des flux. L'absence d'interface graphique dans MiNiFy signifie qu'il faut résoudre le problème de la livraison des flux dans Minifi. Étant donné que MiNiFy est largement utilisé dans l'IOT, il y a beaucoup de composants et le processus de livraison des flux aux instances finales de Minifi doit être automatisé. Vous connaissez ce défi, n'est-ce pas?

Pour résoudre ce problème, un autre sous-projet peut aider — MiNiFi C2 Server. Ce produit est destiné à être le point central de l'architecture de déploiement des configurations. Comment configurer votre environnement est décrit dans cet article sur Habré et l'information est suffisante pour résoudre le problème posé. MiNiFi, associé à C2 Server, met automatiquement à jour sa configuration. Le seul inconvénient de cette approche est qu'il faut créer des modèles sur le C2 Server, un simple commit dans le registre n'est pas suffisant.

La méthode décrite dans l'article précédent fonctionne et est simple à mettre en œuvre, mais il ne faut pas oublier ce qui suit :

  1. Dans Minifi, tous les processeurs de NiFi ne sont pas disponibles.
  2. Les versions des processeurs dans Minifi sont en retard par rapport aux versions des processeurs dans NiFi.

Au moment de la rédaction de cette publication, la dernière version de NiFi est 1.9.2. La version des processeurs de la dernière version de MiNiFi est 1.7.0. Des processeurs peuvent être ajoutés à MiNiFi, mais en raison de la divergence des versions entre les processeurs NiFi et MiNiFi, cela peut ne pas fonctionner.

NiFi CLI

Selon le la description de l'outil sur le site officiel, c'est un outil pour automatiser l'interaction entre NiFi et NiFi Registry dans le domaine de la livraison des flux ou de la gestion des processus. Pour commencer à utiliser cet outil, il est nécessaire de le télécharger d'ici.

Lançons l'utilitaire

.\/bin\/cli.sh
           _     ___  _
 Apache   (_)  .' ..](_)   ,
 _ .--.   __  _| |_  __    )
[ `.-. | [  |'-| |-'[  |  \/  
|  | | |  | |  | |   | | '    '
[___||__][___][___] [___]',  ,'
                           `'
          CLI v1.9.2

Type 'help' to see a list of available commands, use tab to auto-complete.

Pour charger le flux nécessaire depuis le registre, nous devons connaître les identifiants de seau (bucket identifier) et de flux (flow identifier). Ces données peuvent être obtenues soit par CLI, soit via l'interface web de NiFi registry. Dans l'interface web, cela ressemble à :

Automatisation de la livraison des flux dans Apache NiFi

Avec le CLI, cela se fait de la manière suivante :

#> registry list-buckets -u http://nifi-registry:18080

#   Name             Id                                     Description
-   --------------   ------------------------------------   -----------
1   test_bucket   709d387a-9ce9-4535-8546-3621efe38e96   (empty)

#> registry list-flows -b 709d387a-9ce9-4535-8546-3621efe38e96 -u http://nifi-registry:18080

#   Name           Id                                     Description
-   ------------   ------------------------------------   -----------
1   test_flow   d27af00a-5b47-4910-89cd-9c664cd91e85

Lançons l'importation du groupe de processus depuis le registre :

#> nifi pg-import -b 709d387a-9ce9-4535-8546-3621efe38e96 -f d27af00a-5b47-4910-89cd-9c664cd91e85 -fv 1 -u http://nifi:8080

7f522a13-016e-1000-e504-d5b15587f2f3

Un point important — l'hôte sur lequel nous appliquons le groupe de processus peut être n'importe quelle instance de NiFi.

Le groupe de processus a été ajouté avec des processeurs arrêtés, il faut les démarrer.

#> nifi pg-start -pgid 7f522a13-016e-1000-e504-d5b15587f2f3 -u http://nifi:8080

Super, les processeurs ont démarré. Cependant, selon les conditions de la tâche, nous devons faire en sorte que les instances NiFi envoient des données à d'autres instances. Supposons que nous ayons choisi la méthode Push pour transmettre les données vers le serveur. Pour organiser la transmission de données, nous devons activer la transmission des données sur le Remote Process Group (RPG) ajouté, qui est déjà inclus dans notre flux.

Automatisation de la livraison des flux dans Apache NiFi

Je n'ai pas trouvé dans la documentation CLI et d'autres sources comment activer la transmission de données. Si vous savez comment le faire, veuillez l'écrire dans les commentaires.

Puisqu'on a bash et que nous sommes prêts à aller jusqu'au bout, trouvons une solution ! Nous pouvons utiliser l'API NiFi pour résoudre ce problème. Utilisons la méthode suivante, l'ID est pris des exemples ci-dessus (dans notre cas, c'est 7f522a13-016e-1000-e504-d5b15587f2f3). Description des méthodes de l'API NiFi. ici.

Automatisation de la livraison des flux dans Apache NiFi
Dans le corps, il faut passer un JSON d'une forme suivante :

{
    "revision": {
	    "clientId": "value",
	    "version": 0,
	    "lastModifier": "value"
	},
    "state": "value",
    "disconnectedNodeAcknowledged": true
}

Paramètres à remplir pour que ça « fonctionne » :
state — état de transmission des données. Disponible TRANSMITTING pour activer la transmission, STOPPED pour l'arrêter.
version — version du processeur.

La version par défaut sera 0 lors de la création, mais ces paramètres peuvent être obtenus en utilisant la méthode.

Automatisation de la livraison des flux dans Apache NiFi

Pour les amateurs de scripts bash, cette méthode peut sembler utile, mais j'ai du mal — les scripts bash ne sont pas ma tasse de thé. La méthode suivante est plus intéressante et plus pratique à mon avis.

NiPyAPI

NiPyAPI — une bibliothèque pour le langage Python pour interagir avec les instances NiFi. La page de la documentation contient les informations nécessaires pour travailler avec la bibliothèque. Un guide rapide est décrit sur sur lequel je travaille actuellement, il faut atteindre les hôtes derrière le NAT de l'extérieur. En utilisant pour cela des protocoles avec une cryptographie mature, je n'ai pas pu m'empêcher de penser que c'était un peu comme utiliser un canon pour tuer des moineaux. Comme le tunnel est principalement utilisé juste pour percer un trou dans le NAT, le trafic interne est généralement aussi chiffré, on privilégie HTTPS. github.

Notre script pour déployer la configuration — un programme en Python. Passons au codage.
Nous configurons les paramètres pour notre travail futur. Nous aurons besoin des paramètres suivants :

nipyapi.config.nifi_config.host = 'http://nifi:8080/nifi-api' # chemin vers l'instance nifi-api sur laquelle nous déployons le groupe de processus
nipyapi.config.registry_config.host = 'http://nifi-registry:18080/nifi-registry-api' # chemin vers l'API registry nifi-registry
nipyapi.config.registry_name = 'MyBeutifulRegistry' # nom du registre, qui sera utilisé dans l'instance NiFi
nipyapi.config.bucket_name = 'BucketName' # nom du bucket à partir duquel nous tirons le flux
nipyapi.config.flow_name = 'FlowName' # nom du flux que nous tirons

Je vais maintenant énumérer les noms des méthodes de cette bibliothèque, qui sont décrites ici.

Connectons le registre à l'instance NiFi avec l'aide de

nipyapi.versioning.create_registry_client

À cette étape, vous pouvez également ajouter une vérification pour garantir que le registre a déjà été ajouté à l'instance, en utilisant la méthode

nipyapi.versioning.list_registry_clients

Nous trouvons le bucket pour rechercher ultérieurement le flow dans le panier

nipyapi.versioning.get_registry_bucket

Nous recherchons le flow dans le bucket trouvé

nipyapi.versioning.get_flow_in_bucket

Il est ensuite important de comprendre si ce groupe de processus a déjà été ajouté. Un groupe de processus est positionné par coordonnées et il peut arriver qu'un second soit superposé à un composant existant. J'ai vérifié, cela peut arriver 🙂 Pour obtenir tous les groupes de processus ajoutés, nous utilisons la méthode

nipyapi.canvas.list_all_process_groups

et nous pouvons ensuite rechercher, par exemple par nom.

Je ne décrirai pas le processus de mise à jour du modèle, je dirai juste que si des processeurs sont ajoutés dans la nouvelle version du modèle, il n'y a pas de problèmes de messages dans les files d'attente. En revanche, si des processeurs sont supprimés, des problèmes peuvent survenir (nifi n'autorise pas la suppression d'un processeur si des messages sont en attente devant lui). Si cela vous intéresse comment j'ai résolu ce problème — écrivez-moi, s'il vous plaît, nous discuterons de ce point. Mes coordonnées à la fin de l'article. Passons à l'étape d'ajout du groupe de processus.

En déboguant le script, j'ai rencontré une particularité, à savoir que la dernière version du flow n'est pas toujours récupérée, je vous recommande donc de confirmer d'abord cette version :

nipyapi.versioning.get_latest_flow_ver

Déployons le groupe de processus :

nipyapi.versioning.deploy_flow_version

Lançons les processeurs :

nipyapi.canvas.schedule_process_group

Dans la section sur le CLI, il était dit qu'il n'y a pas de transmission de données activée par défaut dans le groupe de processus distant ? Lors de la mise en œuvre du script, j'ai rencontré ce problème également. À ce moment-là, je n'ai pas pu activer la transmission de données via l'API et j'ai décidé d'écrire au développeur de la bibliothèque NiPyAPI pour demander conseil/aide. Le développeur m'a répondu, nous avons discuté du problème et il a écrit qu'il avait besoin de temps pour “vérifier quelque chose”. Et voilà, quelques jours plus tard, je reçois un e-mail contenant une fonction en Python qui résout mon problème de lancement !!! À ce moment-là, la version de NiPyAPI était 0.13.3 et, bien sûr, elle ne comprenait pas cette fonctionnalité. En revanche, dans la version 0.14.0, qui est sortie très récemment, cette fonction a déjà été ajoutée à la bibliothèque. Voici,

nipyapi.canvas.set_remote_process_group_transmission

Ainsi, grâce à la bibliothèque NiPyAPI, nous avons connecté le registre, appliqué le flux et même lancé les processeurs ainsi que le transfert de données. Ensuite, nous pouvons peaufiner le code, ajouter toutes sortes de vérifications, de journalisation, et tout ça. Mais c'est déjà une autre histoire.

Parmi les options d'automatisation que j'ai examinées, la dernière m'a semblé la plus fonctionnelle. Tout d'abord, il s'agit toujours de code en python, dans lequel on peut intégrer du code auxiliaire et bénéficier de tous les avantages du langage de programmation. Deuxièmement, le projet NiPyAPI évolue activement et en cas de problèmes, il est possible d'écrire au développeur. Troisièmement, NiPyAPI est tout de même un outil plus flexible pour interagir avec NiFi dans la résolution de tâches complexes. Par exemple, pour déterminer si les files d'attente de messages sont actuellement vides dans le flux et si l'on peut mettre à jour le groupe de processus.

Voilà, c'est tout. J'ai décrit 3 approches pour automatiser la livraison du flux dans NiFi, les pièges auxquels un développeur peut être confronté et fourni un code fonctionnel pour automatiser la livraison. Si ce sujet vous intéresse autant que moi, écrivez-moi !

Source : habr.com

Acheter un hébergement fiable pour les sites avec protection DDoS, serveurs VPS VDS 🔥 Acheter un hébergement fiable pour les sites avec protection DDoS, serveurs VPS VDS | ProHoster