Aggregera data med fönstertransformer i dataflödesgrafer

En fönstertransform grupperar inkommande meddelanden och producerar ett enda utdata med aggregerade värden när fönstret stängs. Istället för att vidarebefordra varje mätning individuellt kan du beräkna statistik som genomsnitt, minimivärden eller antal och skicka ett konsoliderat resultat nedströms.

För närvarande kan ett fönster stängas baserat på varaktighet, antal, minne eller triggervillkor.

En översikt över dataflödesdiagram och hur transformeringar består i en pipeline finns i Översikt över dataflödesdiagram.

Anmärkning

Icke-durationbaserad fönsterhantering kräver azureiotoperations/graph-dataflow-window:1.1.0 eller senare.

Transformer använder ett uttrycksspråk för att beräkna värden, testvillkor och referensfält. Uttryck hänvisar till indata efter position, inte namn: den första indatan i inputs listan är $1, den andra är $2, och så vidare. Inbyggda funktioner som cToF konvertera och manipulera dessa värden.

För den fullständiga listan över operatorer, funktioner, datatyper och metadatafält, se referensen Expressions.

Fönstertransformer lägger till aggregeringsfunktioner såsom average, min, och max, vilka endast finns tillgängliga i ackumuleringsregler. För hela listan, se Aggregeringsfunktioner.

Förutsättningar

  • En standardregisterslutpunkt med namnet default som pekar mcr.microsoft.com på skapas automatiskt under distributionen.

De Azure CLI exemplen i denna artikel använder miljövariabler så att du kan sätta varje värde en gång och sedan kopiera och klistra in kommandona as-is. Om du använder Azure IoT Operations Codespaces-miljön från quickstart är dessa variabler redan inställda för dig och du kan hoppa över detta steg. Annars, ställ in följande miljövariabler i ditt skal innan du kör kommandona.

Följande skript anger de mest använda miljövariablerna:

Miljövariabel Beskrivning
SUBSCRIPTION_ID ID:t för prenumerationen som innehåller din Azure IoT Operations-instans.
RESOURCE_GROUP Namnet på resursgruppen som innehåller din Azure IoT Operations-instans.
AIO_INSTANCE_NAME Namnet på din Azure IoT Operations-instans. För att lista dina instanser, kör az iot ops list -o table.
CLUSTER_NAME Namnet på det Azure Arc-aktiverade Kubernetes-klustret som hostar din instans.
LOCATION Azure-regionen att använda för nya resurser, till exempel eastus.
SUBSCRIPTION_ID=<subscription-id>
RESOURCE_GROUP=<resource-group-name>
AIO_INSTANCE_NAME=<instance-name>
CLUSTER_NAME=<cluster-name>
LOCATION=<region>

Du behöver bara ställa in de variabler som denna artikel använder. Den här artikeln kan använda ytterligare miljövariabler för resursnamn som du väljer. Artikeln förklarar hur man placerar dem där de introduceras.

Skalningsbegränsning för tillståndskänsliga grafer

Viktigt!

Fönster- och gasreglage transformerar är statele. Varje instans behåller sitt eget tillstånd och instanserna delar inte det tillståndet med varandra. När antalet dataflödesprofiler är fler än ett distribuerar delade prenumerationer meddelanden mellan instanserna, så varje instans ser endast en delmängd av meddelandena. En fönstertransform beräknar sedan aggregationer som medelvärden, summor och räkningar över en partiell dataset, och en throttle-transform upprätthåller den konfigurerade hastighetsgränsen oberoende i varje instans istället för över hela pipelinen.

Sätt antalet data flow-profiler till 1 för alla dataflödesgrafer som använder en fönster- eller throttle-transform. Tillståndslösa dataflödesgrafer som endast använder mapp-, filter-, gren- och koncatenationstransformer kan säkert använda högre instansantal för att öka genomströmningen.

När du ska använda en fönstertransformering

