Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Une transformation de fenêtre regroupe les messages entrants et produit un seul message de sortie avec des valeurs agrégées lorsque la fenêtre se ferme. Au lieu de transférer chaque lecture individuellement, vous pouvez calculer des statistiques telles que des moyennes, des minimums ou des comptages et envoyer un résultat consolidé en aval.
Actuellement, une fenêtre peut se fermer en fonction de la durée, du compte, de la mémoire ou des conditions de déclenchement.
Pour obtenir une vue d’ensemble des graphiques de flux de données et la façon dont les transformations composent dans un pipeline, consultez vue d’ensemble des graphiques de flux de données.
Note
Le fenêtrage non fondé sur la durée nécessite la version azureiotoperations/graph-dataflow-window:1.1.0 ou ultérieure.
Les transformations utilisent un langage d’expressions pour calculer les valeurs, les conditions de test et les champs de référence. Les expressions désignent les entrées par position, et non par le nom : la première entrée de la inputs liste est $1, la seconde est $2, et ainsi de suite. Des fonctions intégrées telles que cToF convertissent et manipulent ces valeurs.
Pour la liste complète des opérateurs, fonctions, types de données et champs de métadonnées, voir la référence Expressions.
Les transformations de fenêtre ajoutent des fonctions d’agrégation telles que average, min, et max, qui ne sont disponibles que dans les règles d’accumulation. Pour la liste complète, voir Fonctions d’agrégation.
Prerequisites
- Instance de Opérations Azure IoT déployée dans un cluster Kubernetes. Pour plus d’informations, consultez Deploy Opérations Azure IoT.
- Un point de terminaison de registre par défaut nommé
defaultqui pointe versmcr.microsoft.comest créé automatiquement pendant le déploiement.
Les exemples Azure CLI de cet article utilisent des variables d’environnement afin de pouvoir définir chaque valeur une fois puis copier-coller les commandes as-is. Si vous utilisez l'environnement Opérations Azure IoT Codespaces du quickstart, ces variables sont déjà définies pour vous et vous pouvez sauter cette étape. Sinon, définissez les variables d’environnement suivantes dans votre shell avant d’exécuter les commandes.
Les scripts suivants définissent les variables d’environnement les plus couramment utilisées :
| Variable d'environnement | Description |
|---|---|
SUBSCRIPTION_ID |
L’identifiant de l’abonnement contenant votre instance Opérations Azure IoT. |
RESOURCE_GROUP |
Le nom du groupe de ressources contenant votre instance Opérations Azure IoT. |
AIO_INSTANCE_NAME |
Le nom de votre instance Opérations Azure IoT. Pour lister vos instances, exécutez az iot ops list -o table. |
CLUSTER_NAME |
Le nom du cluster Kubernetes compatible Azure Arc qui héberge votre instance. |
LOCATION |
La région Azure à utiliser pour de nouvelles ressources, par exemple eastus. |
SUBSCRIPTION_ID=<subscription-id>
RESOURCE_GROUP=<resource-group-name>
AIO_INSTANCE_NAME=<instance-name>
CLUSTER_NAME=<cluster-name>
LOCATION=<region>
Vous n’avez qu’à définir les variables utilisées dans cet article. Cet article peut utiliser des variables d’environnement supplémentaires pour les noms de ressources que vous choisissez. L’article explique comment les placer là où ils sont introduits.
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.
Quand utiliser une transformation de fenêtre
Utilisez une transformation de fenêtre lorsque vous recevez des données de capteur à haute fréquence et souhaitez réduire le volume avant de l’envoyer en aval. Les scénarios courants sont les suivants :
- Moyennes de calcul : un capteur de température publie toutes les secondes, mais votre application cloud a uniquement besoin d’une moyenne de 30 secondes.
- Suivre les extrêmes : vous souhaitez obtenir les mesures minimales et maximales de pression à chaque intervalle d'une minute.
- Compter les événements : vous devez savoir combien d'événements d'ouverture de porte se sont produits au cours des cinq dernières minutes.
- Créez des lots de production : Vous voulez calculer des statistiques pour chaque lot de taille fixe, par exemple tous les 100 paquets provenant d’une ligne de remplissage.
-
Répondre aux changements d’état : Vous souhaitez savoir chaque fois qu’un signal de fonctionnement change, par exemple lorsqu’un mélangeur passe de
runningàdraining.
Fonctionnement du transformateur de fenêtre
La transformation de fenêtre comporte deux étapes internes connectées en séquence :
- Fenêtre : Met les messages en mémoire tampon jusqu’à ce qu’une des conditions de fermeture configurées se déclenche.
- Accumulate: applique vos règles d’agrégation lorsque la fenêtre se ferme. La transformation réduit tous les messages de la fenêtre à un seul message de sortie.
Note
Une transformation de fenêtre doit configurer au moins une condition de fermeture : delay, count, memory, ou triggers.
Configurer les conditions de fermeture des fenêtres
À partir de la version 1.1.0, la transformation de fenêtre ajoute trois clés de configuration en plus de la clé existante delay :
| Clé de configuration | Type de fenêtre | Purpose |
|---|---|---|
delay |
Fenêtre basée sur la durée | Fermez la fenêtre après une durée fixe. |
count |
Fenêtre basée sur le comptage | Fermez la fenêtre après un nombre fixe de messages. |
memory |
Fenêtre basée sur la mémoire | Fermez la fenêtre lorsque la taille de la charge utile en mémoire tampon atteint une limite. |
triggers |
Fenêtre basée sur les déclencheurs | Fermez la fenêtre lorsqu’une expression personnalisée s’évalue à true. |
Fenêtre basée sur la durée
Utilisez cette delay configuration pour fermer la fenêtre après une durée fixe. Ce réglage détermine la durée de chaque fenêtre de roulement.
Note
L’étape de délai aligne les horodatages des messages sur les limites de fenêtre. Si un message arrive 7 secondes après le début d’une fenêtre de 10 secondes, il appartient à la limite de 10 secondes.
Note
Si vous ne fournissez delaypas , la fenêtre utilise un délai d’attente par défaut de 60 secondes comme soupape de sécurité.
- Expérience des opérations
- Service de contrôle d’accès Azure (CLI)
- Bicep
- Kubernetes (débogage uniquement)
Dans la configuration de la transformation de fenêtre, définissez la durée de la fenêtre en secondes. Par exemple, définissez-le sur 30 pour une fenêtre glissante de 30 secondes.
| Propriété | Type | Description |
|---|---|---|
type |
ficelle | Doit être "duration". |
delaySeconds |
uint64 | Nombre de secondes avant la fermeture de la fenêtre. Doit être supérieure à 0. |
Fenêtre basée sur le comptage
Utilisez la count configuration pour fermer la fenêtre après un nombre fixe de messages.
- Expérience des opérations
- Service de contrôle d’accès Azure (CLI)
- Bicep
- Kubernetes (débogage uniquement)
Dans la configuration de transformation de fenêtre, définissez le nombre de messages à 5 et le comportement des messages de frontière sur messageInCurrent.
| Propriété | Type | Description |
|---|---|---|
type |
ficelle | Doit être "messageCount". |
maxMessageCount |
uint64 | Nombre de messages à mettre en mémoire tampon avant la fermeture de la fenêtre. Doit être supérieure à 0. |
boundaryMessage |
ficelle | Que le message qui ferme la fenêtre reste dans la fenêtre courante (messageInCurrent) ou que la fenêtre suivante commence (messageInNext). |
Fenêtre basée sur la mémoire
Utilisez la memory configuration pour fermer la fenêtre lorsque la taille de la charge utile en mémoire tampon atteint une limite.
- Expérience des opérations
- Service de contrôle d’accès Azure (CLI)
- Bicep
- Kubernetes (débogage uniquement)
Dans la configuration de transformation de fenêtre, fixez la taille du tampon à 1048576 octets et le comportement du message frontière sur messageInNext.
| Propriété | Type | Description |
|---|---|---|
type |
ficelle | Doit être "bufferSize". |
maxBufferBytes |
uint64 | Nombre maximal cumulé d’octets de charge utile avant la fermeture de la fenêtre. Doit être supérieure à 0. |
boundaryMessage |
ficelle | Que le message qui ferme la fenêtre reste dans la fenêtre courante (messageInCurrent) ou que la fenêtre suivante commence (messageInNext). |
Fenêtre basée sur les déclencheurs
Utilisez la triggers configuration lorsque la fenêtre devrait se fermer en fonction du contenu du message ou de l’état d’exécution à l’intérieur de la fenêtre actuelle.
- Expérience des opérations
- Service de contrôle d’accès Azure (CLI)
- Bicep
- Kubernetes (débogage uniquement)
Dans la configuration de transformation de fenêtre, ajoutez une règle de déclenchement avec le champ d’entrée temperature, l’expression running_sum($1) + $1 > 100 et le comportement du message de frontière messageInCurrent.
| Propriété | Obligatoire | Description |
|---|---|---|
type |
Oui | Doit être "expression". |
rules |
Oui | Une série de règles de déclenchement. Les règles sont évaluées séquentiellement par message ; La première règle de correspondance ferme la fenêtre. |
datasets |
Non | Des ensembles de données de magasin d’état facultatifs qui référencent le magasin d’état. |
Chaque règle de déclenchement prend en charge ces champs :
| Propriété | Obligatoire | Description |
|---|---|---|
inputs |
Oui | Tableau de références de champs d’entrée. L’expression lie à $1, $2, et ainsi de suite. |
trigger |
Oui | Expression booléenne qui ferme la fenêtre lorsqu’elle évalue à true. |
boundaryMessage |
Oui | Que le message qui ferme la fenêtre reste dans la fenêtre courante (messageInCurrent) ou que la fenêtre suivante commence (messageInNext). |
Le inputs champ prend en charge la même syntaxe d’entrée utilisée ailleurs dans les graphiques de flux de données, incluant les champs simples, ?? les valeurs par défaut, ? $last, $context(key).field, et $metadata.*. Pour plus de détails sur l’utilisation de $context(key), voir Enrichir avec des données externes.
Les expressions trigger peuvent utiliser les fonctions d’expression graphiques régulières et les fonctions d’état d’exécution suivantes qui se réinitialisent lorsque la fenêtre se ferme :
| Fonction | Description |
|---|---|
running_sum($1) |
Somme cumulative de $1 entre les messages précédents dans la fenêtre actuelle. |
running_avg($1) |
Moyenne cumulée de $1 sur les messages précédents. |
running_min($1) |
Valeur minimale de $1 vue dans les messages précédents. Renvoie $1 au premier message (le minimum d’un seul élément est l’élément lui-même). |
running_max($1) |
Valeur maximale de $1 vue dans les messages précédents. Renvoie $1 au premier message (le maximum d’un seul élément est l’élément lui-même). |
running_count($1) |
Compte des messages où $1 il était présent. |
running_count() |
Nombre total de messages (sans filtre de champ). |
first($1) |
Première valeur non vide de $1 dans la fenêtre actuelle. Revient à $1 au premier message. |
changed($1) |
true si $1 diffère de sa valeur dans le message précédent.
false sur le premier message d’une fenêtre (aucune valeur précédente à comparer). |
prev($1) |
La valeur non vide la plus récente de $1 issue d’un message précédent dans la fenêtre courante. Les messages où $1 était vide sont ignorés (la valeur stockée n’est pas écrasée). Retourne $1 dans le premier message d’une fenêtre. |
Note
running_sum($1) et des fonctions similaires renvoient des valeurs issues de messages déjà traités. Pour le message actuel, utilisez $1.
Exemples de règles déclencheuses
Utilisez ces exemples pour voir les modèles courants de inputs et de trigger dans un objet complet de configuration de déclencheur :
- Cet exemple montre une expression de déclenchement régulière. La fenêtre se ferme lorsque le courant
temperaturedépasse 80.
{
"type": "expression",
"rules": [
{
"inputs": ["temperature"],
"trigger": "$1 > 80",
"boundaryMessage": "messageInCurrent"
}
]
}
- Cet exemple montre une expression déclencheur qui permet
running_sum($1) + $1de combiner les messages précédents de la fenêtre courante avec le message actuel, puis de fermer la fenêtre lorsque le seuil est dépassé.
{
"type": "expression",
"rules": [
{
"inputs": ["temperature"],
"trigger": "running_sum($1) + $1 > 100",
"boundaryMessage": "messageInCurrent"
}
]
}
- Cet exemple montre la gestion des entrées null-safe avec
temperature ?? 0, plusmessageInNextpour placer le message de frontière dans la fenêtre suivante.
{
"type": "expression",
"rules": [
{
"inputs": ["temperature ?? 0"],
"trigger": "running_avg($1) > 80",
"boundaryMessage": "messageInNext"
}
]
}
- Cet exemple montre un déclencheur basé sur les métadonnées où la fenêtre se ferme pour une valeur de sujet spécifique de
$metadata.topic.
{
"type": "expression",
"rules": [
{
"inputs": ["$metadata.topic"],
"trigger": "$1 == \"telemetry/high-priority\"",
"boundaryMessage": "messageInCurrent"
}
]
}
- Cet exemple montre des règles déclencheuses utilisant l’enrichissement du jeu de données : il associe le message
factoryIdà une ligne de magasin d’états, litshiftIddepuis$context(factory).shiftId, et ferme la fenêtre lorsque cette valeur de décalage change (changed($1)).
{
"type": "expression",
"datasets": [
{
"key": "factory",
"inputs": ["$source.factoryId", "$context.factoryId"],
"expression": "$1 == $2"
}
],
"rules": [
{
"inputs": ["$context(factory).shiftId"],
"trigger": "changed($1)",
"boundaryMessage": "messageInCurrent"
}
]
}
Dans cet exemple, le jeu de données du magasin d’états représenté par factory est censé contenir des champs comme factoryId et shiftId.
Comportement des frontières
Le boundaryMessage paramètre contrôle ce qui arrive au message qui a provoqué la fermeture d’une fenêtre basée sur le comptage, la mémoire ou le déclencheur :
-
messageInCurrent: inclure le message frontière dans la fenêtre de fermeture. -
messageInNext: fermer d’abord la fenêtre courante, puis commencer la fenêtre suivante avec le message de bordure.
Si messageInNext se déclenche sur le premier message dans une nouvelle fenêtre, la fermeture est empêchée afin qu’une fenêtre vide ne soit pas produite.
Note
La fenêtre basée sur la durée n’utilise boundaryMessagepas . Les limites de durée sont basées sur le temps, pas sur le message, donc il n’y a pas de message de frontière à placer dans la fenêtre actuelle ou la fenêtre suivante.
Combiner les conditions de fermeture
Vous pouvez combiner la durée, le compte, la mémoire et les conditions de déclenchement dans le même graphique.
- La durée est déterminée dans le temps et évaluée par le minuteur.
- Pour chaque message entrant, les conditions pilotées par message sont évaluées dans cet ordre :
Memory > Count > Trigger. - Dans
triggers.rules, les règles sont évaluées séquentiellement et la première règle correspondante l’emporte.
L’ordre d’évaluation dans les conditions pilotées par un message est important pour les résultats d’accumulation lorsqu’un message satisfait plusieurs conditions simultanément. Par exemple, si memory utilise messageInCurrent et count utilise messageInNext, un message qui satisfait les deux conditions suit la configuration mémoire. Le message reste dans la fenêtre courante et contribue à la sortie cumulée de cette fenêtre.
Définir des règles d’accumulation
Chaque règle d’accumulation spécifie comment réduire une fenêtre de messages en une seule valeur de sortie. La clé de configuration est rules.
- Expérience des opérations
- Service de contrôle d’accès Azure (CLI)
- Bicep
- Kubernetes (débogage uniquement)
Dans la configuration de transformation de fenêtre, on ajoute une règle d’accumulation avec entrée temperature, sortie avgTemperature, et fonction average($1)d’agrégation .
| Propriété | Obligatoire | Description |
|---|---|---|
inputs |
Oui | Liste des chemins de champ à lire à partir de chaque message entrant. |
output |
Oui | Chemin du champ pour le résultat agrégé. Chaque règle doit avoir une sortie unique. |
expression |
Oui | Formule qui réduit les valeurs d’entrée dans la fenêtre à un seul scalaire. Doit contenir au moins une fonction d’agrégation. |
description |
Non | Description lisible par l’homme. |
Contrairement aux règles de mappage, expressionest requis pour chaque règle d’accumulation. L’utilisation $1 seule n’est pas valide, car elle référence une collection de valeurs, pas une seule scalaire. Vous devez l’encapsuler dans une fonction d’agrégation comme average($1).
Fonctions d’agrégation
| Fonction | Retours | Comportement de fenêtre vide |
|---|---|---|
average |
Moyenne des valeurs numériques | Error |
sum |
Somme des valeurs numériques | 0,0 |
min |
Valeur numérique minimale | Error |
max |
Valeur numérique maximale | Error |
count |
Nombre de messages où le champ existe | 0 |
first |
Première valeur dans la fenêtre | Error |
last |
Dernière valeur dans la fenêtre | Error |
Chaque fonction prend une variable positionnelle unique comme argument ($1 pour la première entrée, $2 pour la seconde, etc.).
Valeurs non numériques : les fonctions average, sum, min et max ignorent silencieusement les valeurs non numériques.
Fonctions basées sur la présence : count, firstet last fonctionnent sur la présence de champ indépendamment du type valeur.
Combiner des agrégations
Combinez plusieurs fonctions d’agrégation en une seule expression :
- Expérience des opérations
- Service de contrôle d’accès Azure (CLI)
- Bicep
- Kubernetes (débogage uniquement)
Ajoutez une règle avec des entrées temperature et humidity, et une expression average($1) + max($2).
Pour convertir une valeur agrégée, appliquez la fonction de conversion en dehors de l’agrégation. Par exemple, cToF(average($1)) convertit la température moyenne en Fahrenheit.
Chaque fonction d’agrégation doit référencer directement une variable positionnelle unique.
average($1) + max($2) est valide, mais average($1 + $2) ce n’est pas le cas.
Différences par rapport aux règles de carte
| Capacité | Règles de mappage | Règles d’accumulation |
|---|---|---|
| Expression requise | Non | Oui |
| Entrées avec caractères génériques | Soutenu | Non pris en charge |
$metadata Accès |
Soutenu | Non pris en charge |
$context Enrichissement |
Soutenu | Non pris en charge |
? $last Directive |
Soutenu | Non pris en charge |
| Type de contenu de sortie | Correspond à l'entrée | Toujours application/json |
Exemple de configuration complète
Cet exemple montre une configuration complète de fenêtre qui ferme la fenêtre après 30 secondes, 5 messages, 1 048 576 octets en mémoire tampon, ou lorsque running_sum($1) + $1 > 100. L’exemple fixe la boundaryMessage valeur à messageInCurrent pour les trois dernières conditions, et la fenêtre calcule les statistiques de température à sa fermeture.
La condition qui ferme la fenêtre dépend du timing du message, du comptage, de la taille de la charge utile et du contenu. Les exemples suivants montrent la sortie résultante pour chaque condition de fermeture.
Fin de la durée
Si aucune autre condition ne se déclenche en premier et que la fenêtre atteint 30 secondes après avoir reçu ces trois messages :
{ "temperature": 21.5 }
{ "temperature": 23.0 }
{ "temperature": 19.8 }
Le message de sortie est :
{
"avgTemperature": 21.433333333333334,
"minTemperature": 19.8,
"maxTemperature": 23.0,
"readingCount": 3,
"tempRange": 3.2
}
Fin du décompte
Si la fenêtre reçoit ces cinq messages avant que toute autre condition ne soit déclenchée :
{ "temperature": 20.0 }
{ "temperature": 22.0 }
{ "temperature": 21.0 }
{ "temperature": 24.0 }
{ "temperature": 23.0 }
Le message de sortie est :
{
"avgTemperature": 22.0,
"minTemperature": 20.0,
"maxTemperature": 24.0,
"readingCount": 5,
"tempRange": 4.0
}
La mémoire se ferme
Si la taille de la charge utile en mémoire tampon atteint 1 048 576 octets avant qu’une autre condition ne se déclenche, par exemple après ces deux gros messages :
{ "temperature": 21.0, "payloadPad": "<large string>" }
{ "temperature": 22.5, "payloadPad": "<large string>" }
Le message de sortie est :
{
"avgTemperature": 21.75,
"minTemperature": 21.0,
"maxTemperature": 22.5,
"readingCount": 2,
"tempRange": 1.5
}
Le contact se ferme
Si l’expression running_sum($1) + $1 > 100 déclencheuse se déclenche avant toute autre condition, par exemple après ces trois messages :
{ "temperature": 40.0 }
{ "temperature": 35.0 }
{ "temperature": 30.0 }
Le message de sortie est :
{
"avgTemperature": 35.0,
"minTemperature": 30.0,
"maxTemperature": 40.0,
"readingCount": 3,
"tempRange": 10.0
}
- Expérience des opérations
- Service de contrôle d’accès Azure (CLI)
- Bicep
- Kubernetes (débogage uniquement)
Dans l’expérience Opérations, créez un graphique de flux de données avec une transformation de fenêtre :
- Ajoutez une source qui lit à partir de
telemetry/temperature. - Ajouter une transformation de fenêtre. Configurez une fenêtre d’une durée de 30 secondes, une limite de 5 messages, une limite de taille de tampon de 1 048 576 octets et une règle de déclenchement sur
temperatureavec l’expressionrunning_sum($1) + $1 > 100. Pour les conditions de comptage, mémoire et déclenchement, réglez le comportement du message frontière àmessageInCurrent. Ajoutez des règles d’accumulation pour la moyenne, le minimum, le maximum, le nombre et la portée sur letemperatureterrain. - Ajoutez une destination qui envoie vers
telemetry/aggregated.