Agregacja danych za pomocą transformacji okiennych w grafach przepływu danych

Transformacja okna grupuje przychodzące wiadomości i generuje pojedynczą wiadomość wyjściową z zagregowanymi wartościami po zamknięciu okna. Zamiast przekazywać każde odczyty osobno, możesz obliczyć statystyki takie jak średnie, minimum czy liczby i przesłać jeden skonsolidowany wynik dalej.

Obecnie okno może się zamknąć w zależności od czasu trwania, liczby, pamięci lub warunków wyzwalających.

Aby zapoznać się z omówieniem wykresów przepływu danych i sposobem tworzenia przekształceń w potoku, zobacz Omówienie wykresów przepływu danych.

Uwaga / Notatka

Okna niewykorzystujące czasu trwania wymagają wersji azureiotoperations/graph-dataflow-window:1.1.0 lub nowszej.

Transformacje wykorzystują język wyrażeń do obliczania wartości, warunków testowych i pól odniesienia. Wyrażenia odnoszą się do wejść według położenia, a nie nazwy: pierwsze wejście w inputs liście to $1, drugie to $2, i tak dalej. Wbudowane funkcje, takie jak cToF, konwertują i modyfikują te wartości.

Pełną listę operatorów, funkcji, typów danych i pól metadanych można znaleźć w sekcji Expressions.

Transformacje okna udostępniają funkcje agregujące, takie jak average, min i max, które są dostępne tylko w regułach akumulacji. Pełną listę można znaleźć w artykule Funkcje agregacji.

Wymagania wstępne

  • Instancja usługi Operacje Azure IoT wdrożona w klastrze Kubernetes. Aby uzyskać więcej informacji, zobacz Deploy Operacje Azure IoT.
  • Domyślny punkt końcowy rejestru o nazwie default, który wskazuje na mcr.microsoft.com, jest tworzony automatycznie podczas wdrażania.

Azure CLI przykładów w tym artykule używa zmiennych środowiskowych, dzięki czemu można ustawić każdą wartość raz, a następnie skopiować i wkleić polecenia as-is. Jeśli korzystasz ze środowiska Operacje Azure IoT Codespaces z quickstartu, te zmienne są już ustawione i możesz pominąć ten krok. W przeciwnym razie ustaw następujące zmienne środowiskowe w swojej poskoczce przed uruchomieniem poleceń.

Poniższe skrypty określają najczęściej używane zmienne środowiskowe:

Zmienna środowiskowa Opis
SUBSCRIPTION_ID ID subskrypcji zawierającej Twoją instancję Operacje Azure IoT.
RESOURCE_GROUP Nazwa grupy zasobów zawierającej Twoją instancję Operacje Azure IoT.
AIO_INSTANCE_NAME Nazwa Twojej instancji Operacje Azure IoT. Aby wypisać swoje instancje, uruchom az iot ops list -o table.
CLUSTER_NAME Nazwa klastra Kubernetes z włączonym Azure Arc, który hostuje Twoją instancję.
LOCATION Region Azure do wykorzystania dla nowych zasobów, na przykład eastus.
SUBSCRIPTION_ID=<subscription-id>
RESOURCE_GROUP=<resource-group-name>
AIO_INSTANCE_NAME=<instance-name>
CLUSTER_NAME=<cluster-name>
LOCATION=<region>

Wystarczy ustawić zmienne, których używa ten artykuł. W tym artykule może być użytych dodatkowych zmiennych środowiskowych do wyboru nazw zasobów. Artykuł wyjaśnia, jak ustawić je tam, gdzie są wprowadzane.

Ograniczenie skalowania dla wykresów stanowych

Ważna

Transformacje okien i przepustnicy są stanowe. Każda instancja utrzymuje swój własny stan i instancje nie dzielą tego stanu ze sobą. Gdy liczba instancji profilu przepływu danych jest większa niż jedna, subskrypcje współdzielone rozprowadzają wiadomości między instancjami, tak że każda instancja widzi tylko podzbiór tych komunikatów. Transformacja okna oblicza agregacje takie jak średnie, sumy i liczby na częściowym zbiorze danych, a transformacja throttle egzekwuje skonfigurowany limit szybkości niezależnie w każdym przypadku, zamiast w całym potoku.