Använd en fönstertransformering när du tar emot högfrekventa sensordata och vill minska volymen innan du skickar den nedströms. Vanliga scenarier är:

  • Beräkningsmedelvärde: En temperatursensor publicerar varje sekund, men ditt molnprogram behöver bara ett genomsnitt på 30 sekunder.
  • Spåra extremvärden: Du vill ha lägsta och högsta tryckavläsningar under varje minuts intervall.
  • Antal händelser: Du behöver veta hur många dörröppnade händelser som inträffat under de senaste fem minuterna.
  • Skapa produktionsbatcher: Du vill beräkna statistik för varje sats med fast storlek, till exempel för varje 100 paket som kommer från en fyllningslinje.
  • Svara på tillståndsförändringar: Du vill veta när en operativ signal ändras, till exempel när en mixer ändras från running till draining.

Så här fungerar fönstertransformering

Fönstertransformningen har två interna steg anslutna i följd:

  1. Fönster: Buffrar meddelanden tills ett av de konfigurerade stängningsvillkoren aktiveras.
  2. Ackumulera: Tillämpar aggregeringsreglerna när fönstret stängs. Transformen reducerar alla meddelanden i fönstret till ett enda utdatameddelande.

Anmärkning

En fönstertransform måste konfigurera minst ett stängningsvillkor: delay, , countmemory, eller triggers.

Konfigurera fönsterstängningsvillkor

Från och med versionen 1.1.0, lägger fönstertransformen till tre konfigurationsnycklar vid sidan av den befintliga delay nyckeln:

Konfigurationsnyckel Fönstertyp Purpose
delay Varaktighetsbaserat fönster Stäng fönstret efter en fast tid.
count Räknebaserat fönster Stäng fönstret efter ett fast antal meddelanden.
memory Minnesbaserat fönster Stäng fönstret när den buffrade nyttolasten når en gräns.
triggers Triggerbaserat fönster Stäng fönstret när ett anpassat uttryck resulterar i true.

Varaktighetsbaserat fönster

Använd konfigurationen delay för att stänga fönstret efter en fast tid. Denna inställning styr hur länge varje tumbling-fönster varar.

Anmärkning

Fördröjningssteget justerar meddelandetidsstämplar till fönstergränser. Om ett meddelande anländer 7 sekunder in i ett 10-sekundersfönster tillhör det 10-sekundersgränsen.

Anmärkning

Om du inte tillhandahåller delayanvänder fönstret en 60-sekunders standardtimeout som säkerhetsventil.

I konfigurationen för fönstertransformering anger du varaktigheten för fönstret i sekunder. Ställ in det till 30 för ett rullande 30-sekundersfönster till exempel.

Fastighet Type Beskrivning
type snöre Måste vara "duration".
delaySeconds uint64 Antal sekunder innan fönstret stängs. Måste vara större än 0.

Räknebaserat fönster

Använd konfigurationen count för att stänga fönstret efter ett fast antal meddelanden.

I konfigurationen för fönstertransformeringen ställer du in antal meddelanden till 5 och anger beteendet för gränsmeddelanden till messageInCurrent.

Fastighet Type Beskrivning
type snöre Måste vara "messageCount".
maxMessageCount uint64 Antal meddelanden att buffra innan fönstret stängs. Måste vara större än 0.
boundaryMessage snöre Om meddelandet som stänger fönstret stannar kvar i det aktuella fönstret (messageInCurrent) eller startar nästa fönster (messageInNext).

Minnesbaserat fönster

Använd konfigurationen memory för att stänga fönstret när den buffrade nyttolaststorleken når en gräns.

I fönstertransformkonfigurationen, sätt bufferstorleken till 1048576 bytes och ställ in gränsmeddelandebeteendet till messageInNext.

Fastighet Type Beskrivning
type snöre Måste vara "bufferSize".
maxBufferBytes uint64 Maximalt sammanlagt antal byte nyttolast innan fönstret stängs. Måste vara större än 0.
boundaryMessage snöre Om meddelandet som stänger fönstret stannar kvar i det aktuella fönstret (messageInCurrent) eller startar nästa fönster (messageInNext).

