Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
Una trasformazione finestra raggruppa i messaggi in arrivo e produce un unico messaggio di output con valori aggregati quando la finestra si chiude. Invece di inoltrare ogni lettura singolarmente, puoi calcolare statistiche come medie, minimi o conteggi e inviare un risultato consolidato a valle.
Attualmente, una finestra può chiudersi in base alla durata, al conteggio, alla memoria o alle condizioni di trigger.
Per una panoramica dei grafici del flusso di dati e della composizione delle trasformazioni in una pipeline, vedere Panoramica dei grafici del flusso di dati.
Annotazioni
Il windowing non basato sulla durata richiede azureiotoperations/graph-dataflow-window:1.1.0 o versione successiva.
Le trasformazioni utilizzano un linguaggio di espressione per calcolare valori, condizioni di test e campi di riferimento. Le espressioni si riferiscono agli input per posizione, non per nome: il primo input nella inputs lista è $1, il secondo è $2, e così via. Funzioni integrate come cToF convertono e manipolano tali valori.
Per l'elenco completo di operatori, funzioni, tipi di dati e campi di metadati, consulta il riferimento Expressions.
Le trasformate delle finestre aggiungono funzioni di aggregazione come average, min, e max, che sono disponibili solo nelle regole di accumulazione. Per l'elenco completo, vedi Funzioni di aggregazione.
Prerequisiti
- Istanza di Operazioni di Azure IoT distribuita in un cluster Kubernetes. Per altre informazioni, vedere Deploy Operazioni di Azure IoT.
- Un endpoint del Registro di sistema predefinito denominato
defaultche punta amcr.microsoft.comviene creato automaticamente durante la distribuzione.
Gli esempi interfaccia della riga di comando di Azure in questo articolo usano variabili di ambiente così puoi impostare ogni valore una volta e poi copiare e incollare i comandi as-is. Se stai usando l'ambiente Operazioni di Azure IoT Codespaces dal quickstart, queste variabili sono già impostate per te e puoi saltare questo passaggio. Altrimenti, imposta le seguenti variabili di ambiente nella tua shell prima di eseguire i comandi.
I seguenti script impostano le variabili di ambiente più comunemente utilizzate:
| Variabile di ambiente | Descrizione |
|---|---|
SUBSCRIPTION_ID |
L'ID dell'abbonamento che contiene la tua istanza Operazioni di Azure IoT. |
RESOURCE_GROUP |
Il nome del gruppo di risorse che contiene la tua istanza Operazioni di Azure IoT. |
AIO_INSTANCE_NAME |
Il nome della tua istanza Operazioni di Azure IoT. Per elencare le tue istanze, esegui az iot ops list -o table. |
CLUSTER_NAME |
Il nome del cluster Kubernetes abilitato Azure Arc che ospita la tua istanza. |
LOCATION |
La regione Azure da utilizzare per nuove risorse, ad esempio eastus. |
SUBSCRIPTION_ID=<subscription-id>
RESOURCE_GROUP=<resource-group-name>
AIO_INSTANCE_NAME=<instance-name>
CLUSTER_NAME=<cluster-name>
LOCATION=<region>
Devi solo impostare le variabili utilizzate in questo articolo. Questo articolo potrebbe utilizzare variabili ambientali aggiuntive per i nomi delle risorse che scegli. L'articolo spiega come posizionarli dove vengono introdotti.
Limitazione del ridimensionamento per grafici con stato
Importante
Le trasformazioni di finestra e acceleratore sono stateful. Ogni istanza mantiene il proprio stato e le istanze non condividono quello stato tra loro. Quando il conteggio delle istanze del profilo di flusso dati è superiore a uno, gli abbonamenti condivisi distribuiscono i messaggi tra istanze, così che ogni istanza veda solo un sottoinsieme dei messaggi. Una trasformata finestra calcola quindi aggregazioni come medie, somme e conteggi su un dataset parziale, e una trasformata di gas impone il limite di velocità configurato in modo indipendente in ogni istanza invece che su tutta la pipeline.
Imposta il conteggio delle istanze del profilo di flusso dati a 1 per qualsiasi grafico di flusso dati che utilizza una trasformata a finestra o a manetta. I grafici di flusso dati senza stato che utilizzano solo trasformate mappa, filtro, ramo e concatenazione possono utilizzare in sicurezza conteggi di istanze più elevati per aumentare la produttività.
Quando usare una trasformazione di finestra
Usare una trasformazione a finestra quando si ricevono dati del sensore ad alta frequenza e si vuole ridurre il volume prima di inviarli a valle. Gli scenari comuni includono:
- Medie di calcolo: un sensore di temperatura pubblica ogni secondo, ma l'applicazione cloud richiede solo una media di 30 secondi.
- Tenere traccia degli estremi: si desidera che vengano rilevate le letture di pressione minima e massima in ogni intervallo di un minuto.
- Conteggio eventi: è necessario conoscere il numero di eventi di apertura delle porte negli ultimi cinque minuti.
- Crea batch di produzione: vuoi calcolare le statistiche per ogni batch a dimensione fissa, ad esempio ogni 100 pacchetti che escono da una linea di riempimento (fill line).
-
Rispondi ai cambiamenti di stato: vuoi sapere quando un segnale operativo cambia, ad esempio quando un mixer cambia da
runningadraining.
Funzionamento della trasformazione della finestra
La trasformazione della finestra ha due passaggi interni connessi in sequenza:
- Window: Memorizza i messaggi nel buffer fino a quando non si verifica una delle condizioni di chiusura configurate.
- Accumula: applica le regole di aggregazione alla chiusura della finestra. La trasformazione riduce tutti i messaggi nella finestra a un unico messaggio di uscita.
Annotazioni
Una trasformata di finestra deve configurare almeno una condizione di chiusura: delay, count, memory, o triggers.
Configurare le condizioni di chiusura delle finestre
A partire dalla versione 1.1.0, la trasformazione della finestra aggiunge tre tasti di configurazione insieme alla chiave esistente delay :
| Chiave di configurazione | Tipo di finestra | Purpose |
|---|---|---|
delay |
Finestra basata sulla durata | Chiudi la finestra dopo una durata fissa. |
count |
Finestra basata sul numero | Chiudi la finestra dopo un numero fisso di messaggi. |
memory |
Finestra basata sulla memoria | Chiudi la finestra quando la dimensione del carico utile in buffer raggiunge un limite. |
triggers |
Finestra attivata da trigger | Chiudi la finestra quando un'espressione personalizzata ha come risultato true. |
Finestra basata sulla durata
Usa la delay configurazione per chiudere la finestra dopo un periodo fisso. Questa impostazione controlla quanto dura ogni finestra di rotazione.
Annotazioni
Il passaggio di ritardo allinea i timestamp dei messaggi ai limiti della finestra. Se un messaggio arriva 7 secondi dopo l'inizio di una finestra di 10 secondi, appartiene al confine di 10 secondi.
Annotazioni
Se non fornisci delay, la finestra utilizza un timeout predefinito di 60 secondi come valvola di sicurezza.
Nella configurazione della trasformazione della finestra impostare la durata della finestra in secondi. Ad esempio, impostarlo su 30 per una finestra scorrevole di 30 secondi.
| Proprietà | Tipo | Descrizione |
|---|---|---|
type |
string | Deve essere "duration". |
delaySeconds |
uint64 | Numero di secondi prima che la finestra si chiuda. Deve essere maggiore di 0. |
Finestra basata sul numero
Usa la count configurazione per chiudere la finestra dopo un numero fisso di messaggi.
Nella configurazione della trasformazione della finestra, imposta il conteggio dei messaggi su 5 e imposta il comportamento dei messaggi di confine su messaggiInCurrent.
| Proprietà | Tipo | Descrizione |
|---|---|---|
type |
string | Deve essere "messageCount". |
maxMessageCount |
uint64 | Numero di messaggi da buffer prima che la finestra si chiuda. Deve essere maggiore di 0. |
boundaryMessage |
string | Se il messaggio che chiude la finestra rimane nella finestra corrente (messageInCurrent) o inizia la finestra successiva (messageInNext). |
Finestra basata sulla memoria
Usa la memory configurazione per chiudere la finestra quando la dimensione del payload bufferizzato raggiunge un limite.
Nella configurazione della trasformazione della finestra, imposta la dimensione del buffer a 1048576 byte e imposta il comportamento del messaggio di confine su messageInNext.
| Proprietà | Tipo | Descrizione |
|---|---|---|
type |
string | Deve essere "bufferSize". |
maxBufferBytes |
uint64 | Numero massimo cumulativo di byte del payload prima che la finestra si chiuda. Deve essere maggiore di 0. |
boundaryMessage |
string | Se il messaggio che chiude la finestra rimane nella finestra corrente (messageInCurrent) o inizia la finestra successiva (messageInNext). |
Finestra attivata da trigger
Usa la triggers configurazione in cui la finestra dovrebbe chiudersi in base al contenuto del messaggio o allo stato in esecuzione all'interno della finestra corrente.
Nella configurazione della trasformazione finestra, aggiungi una regola trigger con il campo di input temperature, l'espressione running_sum($1) + $1 > 100 e il comportamento del messaggio di confine messageInCurrent.
| Proprietà | Obbligatorio | Descrizione |
|---|---|---|
type |
Sì | Deve essere "expression". |
rules |
Sì | Matrice di regole di attivazione. Le regole vengono valutate in sequenza per ogni messaggio; la prima regola corrispondente chiude la finestra. |
datasets |
No | Dataset opzionali dello state store che fanno riferimento allo state store. |
Ogni regola di attivazione supporta questi campi:
| Proprietà | Obbligatorio | Descrizione |
|---|---|---|
inputs |
Sì | Array di riferimenti ai campi di input. L'espressione lega a $1, $2, e così via. |
trigger |
Sì | Espressione booleana che chiude la finestra quando valuta a true. |
boundaryMessage |
Sì | Se il messaggio che chiude la finestra rimane nella finestra corrente (messageInCurrent) o inizia la finestra successiva (messageInNext). |
Il inputs campo supporta la stessa sintassi di input usata altrove nei grafici di flusso di dati, inclusi campi semplici, ?? valori predefiniti, ? $last, $context(key).field, e $metadata.*. Per maggiori dettagli sull'utilizzo di $context(key), consulta Arricchimento con dati esterni.
Le espressioni trigger possono utilizzare le funzioni regolari di espressione grafica e le seguenti funzioni di stato di esecuzione che si resettano quando la finestra si chiude:
| Funzione | Descrizione |
|---|---|
running_sum($1) |
Somma cumulativa di $1 tra i messaggi precedenti nella finestra corrente. |
running_avg($1) |
Media cumulativa di $1 tra i messaggi precedenti. |
running_min($1) |
Valore minimo di $1 visto nei messaggi precedenti. Restituisce $1 per il primo messaggio (il minimo di un singolo elemento coincide con l'elemento stesso). |
running_max($1) |
Valore massimo di $1 visto nei messaggi precedenti. Ritorna $1 sul primo messaggio (il massimo di un elemento è se stesso). |
running_count($1) |
Conteggio dei messaggi in cui $1 era presente. |
running_count() |
Conteggio totale dei messaggi (senza filtro di campo). |
first($1) |
Primo valore non vuoto di $1 nella finestra corrente. Ritorna $1 al primo messaggio. |
changed($1) |
true se $1 differisce dal suo valore nel messaggio precedente.
false sul primo messaggio di una finestra (nessun valore precedente con cui confrontare). |
prev($1) |
Valore non vuoto più recente di $1 da un messaggio precedente nella finestra corrente. I messaggi in cui $1 era Vuoto vengono saltati (il valore memorizzato non viene sovrascritto). Ritorna $1 al primo messaggio di una finestra. |
Annotazioni
running_sum($1) e funzioni simili restituiscono valori da messaggi già elaborati. Per il messaggio attuale, usa $1.
Esempi di regole trigger
Usa questi esempi per vedere i comuni modelli inputs e trigger in un oggetto completo di configurazione di trigger:
- Questo esempio mostra una normale espressione di attivazione. La finestra si chiude quando la corrente
temperaturesupera 80.
{
"type": "expression",
"rules": [
{
"inputs": ["temperature"],
"trigger": "$1 > 80",
"boundaryMessage": "messageInCurrent"
}
]
}
- Questo esempio mostra un'espressione di trigger che utilizza
running_sum($1) + $1per combinare i messaggi precedenti nella finestra corrente con il messaggio corrente, quindi chiudere la finestra quando la soglia viene superata.
{
"type": "expression",
"rules": [
{
"inputs": ["temperature"],
"trigger": "running_sum($1) + $1 > 100",
"boundaryMessage": "messageInCurrent"
}
]
}
- Questo esempio mostra la gestione dell'input null-safe con
temperature ?? 0, piùmessageInNextper posizionare il messaggio di confine nella finestra successiva.
{
"type": "expression",
"rules": [
{
"inputs": ["temperature ?? 0"],
"trigger": "running_avg($1) > 80",
"boundaryMessage": "messageInNext"
}
]
}
- Questo esempio mostra un trigger basato su metadati in cui la finestra si chiude per un valore specifico dell'argomento da
$metadata.topic.
{
"type": "expression",
"rules": [
{
"inputs": ["$metadata.topic"],
"trigger": "$1 == \"telemetry/high-priority\"",
"boundaryMessage": "messageInCurrent"
}
]
}
- Questo esempio mostra regole trigger che utilizzano l'arricchimento del dataset: abbina il messaggio
factoryIda una riga state-store, leggeshiftIdda$context(factory).shiftId, e chiude la finestra quando il valore di shift cambia (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"
}
]
}
In questo esempio, il dataset di state store rappresentato da factory si prevede contenga campi come factoryId e shiftId.
Comportamento al confine
L'impostazione boundaryMessage controlla cosa accade al messaggio che ha causato la chiusura di una finestra basata sul conteggio, memoria o trigger:
-
messageInCurrent: include il messaggio di confine nella finestra di chiusura. -
messageInNext: chiudi prima la finestra corrente, poi avvia la finestra successiva con il messaggio di confine.
Se messageInNext spara sul primo messaggio di una nuova finestra, la chiusura viene soppressa in modo che non venga emessa una finestra vuota.
Annotazioni
La finestra basata sulla durata non usa boundaryMessage. I confini di durata sono basati sul tempo, non su messaggi, quindi non c'è un messaggio di confine da inserire nella finestra corrente o nella finestra successiva.
Combinare le condizioni di chiusura
Puoi combinare durata, conteggio, memoria e condizioni di trigger nello stesso grafico.
- La durata è determinata dal tempo e valutata dal timer.
- Per ogni messaggio in entrata, le condizioni guidate dai messaggi vengono valutate in questo ordine:
Memory > Count > Trigger. - All'interno di
triggers.rules, le regole vengono valutate in sequenza e prevale la prima regola corrispondente.
L'ordine di valutazione nelle condizioni guidate dal messaggio è importante per i risultati di accumulo quando un messaggio soddisfa più condizioni contemporaneamente. Ad esempio, se memory usa messageInCurrent e count usa messageInNext, un messaggio che soddisfa entrambe le condizioni segue la configurazione di memoria. Il messaggio rimane nella finestra corrente e contribuisce all'output cumulativo della finestra.
Definire le regole di accumulo
Ogni regola di accumulo specifica come ridurre una finestra di messaggi in un singolo valore di output. La chiave di configurazione è rules.
Nella configurazione della trasformazione di finestra, aggiungi una regola di accumulo con input temperature, output avgTemperature e funzione di aggregazione average($1).
| Proprietà | Obbligatorio | Descrizione |
|---|---|---|
inputs |
Sì | Elenco dei percorsi dei campi da leggere da ogni messaggio in arrivo. |
output |
Sì | Percorso del campo per il risultato aggregato. Ogni regola deve avere un output univoco. |
expression |
Sì | Formula che riduce i valori di input nella finestra a un singolo scalare. Deve contenere almeno una funzione di aggregazione. |
description |
No | Descrizione leggibile dagli umani |
A differenza delle regole della mappa, expression è necessario per ogni regola di accumulo. L'uso $1 da solo non è valido perché fa riferimento a una raccolta di valori, non a un singolo scalare. Devi incapsularlo in una funzione di aggregazione come average($1).
Funzione di aggregazione
| Funzione | Restituzioni | Comportamento della finestra vuota |
|---|---|---|
average |
Media dei valori numerici | Error |
sum |
Somma dei valori numerici | 0.0 |
min |
Valore numerico minimo | Error |
max |
Valore numerico massimo | Error |
count |
Numero di messaggi in cui esiste il campo | 0 |
first |
Primo valore nella finestra | Error |
last |
Ultimo valore nella finestra | Error |
Ogni funzione accetta una singola variabile posizionale come argomento ($1 per il primo input, $2 per il secondo e così via).
Valori non numerici: le averagefunzioni , summin, e max ignorano automaticamente i valori non numerici.
Funzioni basate sulla presenza: count, firste last operano sulla presenza del campo indipendentemente dal tipo di valore.
Combinare le aggregazioni
Combinare più funzioni di aggregazione in un'unica espressione:
Aggiungere una regola con gli input temperature e humidity, e l'espressione average($1) + max($2).
Per convertire un valore aggregato, applicare la funzione di conversione all'esterno dell'aggregazione. Ad esempio, cToF(average($1)) converte la temperatura media in Fahrenheit.
Ogni funzione di aggregazione deve fare riferimento direttamente a una singola variabile posizionale.
average($1) + max($2) è valido, ma average($1 + $2) non lo è.
Differenze rispetto alle regole della mappa
| Capability | Le regole della mappa | Regole di accumulo |
|---|---|---|
| Espressione obbligatoria | No | Sì |
| Input con caratteri jolly | Supportato | Non supportato |
$metadata Accesso |
Supportato | Non supportato |
$context Arricchimento |
Supportato | Non supportato |
? $last direttiva |
Supportato | Non supportato |
| Tipo di contenuto di output | Corrisponde all'input | Sempre application/json |
Esempio di configurazione completa
Questo esempio mostra una configurazione completa della finestra che chiude la finestra dopo 30 secondi, 5 messaggi, 1.048.576 byte bufferati, o quando running_sum($1) + $1 > 100. L'esempio imposta il boundaryMessage valore a messageInCurrent per le ultime tre condizioni, e la finestra calcola le statistiche di temperatura quando si chiude.
La condizione che chiude la finestra dipende dal tempismo del messaggio, dal conteggio, dalla dimensione del payload e dal contenuto. I seguenti esempi mostrano l'output risultante per ogni condizione di chiusura.
La durata termina
Se nessun'altra condizione si verifica prima e la finestra temporale raggiunge i 30 secondi dalla ricezione di questi tre messaggi:
{ "temperature": 21.5 }
{ "temperature": 23.0 }
{ "temperature": 19.8 }
Il messaggio di output è:
{
"avgTemperature": 21.433333333333334,
"minTemperature": 19.8,
"maxTemperature": 23.0,
"readingCount": 3,
"tempRange": 3.2
}
Chiusura del conteggio
Se la finestra riceve questi cinque messaggi prima che venga attivata qualsiasi altra condizione:
{ "temperature": 20.0 }
{ "temperature": 22.0 }
{ "temperature": 21.0 }
{ "temperature": 24.0 }
{ "temperature": 23.0 }
Il messaggio di output è:
{
"avgTemperature": 22.0,
"minTemperature": 20.0,
"maxTemperature": 24.0,
"readingCount": 5,
"tempRange": 4.0
}
Chiusura della memoria
Se la dimensione del payload bufferizzato raggiunge 1.048.576 byte prima che venga attivata qualsiasi altra condizione, ad esempio dopo questi due messaggi grandi:
{ "temperature": 21.0, "payloadPad": "<large string>" }
{ "temperature": 22.5, "payloadPad": "<large string>" }
Il messaggio di output è:
{
"avgTemperature": 21.75,
"minTemperature": 21.0,
"maxTemperature": 22.5,
"readingCount": 2,
"tempRange": 1.5
}
Chiusura del grilletto
Se l'espressione running_sum($1) + $1 > 100 del trigger si attiva prima di qualsiasi altra condizione, ad esempio dopo questi tre messaggi:
{ "temperature": 40.0 }
{ "temperature": 35.0 }
{ "temperature": 30.0 }
Il messaggio di output è:
{
"avgTemperature": 35.0,
"minTemperature": 30.0,
"maxTemperature": 40.0,
"readingCount": 3,
"tempRange": 10.0
}
Nell'ambiente Operazioni, crea un grafico del flusso di dati con una trasformazione di finestra:
- Aggiungere una sorgente che legga da
telemetry/temperature. - Aggiungi una trasformazione finestra. Configura una finestra della durata di 30 secondi, un limite di 5 messaggi, un limite di dimensione del buffer di 1.048.576 byte e una regola di trigger su
temperaturecon l'espressionerunning_sum($1) + $1 > 100. Per le condizioni di conteggio, memoria e trigger, imposta il comportamento del messaggio di confine amessageInCurrent. Aggiungi regole di accumulo per media, minimo, massimo, conteggio e range sultemperaturecampo. - Aggiungere una destinazione che invia a
telemetry/aggregated.