Ustaw liczbę instancji profilu przepływu danych na 1 dla każdego wykresu przepływu danych wykorzystującego transformację okna lub przepustzenia. Bezstanowe grafy przepływu danych, które wykorzystują jedynie transformacje mapowania, filtrowania, rozgałęzienia i konkatenacji, mogą bezpiecznie wykorzystywać większą liczbę instancji do zwiększenia przepustowości.

Kiedy należy użyć przekształcenia okna

Użyj przekształcenia okna podczas odbierania danych z czujnika o wysokiej częstotliwości, aby zmniejszyć ilość danych przed ich wysłaniem dalej. Typowe scenariusze obejmują:

  • Średnie obliczeniowe: czujnik temperatury publikuje co sekundę, ale aplikacja w chmurze potrzebuje tylko średniej 30 sekund.
  • Śledź skrajności: potrzebujesz minimalnych i maksymalnych odczytów ciśnienia w każdym interwale jednominutowym.
  • Liczba zdarzeń: Musisz wiedzieć, ile zdarzeń otwarcia drzwi wystąpiło w ciągu ostatnich pięciu minut.
  • Twórz partie produkcyjne: Chcesz obliczyć statystyki dla każdej partii o stałym rozmiarze, na przykład co 100 paczek wychodzących z linii napełniania.
  • Reaguj na zmiany stanu: Chcesz wiedzieć, kiedy sygnał operacyjny się zmienia, na przykład gdy mikser zmienia się z running na draining.

Jak działa transformacja okna

Transformacja okna ma dwa kroki wewnętrzne połączone sekwencyjnie.

  1. Okno: Buforuje komunikaty, aż jeden ze skonfigurowanych warunków zamknięcia okna zostanie spełniony.
  2. Gromadzenie: Stosuje reguły agregacji po zamknięciu okna. Transformacja redukuje wszystkie komunikaty w oknie do jednej wiadomości wyjściowej.

Uwaga / Notatka

Transformacja okna musi skonfigurować co najmniej jeden warunek zamknięcia: delay, count, memory, lub triggers.

Konfiguruj warunki zamknięcia okna

Zaczynając od wersji 1.1.0, transformacja okna dodaje trzy klucze konfiguracyjne obok istniejącego delay klucza:

Klucz konfiguracji Typ okna Purpose
delay Okno oparte na czasie trwania Zamknij okno po ustalonym czasie trwania.
count Okno zliczające Zamknij okno po ustalonej liczbie wiadomości.
memory Okno oparte na pamięci Zamknij okno, gdy rozmiar ładunku buforowanego osiągnie limit.
triggers Okno wyzwalane zdarzeniem Zamknij okno, gdy niestandardowe wyrażenie przyjmie wartość true.

Okno oparte na czasie trwania

Użyj delay tej konfiguracji, aby zamknąć okno po określonym czasie. To ustawienie kontroluje, jak długo trwa każde okno przewracania.

Uwaga / Notatka

Krok opóźnienia dostosowuje znaczniki czasu komunikatu do granic okien. Jeśli wiadomość nadejdzie po 7 sekundach od początku 10-sekundowego okna, należy do granicy 10 sekund.

Uwaga / Notatka

Jeśli nie delaypodasz , okno używa domyślnego 60-sekundowego limitu czasu jako zaworu bezpieczeństwa.

W konfiguracji przekształcania okna ustaw czas trwania okna w sekundach. Na przykład ustaw dla 30 30-sekundowego okna przesuwnego.

Majątek Typ Opis
type ciąg Musi mieć wartość "duration".
delaySeconds uint64 Kilka sekund do zamknięcia okna. Wartość musi być większa niż 0.

Okno zliczające

Użyj count tej konfiguracji, aby zamknąć okno po określonej liczbie wiadomości.

W konfiguracji transformacji okna ustaw Liczba wiadomości na 5 i ustaw zachowanie komunikatu granicznego na messageInCurrent.

Majątek Typ Opis
type ciąg Musi mieć wartość "messageCount".
maxMessageCount uint64 Liczba wiadomości do buforowania przed zamknięciem okna. Wartość musi być większa niż 0.
boundaryMessage ciąg Czy komunikat zamykający okno pozostaje w bieżącym oknie (messageInCurrent) czy rozpoczyna następne okno (messageInNext).

Okno oparte na pamięci