Triggerbaserat fönster

Använd konfigurationen triggers när fönstret ska stängas baserat på meddelandeinnehåll eller körande tillstånd i det aktuella fönstret.

I fönstertransformkonfigurationen, lägg till en triggerregel med inmatningsfält temperature, uttryck running_sum($1) + $1 > 100, och gränsmeddelandebeteende messageInCurrent.

Fastighet Obligatoriskt Beskrivning
type Ja Måste vara "expression".
rules Ja Lista med utlösarregler. Reglerna utvärderas sekventiellt per meddelande; Den första matchningsregeln stänger fönstret.
datasets No Valfria datauppsättningar för tillståndslagret som refererar till tillståndslagret.

Varje triggerregel stöder dessa fält:

Fastighet Obligatoriskt Beskrivning
inputs Ja Lista med referenser till inmatningsfält. Uttrycket binder till $1, $2, och så vidare.
trigger Ja Booleskt uttryck som stänger fönstret när det utvärderar till true.
boundaryMessage Ja Om meddelandet som stänger fönstret stannar kvar i det aktuella fönstret (messageInCurrent) eller startar nästa fönster (messageInNext).

Fältet inputs stöder samma indatasyntax som används på andra ställen i dataflödesgrafer, inklusive vanliga fält, ?? standard, ? $last, $context(key).field, och $metadata.*. För mer information om att använda $context(key), se Berika med extern data.

Triggeruttryck kan använda de reguljära grafuttrycksfunktionerna och följande löpande tillståndsfunktioner som återställs när fönstret stängs:

Function Beskrivning
running_sum($1) Kumulativ summa av $1 mellan tidigare meddelanden i det aktuella fönstret.
running_avg($1) Ackumulerat genomsnitt av $1 över tidigare meddelanden.
running_min($1) Minsta värdet av $1 som setts i tidigare meddelanden. Returnerar $1 för det första meddelandet (minimumet för ett element är elementet självt).
running_max($1) Högsta värdet för $1 i tidigare meddelanden. Returnerar $1 på det första meddelandet (maxvärdet för ett enda element är elementet självt).
running_count($1) Antal meddelanden där $1 förekom.
running_count() Total meddelandeantal (inget fältfilter).
first($1) Första icke-tomma värdet av $1 i det nuvarande fönstret. Returnerar $1 i det första meddelandet.
changed($1) true om $1 skiljer sig från dess värde i föregående meddelande. false på det första meddelandet i ett fönster (inget tidigare värde att jämföra med).
prev($1) Senaste icke-tomma värdet av $1 från ett tidigare meddelande i det aktuella fönstret. Meddelanden där $1 var tomt ignoreras (det lagrade värdet skrivs inte över). Återvänder $1 vid det första meddelandet i ett fönster.

Anmärkning

running_sum($1) och liknande funktioner returnerar värden från tidigare bearbetade meddelanden. För det aktuella meddelandet, använd $1.

Exempel på triggerregler

Använd dessa exempel för att se gemensamma inputs mönster och trigger mönster i ett komplett triggerkonfigurationsobjekt:

  • Detta exempel visar ett vanligt triggeruttryck. Fönstret stängs när strömmen temperature är större än 80.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature"],
      "trigger": "$1 > 80",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Detta exempel visar ett triggeruttryck som använder running_sum($1) + $1 för att kombinera tidigare meddelanden i det aktuella fönstret med det aktuella meddelandet, och sedan stänga fönstret när tröskeln överskrids.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature"],
      "trigger": "running_sum($1) + $1 > 100",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Detta exempel visar nollsäker inmatningshantering med temperature ?? 0, plus messageInNext för att placera gränsmeddelandet i nästa fönster.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature ?? 0"],
      "trigger": "running_avg($1) > 80",
      "boundaryMessage": "messageInNext"
    }
  ]
}
  • Detta exempel visar en metadatabaserad trigger där fönstret stängs för ett specifikt ämnesvärde från $metadata.topic.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["$metadata.topic"],
      "trigger": "$1 == \"telemetry/high-priority\"",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Detta exempel visar triggerregler med dataset-berikning: den matchar meddelandet factoryId med en state-store-rad, läser shiftId från $context(factory).shiftId, och stänger fönstret när skiftvärdet ändras (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"
    }
  ]
}

