Table of Contents
Comprendre le traitement des données en temps réel
Contrairement au traitement par lots, où les données sont recueillies sur une période donnée puis traitées en vrac, le traitement en temps réel exige des latences de sous-secondes. Cette distinction est essentielle dans les cas d'utilisations techniques comme la maintenance prédictive, où un retard dans l'analyse des données de vibration d'une turbine peut entraîner une défaillance catastrophique; ou dans la gestion intelligente du réseau, où les fluctuations de tension doivent être corrigées en millisecondes pour éviter les pannes.
Pour répondre à ces exigences, les modèles de données doivent être conçus avec une compréhension approfondie de la vitesse, de la variété et du volume des données. Les lectures de capteurs provenant d'Internet des objets (IoT) arrivent souvent à des millions d'événements par seconde, contenant chacun des horodatages, des identifiants et des mesures multiples.
Les principaux défis sont le traitement des données hors-commande, la gestion des événements survenus en fin d'arrivée et la sémantique du traitement exactement une fois lorsque les duplicatas ne peuvent pas être tolérés. Un modèle de données bien conçu résume ces complexités, fournissant une interface propre pour les ingénieurs pour interroger et visualiser les données en temps réel.
Principes fondamentaux pour la conception de modèles de données dans les systèmes en temps réel
La conception d'un modèle de données pour les données d'ingénierie en temps réel nécessite un équilibre entre plusieurs principes fondamentaux, qui guident les décisions sur la conception des schémas, les moteurs de stockage et les modèles de requêtes.
Écacité et élasticité
Le modèle de données doit être à l'échelle horizontale pour pouvoir accueillir des volumes croissants de données sans dégradation des performances, ce qui implique souvent la partition des données entre plusieurs nœuds. Par exemple, les données de séries chronologiques peuvent être partitionnées par plage de temps ou par un hachage de l'ID du capteur. L'élasticité permet au système d'ajouter ou de supprimer automatiquement des nœuds au fur et à mesure des changements de charge, ce qui est particulièrement important dans les environnements d'ingénierie où les données éclatent lors d'expériences ou de mises en production.
Faible latence Lire et écrire des chemins
Les applications en temps réel nécessitent des opérations d'écriture et de lecture pour se terminer en millisecondes.Les structures de données qui supportent les écritures en appendice seulement, comme les arbres de fusion log-structurés (LSM), sont courantes dans les bases de données telles que InfluxDB ou TimescaleDB. Pour lire, le modèle doit supporter des balayages efficaces de la plage de temps au fil des fenêtres et des points de recherche pour des états de périphérique spécifiques.
Cohérence et intégrité des données
Dans les contextes techniques, la précision des données n'est pas négociable. Le modèle de données doit imposer des contraintes de cohérence, comme s'assurer qu'une lecture de température se situe dans une plage prédéfinie. Les stratégies de résolution des conflits, comme les derniers wins-write ou les vecteurs de version, sont appliquées lorsque les données arrivent de plusieurs sources.
Flexibilité pour l'adaptation des schémas en évolution
Les projets d'ingénierie ajoutent fréquemment de nouveaux capteurs, modifient les taux d'échantillonnage ou introduisent de nouveaux types de mesures. Un schéma rigide et prédéfini se brise lorsque les données changent. Des modèles de données flexibles, comme les approches de schéma en lecture (par exemple, en utilisant JSONB dans PostgreSQL ou des colonnes dynamiques dans Cassandra), permettent aux ingénieurs d'ingérer des données sans modifier le schéma de stockage.
Choisir les bonnes structures de données et moteurs de stockage
Le choix des structures de données a une incidence directe sur la capacité du système à traiter les données en temps réel. Ci-dessous sont les structures les plus couramment utilisées dans les modèles de données d'ingénierie, ainsi que leurs compromis.
Bases de données série chronologique
Les bases de données de séries chronologiques (TSDB) sont conçues pour stocker et interroger les points de données séquentiels indexés par le temps. Elles compressent généralement les données efficacement en utilisant l'encodage delta et l'encodage de la longueur d'exécution, réduisant ainsi les coûts de stockage. Les TSDB prennent également en charge les politiques de réduction de l'échantillonnage et de rétention qui regroupent automatiquement ou suppriment les données anciennes.
Magasins de valeurs clés
Les magasins à valeur clé sont excellents pour les recherches en temps réel de l'état ou de la configuration des appareils. Ils offrent une latence extrêmement faible pour les lectures et les écritures ponctuelles. Dans les modèles de données techniques, la clé est souvent un composite d'ID et d'horodatage des appareils, tandis que la valeur est un blob sérialisé de lectures de capteurs.
Magasins autochtones à traitement par flux
Les technologies comme Apache Kafka , ou Apache Flink , permettent de traiter et de stocker les données dans le flux lui-même. Cette architecture réduit le besoin de bases de données séparées lorsque le cas d'utilisation primaire est analytique en temps réel et alerte. Par exemple, un modèle de données mis en œuvre avec Kafka Streams peut maintenir, dans un magasin local, les dix dernières minutes de données de vibration pour chaque machine et déclencher une alerte lorsque la moyenne mobile dépasse un seuil.
Approches hybrides
De nombreux systèmes d'ingénierie utilisent une stratégie hybride : un processeur de flux pour l'analyse en temps réel, un TSDB pour le stockage historique et un stock de valeurs clés pour l'état actuel. Cette architecture fournit peu de latence pour les tableaux de bord opérationnels tout en permettant une analyse historique profonde. Le modèle de données doit définir comment les flux de données entre ces couches, souvent en utilisant des modèles de capture de données de changement (CDC) ou de double écriture.
Stratégies de conception pour les modèles de données d'ingénierie
Des modèles de données efficaces pour les données d'ingénierie en temps réel sont conçus avec des stratégies spécifiques qui répondent aux contraintes uniques du domaine.
Modèles de dispositifs et capteurs
Dans un modèle relationnel, vous pouvez avoir une table avec des métadonnées (emplacement, fabricant, date d'installation) et une table avec du temps, du type de capteur et de la valeur. Cependant, dans des scénarios en temps réel, la table de mesures peut croître rapidement des milliards de lignes. Une meilleure conception est d'utiliser un modèle de séries chronologiques où chaque mesure est stockée comme une ligne avec un horodatage, un identifiant de périphérique et une charge utile de paires de valeurs clés pour différentes mesures.
Exemple d'enregistrement de mesure à plat :
timestamp: 2025-03-09T14:30:01.234Z, device id: "sensor-42", métriques: {"température": 68.2, "humidité": 45.1, "pression": 1013.2}
Normalisation par rapport à la dénormalisation
La normalisation réduit la redondance des données et améliore les performances d'écriture en stockant les métadonnées séparément. Dans les systèmes en temps réel, cependant, souvent, l'intégration du flux de mesure avec les métadonnées des appareils peut introduire une latence. La dénormalisation est souvent préférée pour les requêtes de chemin chaud. Par exemple, l'emplacement du périphérique directement dans la ligne de mesure élimine une jointure pendant l'alerte. L'échange est accru le stockage et l'incohérence potentielle lorsque les métadonnées des appareils changent (p. ex., un capteur est déplacé).
Partitionnement et rembourrage
La partition des données est essentielle pour l'évolutivité. La partition basée sur le temps est la plus courante pour les données de séries chronologiques : chaque partition couvre un intervalle de temps spécifique (par exemple, une heure ou une journée). Cela permet au système de laisser tomber rapidement les anciennes partitions et d'effectuer des requêtes de portée efficacement. La partition basée sur l'ID des périphériques distribue la charge uniformément entre les nœuds, mais elle peut conduire à des points chauds si certains appareils génèrent beaucoup plus de données que d'autres. Une combinaison de temps et de hachage des périphériques fonctionne bien.
Indexation pour la performance de requête
Les stratégies d'indexation doivent être adaptées aux modèles de requêtes les plus courants : « obtenir toutes les données pour l'appareil X au cours de la dernière heure » ou « trouver tous les appareils dont la température dépasse 100°C à la dernière minute ». Un index basé sur le temps combiné avec un index de balise de périphérique est typique. Les techniques avancées comprennent l'utilisation d'un index de liste de saut pour les bases de données de séries chronologiques ou d'un index bitmap pour les balises de faible cardinalité.
Mise en œuvre avec les technologies de traitement de flux
Les modèles de données d'ingénierie en temps réel sont souvent construits sur des cadres de traitement de flux qui fournissent exactement une fois sémantique, tolérance aux défauts et gestion de l'état.
Apache Kafka
Le modèle de données pour les sujets de Kafka devrait s'aligner avec les consommateurs en aval. Par exemple, chaque type d'appareil peut avoir son propre sujet, ou tous les appareils partagent un seul sujet avec une partition par groupe de périphérique. Le schéma de message (par exemple, Avro ou Protobuf) comprend un horodatage, un identifiant de périphérique et la charge utile des mesures. Compaction peut être activé pour conserver seulement la dernière valeur pour chaque clé, ce qui est utile pour les mises à jour de l'état de l'appareil.
Lien Apache
Le modèle de données de Flink est défini par les types d'événements et les descripteurs d'état. Par exemple, pour détecter les modèles de vibrations anormales, Flink maintient un état qui stocke les 100 dernières lectures d'accélération par appareil. Le modèle de données devrait être conçu pour minimiser la taille de l'état; utiliser des dictionnaires pour les ID de capteurs et compresser les champs répétés. Flink prend également en charge le traitement événement-temps, de sorte que le modèle de données doit inclure l'horodatage événement (pas l'horodatage de traitement) pour la fenêtre correcte.
Flèches d'Apache
Le modèle de données peut être représenté par un DataFrame ou Dataset, avec des schémas définis en code. Bien que le micro-batch introduit une latence plus élevée que le streaming pur (par exemple Flink), il est plus facile d'utiliser pour les charges de travail analytiques qui doivent joindre des flux avec des tables historiques. Le modèle de données devrait tenir compte du mécanisme de contrôle utilisé par Spark pour maintenir exactement une sémantique, qui écrit état dans un répertoire de points de contrôle.
Intégration des bases de données
Les processeurs de flux écrivent souvent dans une base de données en temps réel. Le modèle de données doit définir la correspondance entre le flux d'événements et le schéma de base de données. Par exemple, un travail Flink lit les données brutes de capteur de Kafka, applique un certain filtrage, et écrit à InfluxDB en utilisant son protocole de ligne. Le schéma de base de données , les noms de mesure, les balises et les champs doivent être conçus pour correspondre aux requêtes que les tableaux de bord exécuteront.
Étude de cas : Modèle de données pour un système de maintenance prédictive en temps réel
Considérez une usine de 10 000 machines, chacune équipée de capteurs mesurant la température, les vibrations et la vitesse de rotation. L'objectif est de prédire les défaillances 30 minutes à l'avance et déclencher des alertes de maintenance.
Le modèle de données est conçu comme suit:
- Couche d'ingestion:[ Chaque machine envoie un message JSON à chaque seconde à un sujet Kafka partitionné par groupe de machine. Le message comprend un horodatage, un identifiant de machine et trois paramètres.
- Traitement de la commande :[ Un travail Flink consomme le sujet. Il maintient une fenêtre coulissante de 30 minutes par machine en utilisant le magasin d'état Flink. L'état est cliqué par l'ID de la machine et stocké comme une liste des 1800 dernières lectures (30 minutes x 60 secondes). Pour chaque nouvelle lecture, le travail calcule une moyenne mobile et un écart-type pour chaque métrique. Si le z-score dépasse 3, il envoie une alerte à un sujet Kafka séparé.
- Base de données: Le travail Flink écrit aussi chaque lecture brute à TimescaleDB. Le schéma de table utilise une hypertable partitionnée par le temps (1 heure de morceaux) et indexée par ID machine. Les mots clés comme groupe machine et emplacement sont stockés dans une table de métadonnées séparée, rejointe uniquement pour les requêtes analytiques.
- Real-Time Dashboard:[ Le tableau de bord interroge TimescaleDB pour la dernière heure de données par machine, en utilisant un agrégat continu qui précalcule min, max et avg par minute. Les règles d'alerte sont évaluées par le processeur de flux, et non par la base de données, pour garder la latence en dessous de 100 ms.
Ce modèle hybride permet de concilier la nécessité d'alertes à faible latence (via le traitement de flux) avec une analyse historique flexible (via une base de données série temporelle).Le modèle de données reste simple : un hypertable unique pour les données brutes, avec des index optimisés pour le motif de requête le plus courant (intervalle de temps + identifiant machine).
Meilleures pratiques de déploiement de la production
Passer de la conception à la production exige une attention à la surveillance, à l'évolution des schémas et à la gestion des coûts.
Surveillance et profil Demande de renseignements Performance
Utilisez des outils spécifiques à une base de données (p. ex. TimescaleDB=s , InfluxDB=s request inspector) pour identifier les requêtes lentes. Surveillez le débit et la latence d'écriture; si vous écrivez des pics de latence, envisagez d'augmenter le nombre de partitions ou de mettre en accord la stratégie de compactage.
Plan pour Schema Evolution
Utilisez des registres de schémas (comme le registre du schéma de confluent) pour gérer les schémas Avro ou Protobuf. Pour les bases de données qui supportent l'évolution du schéma (par exemple, ajouter de nouveaux champs à une colonne JSONB), assurez la compatibilité en arrière. Évitez les changements destructeurs aux tables de production; au lieu d'ajouter de nouvelles colonnes ou créer de nouvelles tables et migrer les données asynchrones.
Optimiser les coûts
Les données de séries chronologiques peuvent être coûteuses à stocker à haute granularité. Mettre en œuvre des politiques de conservation pour supprimer automatiquement les données plus anciennes qu'un certain seuil.Utiliser l'échantillonnage : stocker les données brutes pendant 7 jours, puis les moyennes d'une minute pendant 30 jours, puis les moyennes horaires pendant 1 an.
Essai avec volume réel de données
Simulez le taux de données attendu dans un environnement de mise en scène avant d'aller à la production. Mesurez la distribution de latence (p50, p99, p999) pour les deux écritures et les lectures. Assurez-vous que le modèle de données peut gérer les charges de pointe (p. ex., pendant le démarrage de la machine lorsque de nombreux capteurs envoient des données simultanément).
Tendances futures de la modélisation des données d'ingénierie en temps réel
Les nouvelles tendances comprennent l'utilisation de bases de données accélérées par le GPU[ pour l'analyse en temps réel sur de gros ensembles de données, et l'adoption d'un calcul de bord où les modèles de données doivent travailler sur des appareils à ressources limitées. Une autre tendance est l'intégration des modèles ML directement dans le pipeline de données, exigeant des modèles de données qui peuvent servir de vecteurs et de prédictions de caractéristiques aux côtés de données brutes de capteurs.
Les ingénieurs doivent rester informés des avancées dans le streaming SQL (par exemple, matérialisez, RisingWave) qui permettent l'analyse en temps réel avec SQL standard, réduisant le besoin de code de traitement de flux personnalisé. Ces outils font appliquer un modèle de données déclaratif qui gère automatiquement l'état et les indices.
Conclusion
La conception de modèles de données pour le traitement des données en temps réel est une tâche complexe mais enrichissante. En respectant les principes de l'évolutivité, de la faible latence, de flexibilité et de cohérence, et en choisissant les bonnes structures de données et les technologies de traitement des flux, les ingénieurs peuvent construire des systèmes qui fournissent des informations opportunes et maintiennent la continuité opérationnelle. La clé est de comprendre les modèles de requête spécifiques et les exigences de latence de votre application, prototype avec des données réelles, et itérer sur le modèle à mesure que le paysage d'ingénierie évolue.