Użyj memory tej konfiguracji, aby zamknąć okno, gdy rozmiar ładunku buforowanego osiągnie limit.

W konfiguracji przekształcenia okna ustaw Rozmiar bufora na 1048576 bajtów i ustaw działanie komunikatu granicznego na messageInNext.

Majątek Typ Opis
type ciąg Musi mieć wartość "bufferSize".
maxBufferBytes uint64 Maksymalna skumulowana liczba bajtów ładunku przed zamknięciem okna. Wartość musi być większa niż 0.
boundaryMessage ciąg Czy komunikat zamykający okno pozostaje w bieżącym oknie (messageInCurrent) czy rozpoczyna następne okno (messageInNext).

Okno wyzwalane zdarzeniem

Użyj konfiguracji triggers, gdy okno powinno zostać zamknięte na podstawie treści wiadomości lub stanu uruchomienia w bieżącym oknie.

W konfiguracji transformacji okna dodaj regułę wyzwalacza z polem wejściowym temperature, wyrażeniem running_sum($1) + $1 > 100 i zachowaniem komunikatu granicznego messageInCurrent.

Majątek Wymagane Opis
type Yes Musi mieć wartość "expression".
rules Yes Tablica reguł wyzwalania. Reguły są oceniane sekwencyjnie dla każdej wiadomości; Pierwsza zasada dopasowania zamyka okno.
datasets Nie. Opcjonalne zbiory danych do przechowywania stanów, które odwołują się do magazynu stanów.

Każda reguła wyzwalająca wspiera następujące pola:

Majątek Wymagane Opis
inputs Yes Tablica odniesień do pól wejściowych. Wyrażenie wiąże się z $1, $2, i tak dalej.
trigger Yes Wyrażenie logiczne, które zamyka okno, gdy jest obliczane jako true.
boundaryMessage Yes Czy komunikat zamykający okno pozostaje w bieżącym oknie (messageInCurrent) czy rozpoczyna następne okno (messageInNext).

Pole inputs obsługuje taką samą składnię danych wejściowych, jak ta stosowana w innych miejscach grafów przepływu danych, w tym zwykłe pola, wartości domyślne ??, ? $last, $context(key).field oraz $metadata.*. Więcej informacji na temat używania $context(key) można znaleźć w temacie Wzbogacanie przy użyciu danych zewnętrznych.

Wyrażenia wyzwalające mogą wykorzystywać funkcje wyrażeń grafu regularnego oraz następujące funkcje stanu działania, które resetują się po zamknięciu okna:

Function Opis
running_sum($1) Skumulowana suma $1 z poprzednich wiadomości w bieżącym oknie.
running_avg($1) Średnia skumulowana z $1 poprzednich wiadomości.
running_min($1) Minimalna wartość $1 widoczna w poprzednich wiadomościach. Zwraca $1 dla pierwszej wiadomości (minimum zbioru jednoelementowego jest tym elementem).
running_max($1) Maksymalna wartość $1 widziana w poprzednich wiadomościach. Zwraca $1 dla pierwszego komunikatu (maksimum zbioru jednoelementowego jest tym elementem).
running_count($1) Liczba wiadomości, w których element $1 był obecny.
running_count() Całkowita liczba wiadomości (bez filtra polowego).
first($1) Pierwsza niepusta wartość $1 w bieżącym oknie. Wraca $1 po pierwszej wiadomości.
changed($1) true , jeśli $1 różni się od wartości w poprzedniej wiadomości. false w przypadku pierwszej wiadomości w oknie (brak wcześniejszej wartości do porównania).
prev($1) Najnowsza niepusta wartość $1 z wcześniejszego komunikatu w bieżącym oknie. Wiadomości, gdzie $1 było puste, są pomijane (wartość zapisana nie jest nadpisana). Wraca $1 po pierwszej wiadomości w oknie.

Uwaga / Notatka

running_sum($1) a podobne funkcje zwracają wartości z wcześniej przetworzonych wiadomości. Dla bieżącej wiadomości użyj $1.

Przykłady reguł wyzwalania

