Construire un pipeline de supervision de flotte IoT
Ce tutoriel vous guide à travers un pipeline complet en temps réel sur la Data Platform : streaming de télémétrie de capteurs IoT simulés depuis une action Custom vers
Objectif
Ce tutoriel vous guide à travers un pipeline complet en temps réel sur la Data Platform : streaming de télémétrie de capteurs IoT simulés depuis une action Custom vers Apache Kafka, ingestion dans une table du Lakehouse Manager, et visualisation de la santé de la flotte sur un dashboard Superset en direct.
À la fin, vous disposerez d'un dashboard fonctionnel qui montre la santé des devices, les tendances environnementales, les anomalies, et les temps d'arrêt, piloté par un scénario de panne scripté afin que les données racontent toujours une histoire.
Introduction
Prérequis
Pour suivre ce tutoriel, vous avez besoin de :
- Un projet Data Platform OVHcloud avec le Data Processing Engine, le Lakehouse Manager, et une configuration Superset/Trino disponibles.
- Un broker Apache Kafka sur lequel vous pouvez écrire. Ce tutoriel utilise un service Kafka managé OVH avec authentification par certificat client (mTLS), mais l'approche fonctionne avec n'importe quel broker une fois les paramètres de connexion adaptés.
- Un topic Kafka nommé
iot_readingscréé sur votre broker (l'auto-création de topic est souvent désactivée).
Nous recommandons de d'abord terminer le premier tutoriel de démarrage et le tutoriel Streamer des données depuis Apache Kafka. Ce guide suppose que vous êtes à l'aise avec les principaux composants de la platform.
Ce que vous allez construire
Le scénario
Le producer simule une flotte de capteurs environnementaux (température, humidité, CO₂, PM2.5) répartis sur quatre sites dans trois régions. Une machine à états par device pilote un cycle de panne réaliste afin que les données ne soient jamais plates :
Deux devices sont scriptés pour tomber en panne selon un calendrier fixe, afin que chaque exécution montre de manière fiable l'histoire complète. Les relevés suivent un « cycle quotidien » compressé sur dix minutes avec du bruit aléatoire, de sorte que les tendances paraissent crédibles. Lorsqu'un device passe OFFLINE, il cesse complètement d'émettre, ce qui apparaît plus tard comme un véritable trou dans les données.
À propos du format des messages. Le connecteur Kafka ne lit que les champs JSON de premier niveau, et il infère le schéma de la table en échantillonnant les messages. Chaque message émet donc tous les champs avec des types stables (les valeurs numériques sont toujours des floats, error_code est la chaîne "NONE" quand tout va bien). Un champ parfois absent, ou parfois un entier et parfois un float, casserait le schéma inféré.
Étape 1 : stocker vos identifiants Kafka dans un bucket
Plutôt que de coller des certificats ou des mots de passe dans la source de l'action, stockez-les dans un bucket du Lakehouse Manager et récupérez-les au runtime. Une action Custom est auto-authentifiée à son projet, elle peut donc lire le bucket sans identifiants supplémentaires, et rien de sensible ne vit dans votre code.
- Dans le Lakehouse Manager, créez un bucket nommé
iot-demo-certs. - Chargez-y vos trois fichiers mTLS :
Le producer télécharge ces fichiers au démarrage à l'aide du connecteur bucket du SDK.
Ne commitez jamais de fichiers de certificat ou de clé dans un repository, et ne les collez jamais dans la source de l'action. Conservez-les uniquement dans le bucket.
Étape 2 : créer l'action Custom producer
Téléchargez le producer et ajoutez-le comme action Custom :
- Allez dans Data Processing Engine > Actions > New > Custom.
- Chargez directement le fichier
producer.pytéléchargé. - Définissez la fonction d'entrée sur
customfuncdans le panneau Information de l'action. - Ajoutez
kafka-pythonaux Python Requirements de l'action. - Réglez le mode d'exécution sur Always-up. Le producer boucle en continu, donc le mode Serverless expirerait par timeout.
Utilisez kafka-python, pas kafka. Le package kafka brut sur PyPI est une distribution abandonnée, Python 2 uniquement, et échoue sur les workers avec invalid syntax (simple.py, line 54). C'est le package kafka-python qui fournit l'espace de noms from kafka import ....
Paramètres que vous pouvez ajuster
Toute la configuration se trouve en haut du fichier. Ceux que vous serez le plus susceptible de modifier :
Points clés du code
Vous n'avez pas besoin de lire tout le fichier pour l'exécuter, mais certaines parties méritent d'être connues.
Bloc de connexion. Définissez ici l'endpoint de votre broker et la méthode d'authentification, et choisissez si l'exécution est finie :
Identifiants depuis un bucket. Les certificats sont récupérés au runtime, jamais codés en dur, si bien que rien de sensible ne vit dans l'action :
La flotte. Modifiez les sites ou le nombre par site pour redimensionner la simulation :
Le scénario de panne. Deux devices tombent en panne selon un calendrier fixe, et ces durées fixent le rythme de l'histoire OK vers DEGRADED vers FAULT vers OFFLINE :
Lancez l'action. Chaque message porte les douze champs avec des types stables, si bien que le schéma est complet dès que des données commencent à circuler. Le premier device scripté se dégrade après environ quinze secondes, donc les échantillons DEGRADED et FAULT apparaissent rapidement. Surveillez les logs de l'action pour la ligne de heartbeat : sent ~N messages; fleet: {...}.
Étape 3 : connecter Kafka avec le connecteur source
Allez dans Connectors > Sources > Apache Kafka et créez une connexion. Pointez-la vers l'endpoint de votre broker et le topic iot_readings.
Pour un broker mTLS, authentifiez-vous avec les champs SSL du connecteur : chargez votre certificat CA, votre certificat client, et votre clé client, et laissez le nom d'utilisateur et le mot de passe vides.
Le formulaire du connecteur expose des champs SSL même lorsque le broker utilise des certificats clients. Ce sont eux qui font fonctionner le handshake mTLS. Pour tous les détails, consultez la référence du connecteur Kafka.
Étape 4 : extraire les métadonnées et créer la table
Une fois les données en circulation, ouvrez l'Analyzer et extrayez les métadonnées.
Vous verrez seize attributs : les douze champs du producer, plus quatre colonnes d'enveloppe ajoutées par Kafka (timestamp, date, offset_r, partition). Pour une simple table en ajout seul, vous pouvez ignorer les quatre colonnes d'enveloppe. Ne conservez partition et offset_r que si vous voulez une sémantique d'upsert ou de déduplication.
Le producer émet ts sous forme de chaîne ISO-8601. S'il est inféré comme une chaîne, vous pouvez soit le définir comme type Timestamp dans la table, soit le conserver comme chaîne et le parser dans Trino (voir Bon à savoir). Construisez la table du Lakehouse Manager iot_readings à partir du schéma obtenu.
La liste complète des colonnes :
OFFLINE n'apparaît jamais comme valeur de ligne. Un device hors ligne cesse d'émettre, donc le temps d'arrêt apparaît comme un trou (aucune ligne pour ce device_id). Vous le détectez en comparant le ts le plus récent de chaque device à l'heure actuelle.
Étape 5 : charger le flux dans votre table
Créez une action Load dans le Data Processing Engine, en associant le topic iot_readings à votre table iot_readings. Exécutez-la en mode Always-up afin qu'elle continue à consommer le flux en direct.
L'action Load continue de s'exécuter tant que des données existent dans Kafka, jusqu'à son timeout. C'est le comportement attendu pour un chargement en streaming.
Une action Load qui se termine normalement lance automatiquement une mise à jour des métadonnées, de sorte que le nombre de lignes de la table se rafraîchit tout seul. Dans cette configuration en streaming, l'action s'exécute en Always-up et est arrêtée manuellement ou expire par timeout plutôt que de se terminer proprement, si bien que cette mise à jour automatique peut ne pas s'exécuter. Dans ce cas, lancez une action Update Metadata afin que le nombre de lignes dans le Lakehouse Manager reflète les données ingérées.
Étape 6 : requêter la table depuis Superset
Connectez Superset à votre table via Trino. Si vous n'avez pas encore déployé Superset, suivez Déployer Apache Superset.
Dans Superset, ajoutez une connexion à une base de données avec une URI SQLAlchemy de la forme trino://<user>@<trino-host>:<port>/<catalog>. Testez la connexion avant de construire des graphiques. C'est l'étape où les configurations de bout en bout échouent le plus souvent.
Si ts est arrivé sous forme de chaîne, parsez-le dans Trino avec from_iso8601_timestamp(ts), qui retourne un TIMESTAMP(3) WITH TIME ZONE. Créez un dataset virtuel qui expose la colonne parsée et marquez-la comme la colonne temporelle principale du dataset, afin qu'elle devienne disponible comme axe X des graphiques de séries temporelles.
Étape 7 : construire les dashboards
Le dashboard fini offre une vue en un seul écran de toute la flotte.
Un ensemble pratique de panneaux :
Pour les datasets virtuels complets et le SQL Trino exact et la configuration de chaque panneau, suivez la page complémentaire :
Construire les dashboards Superset
Quelques astuces Superset qui font gagner du temps :
- Sur les bar charts, placez la catégorie dans l'axe X et laissez la case Dimensions vide. Dimensions ne sert qu'à découper chaque barre en sous-séries.
- Les lignes de seuil, comme un niveau d'alerte batterie, sont disponibles sur les graphiques de séries temporelles via des couches d'annotation. Sur un bar chart catégoriel, triez par ordre croissant à la place. Sur un table, utilisez le conditional formatting.
- Les valeurs de CO₂ (autour de 480) écrasent la temperature (autour de 21) et le PM2.5 (autour de 9) sur un axe partagé. Affichez le CO₂ séparément, ou utilisez un axe secondaire, afin que les petites métriques restent lisibles.
- Le tableau des temps d'arrêt est vide lorsque la flotte est en bonne santé, ce qui est normal. Construisez-le sans filtre afin qu'il liste chaque device, puis mettez en évidence les lignes obsolètes avec le conditional formatting, afin que le panneau ne paraisse jamais cassé.
Bon à savoir
- Le producer ne se termine jamais par conception. C'est une boucle infinie en mode Always-up. Le timeout par défaut de l'action est d'environ deux heures. Si vous l'arrêtez manuellement, il apparaît comme « stopped » plutôt que comme un succès, car il n'y a pas de fin naturelle. Pour obtenir une exécution propre en succès, définissez
MAX_RUNTIME_SECSsur un nombre de secondes. La boucle vide alors ses buffers, se ferme, enregistre un total, et se termine. Laissez-le àNonepour un dashboard en direct en continu. - Rafraîchissez le nombre de lignes avec Update Metadata. Une action Load qui se termine normalement met à jour automatiquement les métadonnées de la table. Dans cette configuration en streaming, l'action s'exécute en Always-up et est arrêtée ou expire par timeout au lieu de se terminer proprement, si bien que, comme indiqué à l'étape 5, vous devrez peut-être lancer une action Update Metadata pour que le nombre de lignes reflète ce qui a été ingéré.
- Parsez
tspour les séries temporelles. Si la colonne est une chaîne, parsez-la avecfrom_iso8601_timestamp(ts)et marquez-la comme colonne temporelle dans Superset. - Prémunissez-vous contre les lignes non parsables. Si un message de test a laissé une ligne dont
tsn'est pas un timestamp valide,from_iso8601_timestamplèveINVALID_FUNCTION_ARGUMENT. Enveloppez le parsing danstry(from_iso8601_timestamp(ts))afin que la ligne invalide prenne la valeur null, ou supprimez la ligne.
Ce que vous avez construit
Vous disposez maintenant d'un pipeline de streaming complet : une action Custom produisant de la télémétrie simulée, Kafka la transportant, un connecteur et une action Load l'ingérant dans le Lakehouse Manager, et Superset la visualisant en direct via Trino. Au fur et à mesure que les pannes scriptées se déroulent, le dashboard traverse toute l'histoire : une flotte en bonne santé, un device qui se dégrade puis tombe en panne, une fenêtre de temps d'arrêt qui apparaît dans le tableau des temps d'arrêt et comme un trou dans la tendance, et enfin une réparation qui ramène le device à la normale.
À partir de là, vous pouvez étendre le modèle avec une table de dimension device séparée, ajouter davantage de métriques, ou adapter le producer pour rejouer vos propres données de capteurs réelles.
Ce modèle passe à l'échelle. Les mêmes briques de base se composent en un pipeline beaucoup plus vaste. Comme un topic correspond à une table, vous gérez plusieurs flux de données en répétant les éléments : ajoutez davantage de topics Kafka, extrayez les métadonnées de chacun vers sa propre table du Lakehouse Manager, et exécutez une action Load par mapping topic-vers-table. Un seul producer peut émettre vers plusieurs topics, et Superset peut joindre les tables résultantes via Trino. Ainsi, une flotte multi-topics et multi-tables (par exemple, des flux séparés pour les relevés environnementaux, la consommation d'énergie, et les événements de device) n'est jamais que ce tutoriel appliqué plusieurs fois, alimentant un même ensemble de dashboards.
Aller plus loin
Si vous avez besoin d'une formation ou d'une assistance technique pour la mise en oeuvre de nos solutions, contactez votre commercial ou cliquez sur ce lien pour obtenir un devis et demander une analyse personnalisée de votre projet à nos experts de l’équipe Professional Services.
Posez vos questions, faites-nous part de vos commentaires et interagissez directement avec l’équipe qui développe la Data Platform sur le canal Discord dédié.
Si vous avez besoin d'une assistance concernant vos services OVHcloud, créez une demande depuis notre centre d'aide.
Rejoignez notre communauté d'utilisateurs.