Bonjour à tous !

La tâche consiste à déployer le flux, illustré dans l'image ci-dessus, sur N serveurs avec . 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 et il est important de ne pas oublier de configurer l'instance NiFi pour autoriser S2S, voir .
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 :
- 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
- 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 . 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 :
- le temps nécessaire pour actualiser les flux augmente. Il faut se connecter à tous les serveurs
- des erreurs d'actualisation des modèles surviennent. Ici, cela a été mis à jour, mais là, cela a été oublié
- 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 :
- Utiliser MiNiFi au lieu de NiFi
- NiFi CLI
- NiPyAPI
Utilisation de 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 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 :
- Dans Minifi, tous les processeurs de NiFi ne sont pas disponibles.
- 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 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 .
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 à :

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.

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. .

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.

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. contient les informations nécessaires pour travailler avec la bibliothèque. Un guide rapide est décrit sur 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 .
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_clientsNous trouvons le bucket pour rechercher ultérieurement le flow dans le panier
nipyapi.versioning.get_registry_bucketNous recherchons le flow dans le bucket trouvé
nipyapi.versioning.get_flow_in_bucketIl 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_groupset 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_verDéployons le groupe de processus :
nipyapi.versioning.deploy_flow_versionLançons les processeurs :
nipyapi.canvas.schedule_process_groupDans 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_transmissionAinsi, 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,
Source : habr.com