Użyj tych przykładów, aby zobaczyć wspólne inputs i trigger wzorce w kompletnym obiekcie konfiguracji wyzwalacza:

  • Ten przykład pokazuje zwykłe wyrażenie wyzwalacza. Okno zamyka się, gdy prąd temperature przekracza 80.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature"],
      "trigger": "$1 > 80",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Ten przykład pokazuje wyrażenie wyzwalacza, które używa running_sum($1) + $1, aby połączyć wcześniejsze wiadomości w bieżącym oknie z bieżącą wiadomością, a następnie zamknąć okno, gdy próg zostanie przekroczony.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature"],
      "trigger": "running_sum($1) + $1 > 100",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Ten przykład pokazuje bezpieczną obsługę danych wejściowych z uwzględnieniem wartości null przy użyciu temperature ?? 0 oraz messageInNext, aby umieścić komunikat graniczny w następnym oknie.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature ?? 0"],
      "trigger": "running_avg($1) > 80",
      "boundaryMessage": "messageInNext"
    }
  ]
}
  • Ten przykład pokazuje wyzwalacz oparty na metadanych, w którym okno zamyka się dla określonej wartości tematu pochodzącej z $metadata.topic.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["$metadata.topic"],
      "trigger": "$1 == \"telemetry/high-priority\"",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Ten przykład pokazuje reguły wyzwalania z wykorzystaniem wzbogacania danych: dopasowuje komunikat factoryId do wiersza magazynu stanu, odczytuje shiftId z $context(factory).shiftId i zamyka okno, gdy ta wartość przesunięcia się zmienia (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"
    }
  ]
}

W tym przykładzie zbiór danych stanu reprezentowany przez factory oczekuje się, że zawiera pola takie jak factoryId i shiftId.

Zachowanie graniczne

Ustawienie boundaryMessage określa, co stanie się z wiadomością, która spowodowała zamknięcie okna opartego na liczbie, pamięci lub wyzwalaczu:

  • messageInCurrent: uwzględnij komunikat o granicy w oknie zamykającym.
  • messageInNext: najpierw zamknij bieżące okno, a następnie rozpocznij następne okno z komunikatem o granicy.

Jeśli messageInNext zostanie wywołane przy pierwszej wiadomości w nowym oknie, operacja zamknięcia jest pomijana, aby nie wygenerować pustego okna.

Uwaga / Notatka

Okno oparte na czasie trwania nie używa boundaryMessage. Granice czasu trwania są oparte na czasie, a nie na wiadomościach, więc nie ma wiadomości granicznej, którą można umieścić w obecnym lub następnym oknie.

Połącz warunki zamknięcia

Możesz połączyć czas trwania, liczbę danych, pamięć i warunki wyzwalające w tym samym wykresie.

  • Czas trwania jest sterowany czasem i oceniany przez timer.
  • Dla każdej przychodzącej wiadomości warunki sterowane wiadomością są oceniane w następującej kolejności: Memory > Count > Trigger.
  • W obrębie triggers.rules reguły są sprawdzane kolejno, a zastosowanie ma pierwsza pasująca reguła.

Kolejność oceny w warunkach wywołanych wiadomością ma znaczenie dla wyników kumulacji, gdy wiadomość spełnia wiele warunków jednocześnie. Na przykład, jeśli memory używa messageInCurrent i count używa messageInNext, komunikat spełniający oba warunki podąża za konfiguracją pamięci. Wiadomość pozostaje w obecnym oknie i ma wpływ na wynik akumulacji tego okna.

Definiowanie reguł akumulacyjnych

Każda reguła akumulowania określa, jak zmniejszyć okno komunikatów do pojedynczej wartości wyjściowej. Kluczem konfiguracji jest rules.

W konfiguracji transformacji okna dodaj regułę akumulacji z wejściem temperature, wyjściem avgTemperature i funkcją agregacji average($1).

Majątek Wymagane Opis
inputs Yes Lista ścieżek pól do odczytu z każdej przychodzącej wiadomości.
output Yes Ścieżka pola dla zagregowanego wyniku. Każda reguła musi mieć unikatowe dane wyjściowe.
expression Yes Formuła, która agreguje wartości wejściowe w oknie do pojedynczego skalara. Musi zawierać co najmniej jedną funkcję agregacji.
description Nie. Czytelny dla człowieka opis.

W przeciwieństwie do reguł mapy, expression jest wymagany dla każdej reguły akumulowania. Użycie samego $1 nie jest prawidłowe, ponieważ odwołuje się do kolekcji wartości, a nie pojedynczego skalara. Musisz opakowować go w funkcji agregacji, takiej jak average($1).