I detta exempel förväntas tillståndslagringsdatasetet som representeras av factory innehålla fält som factoryId och shiftId.

Gränsbeteende

Inställningen boundaryMessage styr vad som händer med meddelandet som orsakade att ett räknings-, minnes- eller triggerbaserat fönster stängdes:

  • messageInCurrent: inkludera gränsmeddelandet i stängningsfönstret.
  • messageInNext: stäng först det aktuella fönstret och starta sedan nästa fönster med gränsmeddelandet.

Om messageInNext utlöses vid det första meddelandet i ett nytt fönster förhindras stängningen så att ett tomt fönster inte skickas.

Anmärkning

Varaktighetsbaserat fönster använder inte boundaryMessage. Varaktighetsgränser är tidsbaserade, inte meddelandebaserade, så det finns inget gränsmeddelande att placera i nuvarande eller nästa fönster.

Kombinera avslutningsvillkor

Du kan kombinera varaktighet, antal, minne och trigger-villkor i samma graf.

  • Varaktigheten är tidsdriven och utvärderas av timern.
  • För varje inkommande meddelande utvärderas meddelandedrivna villkor i denna ordning: Memory > Count > Trigger.
  • Inom triggers.rules, utvärderas reglerna sekventiellt och den första matchningsregeln vinner.

Utvärderingsordningen i meddelandedrivna villkor är viktig för ackumuleringsresultaten när ett meddelande uppfyller flera villkor samtidigt. Till exempel, om memory använder messageInCurrent och count använder messageInNext, följer ett meddelande som uppfyller båda villkoren minneskonfigurationen. Meddelandet stannar kvar i det aktuella fönstret och bidrar till det fönstrets ackumuleringsutdata.

Definiera ackumuleringsregler

Varje ackumuleringsregel anger hur du minskar ett fönster med meddelanden till ett enda utdatavärde. Konfigurationsnyckeln är rules.

I fönstertransformkonfigurationen, lägg till en ackumuleringsregel med inmatning temperature, utdata avgTemperatureoch aggregeringsfunktion average($1).

Fastighet Obligatoriskt Beskrivning
inputs Ja Lista över fältvägar som ska läsas från varje inkommande meddelande.
output Ja Fältsökväg för det aggregerade resultatet. Varje regel måste ha unika utdata.
expression Ja Formel som minskar indatavärden i fönstret till en enda skalär. Måste innehålla minst en aggregeringsfunktion.
description No Beskrivning som kan läsas av människor.

Till skillnad från kartregler expression är ett krav för varje ackumuleringsregel. Enbart användning $1 är inte giltigt eftersom det refererar till en samling värden, inte en enda skalär. Du måste omsluta den i en aggregeringsfunktion som average($1).

Sammansättningsfunktioner

Function Retur Beteende för tomma fönster
average Medelvärde av numeriska värden Error
sum Summa av numeriska värden 0,0
min Minsta numeriska värde Error
max Maximalt numeriskt värde Error
count Antal meddelanden där fältet finns 0
first Första värdet i fönstret Error
last Sista värdet i fönstret Error

Varje funktion tar en enskild positionsvariabel som argument ($1 för den första indatan, $2 för den andra och så vidare).

Icke-numeriska värden: funktionerna average, sum, minoch max hoppar tyst över icke-numeriska värden.

Närvarobaserade funktioner: count, firstoch last fungerar på fältnärvaro oavsett värdetyp.

Kombinera sammansättningar

Kombinera flera aggregeringsfunktioner i ett enda uttryck:

Lägg till en regel med indata temperature och humidity, och uttryck average($1) + max($2).

