Data flow graphs in Opérations Azure IoT

Un graphe de flux de données est un pipeline configurable qui traite les données au fur et à mesure qu’elles circulent dans Opérations Azure IoT. Un flux de données standard suit une séquence fixe d’enrichissement, de filtre, de cartographie, mais un graphe de flux de données permet de composer des transformations dans n’importe quel ordre, de ramifier en chemins parallèles et d’agréger les données sur des fenêtres temporelles.

La DataflowGraph ressource personnalisée Kubernetes définit un graphe de flux de données. À l’intérieur de la ressource, vous connectez les sources, transformations et destinations pour construire des pipelines de traitement adaptés à votre scénario.

Important

Les graphiques de flux de données prennent actuellement en charge uniquement les points de terminaison MQTT, Kafka et OpenTelemetry. D'autres types de points de terminaison tels que lac de données, Microsoft Fabric OneLake, Azure Data Explorer et stockage local ne sont pas pris en charge.

Flux de données et graphiques de flux de données

Les opérations Azure IoT offrent deux façons de traiter les données dans un pipeline :

Capacité Flux de données Graphiques de flux de données
Forme de pipeline Corrigé : enrichir, filtrer, mapper Flexible : tout ordre, branchement, fusion
Types de transformation Carte, filtre, enrichissement Cartographie, filtre, branchement, concaténation, fenêtre, accélération, enrichissement
Agrégation basée sur le temps Non disponible Transformations de fenêtre avec fenêtres tumbling
Routage conditionnel Non disponible Transformations de Branch et Concatenate
Prise en charge des points de terminaison Tous les types de points de terminaison MQTT, Kafka et OpenTelemetry uniquement

Pour les nouveaux projets qui utilisent des types de points de terminaison pris en charge, nous recommandons les graphes de flux de données. Les flux de données restent entièrement pris en charge pour tous les scénarios et prennent en charge la gamme complète de types de points de terminaison.

Transformations disponibles

Chaque transformation est une étape de traitement préconstruite que vous configurez avec des règles et que vous chaînez avec d’autres transformations à l’intérieur d’une DataflowGraph ressource.

Transformer Artéfact Description
Carte azureiotoperations/graph-dataflow-map:1.0.0 Renommez, restructurez, calculez et copiez des champs.
Filtrer azureiotoperations/graph-dataflow-filter:1.0.0 Supprimez les messages qui correspondent à une condition.
Branche azureiotoperations/graph-dataflow-branch:1.0.0 Acheminer chaque message vers un chemin true ou false en fonction d’une condition.
Concatenate azureiotoperations/graph-dataflow-concatenate:1.0.0 Fusionnez deux chemins d’accès ou plus en un.
Window azureiotoperations/graph-dataflow-window:1.0.0 Collecter des messages sur un intervalle de temps, puis agréger.
limitation du débit azureiotoperations/graph-dataflow-throttle:1.0.0 Limitez le taux de messages par schéma de sujet MQTT.

Toutes les transformations partagent un langage d’expression pour les opérateurs, les fonctions et les références de champ. Vous pouvez égalementenrichir les messages avec des données externes provenant d’un état de stockage dans les transformations map, filter et branch.

Tip

Les expressions utilisent des variables positionnelles, donc $1 est la première entrée, $2 la seconde, et ainsi de suite. La référence Expressions liste les fonctions intégrées comme cToF et couvre chaque opérateur, fonction et champ de métadonnées disponible pour les transformations.

Comment les transformations se composent dans un graphe de flux de données

Transformations se connectent dans une séquence à l’intérieur d’une DataflowGraph ressource : >Transformation source A > Transform B > ... > Destination.

Les transformations de branchement divisent le flux en chemins parallèles, et les transformations de concaténation les rassemblent.

Vous pouvez chaîner n’importe quel nombre de transformations dans n’importe quel ordre. Un pipeline avec une transformation de carte unique est aussi valide que celui qui filtre, branche, mappe chaque chemin différemment, fusionne, puis agrège au fil d’une fenêtre de temps.

Comment fonctionne la configuration des graphes de flux de données

Chaque transformation dans un graphe de flux de données fait référence à un artefact préconstruit extrait d’un registre conteneur. Vous configurez la transformation en passant des règles au format JSON via la configuration section de la ressource de graphe.

Lorsque vous déployez Opérations Azure IoT, il crée automatiquement un terminau de registre par défaut nommé default qui pointe vers mcr.microsoft.com. Les transformations intégrées utilisent ce point de terminaison pour extraire des artefacts de Microsoft Container Registry. Vous n’avez pas besoin de configuration supplémentaire pour le registre.

Une ressource de graphe de flux de données définit trois types d’éléments — une source, une ou plusieurs transformations (chacune avec nodeType: Graph), et une destination — et un ensemble de nodeConnections ces éléments décrivant comment les données circulent entre eux. Chaque transformation configuration passe ses règles sous forme de chaîne JSON sous la rules tonalité.

Pour un exemple complet et exécutable qui lit des données de température, convertit Celsius en Fahrenheit avec une transformation de map, et publie le résultat — dans l’expérience Opérations, Azure CLI, Bicep et Kubernetes — voir Créer un graphique de flux de données. Dans les articles pratiques qui suivent, des exemples se concentrent sur les règles de transformation elles-mêmes.

Configurer les schémas sur les connexions des nœuds