Funkcje agregacji

Function Zwroty Zachowanie pustego okna
average Średnia wartości liczbowych Błąd
sum Suma wartości liczbowych 0,0
min Minimalna wartość liczbowa Błąd
max Maksymalna wartość liczbowa Błąd
count Liczba komunikatów, w których istnieje pole 0
first Pierwsza wartość w oknie Błąd
last Ostatnia wartość w oknie Błąd

Każda funkcja przyjmuje pojedynczą zmienną pozycyjną jako argument ($1 dla pierwszego wejściowego, $2 drugiego itd.).

Wartości nieliczbowe: averagefunkcje , sum, mini max w trybie dyskretnym pomijają wartości nieliczbowe.

Funkcje oparte na obecności: count, firsti last działają na obecności pól niezależnie od typu wartości.

Łączenie agregacji

Połącz wiele funkcji agregacji w jednym wyrażeniu:

Dodaj regułę z danymi wejściowymi temperature i humidity, i wyrażeniem average($1) + max($2).

Aby przekonwertować zagregowaną wartość, zastosuj funkcję konwersji poza agregacją. Na przykład cToF(average($1)) konwertuje średnią temperaturę na Fahrenheita.

Każda funkcja agregacji musi bezpośrednio odwoływać się do pojedynczej zmiennej pozycyjnej. average($1) + max($2) jest prawidłowa, ale average($1 + $2) nie jest.

Różnice względem reguł mapy

Zdolność Mapuj reguły Reguły gromadzenia
Wymagane wyrażenie Nie. Yes
Dane wejściowe z symbolami wieloznacznymi Wsparte Niewspierane
$metadata Dostęp Wsparte Niewspierane
$context Wzbogacanie Wsparte Niewspierane
? $last dyrektywa Wsparte Niewspierane
Typ zawartości wyjściowej Dopasuj dane wejściowe Zawsze application/json

Pełny przykład konfiguracji

Ten przykład pokazuje pełną konfigurację okna, która zamyka okno po 30 sekundach, 5 wiadomości, 1 048 576 bajtów buforowanych lub gdy running_sum($1) + $1 > 100. Przykład ustawia boundaryMessage wartość dla messageInCurrent ostatnich trzech warunków, a okno oblicza statystyki temperatury po zamknięciu okna.

To, który warunek zamyka okno, zależy od czasu wysłania wiadomości, liczby sygnałów, wielkości ładunku i zawartości. Poniższe przykłady pokazują wynik dla każdego warunku zamknięcia.

Czas trwania się kończy

Jeśli żaden inny warunek nie zajdzie wcześniej, a okno osiągnie 30 sekund po otrzymaniu tych trzech komunikatów:

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

Komunikat wyjściowy to:

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

Licznik zamyka się

Jeśli okno otrzyma te pięć komunikatów, zanim zostanie spełniony jakikolwiek inny warunek:

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

Komunikat wyjściowy to:

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

Pamięć się zamyka

Jeśli rozmiar ładunku buforowanego osiągnie 1 048 576 bajtów przed wystąpieniem jakiegokolwiek innego warunku, na przykład po tych dwóch dużych komunikatach:

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

Komunikat wyjściowy to:

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

Wyzwalacz zamyka się

Jeśli wyrażenie running_sum($1) + $1 > 100 wyzwalające wywoła się przed jakimkolwiek innym warunkiem, na przykład po tych trzech komunikatach:

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

Komunikat wyjściowy to:

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

W zakresie operacji utwórz graf przepływu danych z transformacją okna:

  1. Dodaj źródło , które odczytuje z elementu telemetry/temperature.
  2. Dodaj przekształcenie window. Skonfiguruj 30-sekundowe okno czasowe, limit liczby wiadomości wynoszący 5, limit rozmiaru bufora wynoszący 1 048 576 bajtów oraz regułę wyzwalacza dla temperature z wyrażeniem running_sum($1) + $1 > 100. Dla warunków zliczania, pamięci i wyzwalania ustaw sposób działania komunikatu granicznego na messageInCurrent. Dodaj reguły agregacji dla średniej, minimum, maksimum, liczności i zakresu dla pola temperature.
  3. Dodaj docelowy punkt, który wysyła do telemetry/aggregated.

Następne kroki