Om du vill konvertera ett aggregerat värde använder du konverteringsfunktionen utanför aggregeringen. Konverterar till exempel cToF(average($1)) medeltemperaturen till Fahrenheit.

Varje sammansättningsfunktion måste referera till en enskild positionsvariabel direkt. average($1) + max($2) är giltigt, men average($1 + $2) är inte det.

Skillnader från kartregler

Förmåga Mappningsregler Ackumuleringsregler
Uttryck krävs No Ja
Jokerteckenindata Stöds Stöds ej
$metadata Tillgång Stöds Stöds ej
$context Anrikning Stöds Stöds ej
? $last direktiv Stöds Stöds ej
Utdatainnehållstyp Motsvarar indata Alltid application/json

Fullständigt konfigurationsexempel

Detta exempel visar en komplett fönsterkonfiguration som stänger fönstret efter 30 sekunder, 5 meddelanden, 1 048 576 buffrade bytes, eller när running_sum($1) + $1 > 100. Exemplet sätter boundaryMessage värdet till messageInCurrent för de tre sista villkoren, och fönstret beräknar temperaturstatistik när det stängs.

Vilket villkor som stänger fönstret beror på tidpunkten för meddelandena, antalet meddelanden, nyttolastens storlek och innehållet. Följande exempel visar den resulterande utgången för varje stängningsvillkor.

Stängningstid

Om inget annat villkor utlöses först och tidsfönstret når 30 sekunder efter att ha tagit emot dessa tre meddelanden:

{ "temperature": 21.5 }
{ "temperature": 23.0 }
{ "temperature": 19.8 }

Utdatameddelandet är:

{
  "avgTemperature": 21.433333333333334,
  "minTemperature": 19.8,
  "maxTemperature": 23.0,
  "readingCount": 3,
  "tempRange": 3.2
}

Räkningen avslutas

Om fönstret tar emot dessa fem meddelanden innan något annat villkor aktiveras:

{ "temperature": 20.0 }
{ "temperature": 22.0 }
{ "temperature": 21.0 }
{ "temperature": 24.0 }
{ "temperature": 23.0 }

Utdatameddelandet är:

{
  "avgTemperature": 22.0,
  "minTemperature": 20.0,
  "maxTemperature": 24.0,
  "readingCount": 5,
  "tempRange": 4.0
}

Minnet stängs

Om den buffrade nyttolaststorleken når 1 048 576 byte innan något annat villkor aktiveras, till exempel efter dessa två stora meddelanden:

{ "temperature": 21.0, "payloadPad": "<large string>" }
{ "temperature": 22.5, "payloadPad": "<large string>" }

Utdatameddelandet är:

{
  "avgTemperature": 21.75,
  "minTemperature": 21.0,
  "maxTemperature": 22.5,
  "readingCount": 2,
  "tempRange": 1.5
}

Utlösaren sluter

Om triggeruttrycket running_sum($1) + $1 > 100 aktiveras före något annat villkor, till exempel efter dessa tre meddelanden:

{ "temperature": 40.0 }
{ "temperature": 35.0 }
{ "temperature": 30.0 }

Utdatameddelandet är:

{
  "avgTemperature": 35.0,
  "minTemperature": 30.0,
  "maxTemperature": 40.0,
  "readingCount": 3,
  "tempRange": 10.0
}

I driftmiljön skapar du ett dataflödesdiagram med en fönstertransformering:

  1. Lägg till en källa som läser från telemetry/temperature.
  2. Lägg till en fönstertransformering. Konfigurera ett 30-sekunders varaktighetsfönster, en gräns på 5 meddelanden, en buffertstorleksgräns på 1 048 576 byte och en triggerregel med temperature uttryck running_sum($1) + $1 > 100. För räknings-, minnes- och triggervillkoren, sätt gränsmeddelandets beteende till .messageInCurrent Lägg till ackumuleringsregler för genomsnitt, min, max, antal och intervall på fältet temperature.
  3. Lägg till ett mål som skickar till telemetry/aggregated.

Nästa steg