Les graphes de flux de données gèrent les schémas différemment des flux de données. Au lieu de définir le schéma sur la source ou la transformation, vous configurez des schémas sur les connexions de nœud entre les nœuds du graphique. Les transformations de branchement et de filtre peuvent éventuellement valider les données d’exécution par rapport aux schémas attachés aux connexions de nœuds.

Chaque entrée du nodeConnections tableau peut inclure un schema sur le from côté d’une connexion. Ce schéma décrit le format attendu des données qui circulent entre ces deux nœuds :

nodeConnections: [
  {
    from: {
      name: 'source'
      schema: {
        schemaRef: 'aio-sr://my-namespace/sensor-data:1'
        serializationFormat: 'Json'
      }
    }
    to: {
      name: 'transform'
    }
  }
]

La schemaRef valeur utilise le format aio-sr://<namespace>/<name>:<version> et pointe vers un schéma stocké dans le registre du schéma. Parce que les graphes de flux de données ne prennent en charge que les terminaux MQTT, Kafka et OpenTelemetry, le format de sérialisation supporté est Json.

Le tableau suivant résume comment la configuration des schémas diffère entre les flux de données et les graphiques de flux de données :

Aspect Flux de données Graphiques de flux de données
Emplacement du schéma Sur la source (sourceSettings.schemaRef) et la transformation (builtInTransformationSettings.schemaRef) Sur les connexions de nœud (nodeConnections[].from.schema)
Formats de destination pris en charge JSON, Parquet, Delta JSON
Validation du runtime Non supporté pour les schémas sources Optionnel sur les connexions de nœuds via les transformations de branchement et de filtre

Note

Pour les graphes de flux de données, JSON est actuellement le seul format de destination pris en charge, malgré les formats listés dans la documentation de référence de l’API REST.

Pour les définitions de schémas de message, les formats et la manière de les téléverser, voir Comprendre les schémas de messages.

Transformations intégrées vs. transformations WASM

Les graphiques de flux de données prennent en charge deux types de transformations :

  • Les transformations intégrées sont prédéfinies par Microsoft (mappage, filtre, branchement, concaténation, fenêtre, limitation de débit). Vous les configurez avec des règles. Aucun codage n’est requis.
  • Les transformations WASM sont des modules WebAssembly personnalisés que les développeurs créent et déploient. Utilisez-les quand vous avez besoin d’une logique que les transformations intégrées ne couvrent pas.

Les deux types de transformations s’exécutent dans la même DataflowGraph ressource, et vous pouvez les mélanger dans un seul pipeline. Pour plus d’informations sur la création et le déploiement de transformations personnalisées, consultez Utiliser des transformations WASM dans des graphiques de flux de données.

Gestion des erreurs dans les graphiques de flux de données

Lorsqu’une transformation rencontre une erreur lors du traitement d’un message (par exemple, un champ manquant ou une expression invalide), la transformation perd le message et enregistre une erreur. Le pipeline continue de traiter les messages suivants.

Causes courantes des erreurs de traitement :

  • Un champ référencé dans une règle n’existe inputs pas dans le message.
  • Une expression de filtre ou de branche retourne une valeur non booléenne.
  • Une expression fait référence à un type de données incompatible (par exemple, un objet JSON en arithmétique).
  • Un magasin d’état utilisé pour l’enrichissement est inaccessible.

Pour surveiller les erreurs de traitement, vérifiez les journaux des pods pour le graphe de flux de données ou utilisez les points de terminaison de métriques. Pour plus d’informations, consultez Configurer l’observabilité et la surveillance.

Limitation de mise à l’échelle des graphes à état

Important

Les transformations de fenêtre et de manette des gaz sont avec état. Chaque instance maintient son propre état et les instances ne partagent pas cet état entre elles. Lorsque le nombre d’instances du profil de flux de données est supérieur à un, les abonnements partagés répartissent les messages entre les instances, de sorte que chaque instance ne voit qu’un sous-ensemble des messages. Une transformée de fenêtre calcule ensuite des agrégations telles que les moyennes, les sommes et les comptages sur un jeu de données partiel, et une transformée de puissance impose la limite de débit configurée indépendamment dans chaque instance au lieu de couvrir l’ensemble du pipeline.

Réglez le nombre d’instances du profil de flux de données à 1 pour tout graphique de flux de données utilisant une transformation de fenêtre ou de régulation. Les graphes de flux de données sans état qui utilisent uniquement les transformations de cartographie, filtre, branchement et concaténation peuvent utiliser en toute sécurité des comptes d’instances plus élevés pour augmenter le débit.

Conseils de performance pour les graphiques de flux de données

Chaque transformation du pipeline ajoute une surcharge de traitement. Gardez ces instructions à l’esprit :

  • Préférez moins de transformations avec plus de règles. Si vous avez de nombreuses règles de transformation qui fonctionnent sur la même structure, placez-les dans une transformation de carte unique plutôt que de créer des transformations distinctes pour chaque règle.
  • Utilisez plusieurs transformations lorsque la logique est distincte. Les transformations distinctes sont logiques lorsque différentes étapes de traitement sont fondamentalement différentes (filtrage par rapport au mappage et à l’agrégation).
  • Conservez les règles associées ensemble. Une transformation de carte unique peut gérer le changement de nom, la restructuration, les champs calculés et les transformations de métadonnées en même temps.