Implementacja agregatów zdefiniowanych przez użytkownika JavaScript w Azure Stream Analytics

Usługa Azure Stream Analytics obsługuje agregacje definiowane przez użytkownika (UDA) napisane w języku JavaScript, dzięki czemu można implementować złożoną, stanową logikę biznesową. W UDA masz pełną kontrolę nad strukturą danych stanów, akumulacją stanów, dekumulacją stanów oraz obliczeniami wyników agregowanych.

Użyj JavaScript UDA, gdy wbudowane funkcje agregacji nie spełniają Twoich potrzeb i chcesz agregować zdarzenia okienkowe za pomocą własnego algorytmu.

Ten artykuł pokazuje, jak stworzyć UDA i jak wywołać ją za pomocą operacji okienkowych w zapytaniu Stream Analytics.

Wymagania wstępne

Przed rozpoczęciem upewnij się, że masz następujące elementy:

Wybierz typ agregatu zdefiniowany przez użytkownika w JavaScript

Agregat zdefiniowany przez użytkownika opiera się na specyfikacji okna czasowego, aby agregować zdarzenia w tym oknie i generować pojedynczą wartość wyniku. Stream Analytics obsługuje dwa typy interfejsów UDA: AccumulateOnly oraz AccumulateDeaccumulate. Oba typy działają z przewracaniem się, podskakiwaniem, ślizganiem i oknami sesji. Wybierz typ na podstawie używanego algorytmu.

Agregaty AccumulateDeaccumulate działają lepiej niż agregaty AccumulateOnly, gdy używasz ich z oknami hopping, sliding i session windows, ponieważ Stream Analytics może usuwać zdarzenia ze stanu zamiast je ponownie obliczać.

Agregacje wyłącznie kumulujące

Agregaty AccumulateOnly mogą jedynie dodawać nowe zdarzenia do swojego stanu. Algorytm nie pozwala na dekumulację wartości. Wybierz ten typ, gdy nie możesz usunąć informacji o zdarzeniu z wartości stanu. Poniższy kod to szablon JavaScript dla agregatów AccumulateOnly:

// Sample UDA which state can only be accumulated.
function main() {
    this.init = function () {
        this.state = 0;
    }

    this.accumulate = function (value, timestamp) {
        this.state += value;
    }

    this.computeResult = function () {
        return this.state;
    }
}

Akumulacja agregacji skumulowanych

AccumulateDeaccumulate usuwa ze stanu wcześniej zakumulowaną wartość. Na przykład możesz usunąć parę klucz-wartość z listy wartości zdarzeń lub odejmować wartość od sumy agregatu. Poniższy kod to szablon JavaScript dla agregatów AccumulateDeaccumulate:

// Sample UDA which state can be accumulated and deaccumulated.
function main() {
    this.init = function () {
        this.state = 0;
    }

    this.accumulate = function (value, timestamp) {
        this.state += value;
    }

    this.deaccumulate = function (value, timestamp) {
        this.state -= value;
    }

    this.deaccumulateState = function (otherState){
        this.state -= otherState.state;
    }

    this.computeResult = function () {
        return this.state;
    }
}

Zrozum deklarację funkcji JavaScript

Deklaracja obiektu Function definiuje każdą UDA w systemie JavaScript. Poniższa lista opisuje główne elementy definicji UDA.

Alias funkcji

Aliasem funkcji jest identyfikator UDA. Gdy wywołujesz UDA w zapytaniu Stream Analytics, zawsze używaj aliasu razem z prefiksem uda. .

Typ funkcji

Dla UDA ustaw typ funkcji na JavaScript UDA.

Typ danych wyjściowych

Ustaw typ wyjścia na konkretny typ, który obsługuje zadanie Stream Analytics, lub na Any, jeśli chcesz obsłużyć ten typ w swoim zapytaniu.

Nazwa funkcji

Nazwa obiektu Function. Nazwa funkcji musi odpowiadać aliasowi UDA.

Metoda: init()

Metoda init() inicjalizuje stan agregatu. Stream Analytics wywołuje tę metodę na początku okna.

Metoda: accumulate()

Metoda accumulate() oblicza stan UDA na podstawie poprzedniego stanu oraz wartości bieżących zdarzeń. Analiza strumieni używa tej metody, gdy zdarzenie wchodzi w okno czasowe (TumblingWindow, HoppingWindow, SlidingWindow, lub SessionWindow).

Metoda: deaccumulate()

Metoda deaccumulate() ta przelicza stan na podstawie poprzedniego stanu oraz wartości bieżących zdarzeń. Usługa Stream Analytics wywołuje tę metodę, gdy zdarzenie opuszcza element SlidingWindow lub SessionWindow.

Metoda: deaccumulateState()

Metoda deaccumulateState() ponownie oblicza stan na podstawie poprzedniego stanu oraz stanu przeskoku. Analiza strumieni wywołuje tę metodę, gdy zestaw zdarzeń opuszcza HoppingWindow.

Metoda: computeResult()

Metoda computeResult() zwraca łączny wynik na podstawie aktualnego stanu. Analiza strumieni wywołuje tę metodę na końcu okna czasowego (TumblingWindow, HoppingWindow, SlidingWindow, lub SessionWindow).

Przegląd obsługiwanych typów danych wejściowych i wyjściowych

Agregaty zdefiniowane przez użytkownika JavaScript wykorzystują te same konwersje typów wejściowych i wyjściowych co funkcje zdefiniowane przez użytkownika JavaScript (UDF). Aby uzyskać pełne mapowanie między typami danych Stream Analytics a typami danych JavaScript, zobacz sekcję Stream Analytics i konwersję typów JavaScript w Integrate JavaScript UDFS.

Dodaj obiekt UDA języka JavaScript w portalu Azure

W tej sekcji tworzysz UDA, która oblicza średnią ważoną czasowo. Aby utworzyć UDA w JavaScript w istniejącym stanowisku Stream Analytics, należy postępować zgodnie z następującymi krokami:

  1. Zaloguj się do portalu Azure i przejdź do swojej pracy w Stream Analytics.

  2. W sekcji Topologia zadań wybierz Funkcje.

  3. Wybierz Dodaj, a następnie wybierz JavaScript UDA.

  4. Na stronie Nowa funkcja w edytorze pojawia się domyślny szablon UDA.

  5. Wprowadź TWA jako alias funkcji, a następnie zastąp implementację funkcji następującym kodem:

    // Sample UDA which calculates the time-weighted average of incoming values.
    function main() {
        this.init = function () {
            this.totalValue = 0.0;
            this.totalWeight = 0.0;
        }
    
        this.accumulate = function (value, timestamp) {
            this.totalValue += value.level * value.weight;
            this.totalWeight += value.weight;
    
        }
    
        // Uncomment the following block for an AccumulateDeaccumulate implementation.
        /*
        this.deaccumulate = function (value, timestamp) {
            this.totalValue -= value.level * value.weight;
            this.totalWeight -= value.weight;
        }
    
        this.deaccumulateState = function (otherState){
            this.totalValue -= otherState.totalValue;
            this.totalWeight -= otherState.totalWeight;
        }
        */
    
        this.computeResult = function () {
            if(this.totalValue == 0) {
                result = 0;
            }
            else {
                result = this.totalValue/this.totalWeight;
            }
            return result;
        }
    }
    
  6. Wybierz Zapisz. Twoje UDA pojawia się na liście funkcji.

  7. Wybierz nową funkcję TWA, aby przejrzeć jej definicję.

Wywoływanie funkcji UDA języka JavaScript w zapytaniu usługi Stream Analytics

W portalu Azure otwórz swoje zadanie i edytuj zapytanie. Wywołaj funkcję TWA() z użyciem wymaganego prefiksu uda.. Na przykład:

WITH value AS
(
    SELECT
    NoiseLevelDB as level,
    DurationSecond as weight
FROM
    [YourInputAlias] TIMESTAMP BY EntryTime
)
SELECT
    System.Timestamp as ts,
    uda.TWA(value) as NoiseDoseTWA
FROM value
GROUP BY TumblingWindow(minute, 5)

Przetestuj zapytanie z UDA

Stwórz lokalny plik JSON z następującą zawartością, prześlij go jako przykładowe wejście do swojego zadania Stream Analytics, a następnie przetestuj poprzedzające zapytanie:

[
  {"EntryTime": "2017-06-10T05:01:00-07:00", "NoiseLevelDB": 80, "DurationSecond": 22.0},
  {"EntryTime": "2017-06-10T05:02:00-07:00", "NoiseLevelDB": 81, "DurationSecond": 37.8},
  {"EntryTime": "2017-06-10T05:02:00-07:00", "NoiseLevelDB": 85, "DurationSecond": 26.3},
  {"EntryTime": "2017-06-10T05:03:00-07:00", "NoiseLevelDB": 95, "DurationSecond": 13.7},
  {"EntryTime": "2017-06-10T05:03:00-07:00", "NoiseLevelDB": 88, "DurationSecond": 10.3},
  {"EntryTime": "2017-06-10T05:05:00-07:00", "NoiseLevelDB": 103, "DurationSecond": 5.5},
  {"EntryTime": "2017-06-10T05:06:00-07:00", "NoiseLevelDB": 99, "DurationSecond": 23.0},
  {"EntryTime": "2017-06-10T05:07:00-07:00", "NoiseLevelDB": 108, "DurationSecond": 1.76},
  {"EntryTime": "2017-06-10T05:07:00-07:00", "NoiseLevelDB": 79, "DurationSecond": 17.9},
  {"EntryTime": "2017-06-10T05:08:00-07:00", "NoiseLevelDB": 83, "DurationSecond": 27.1},
  {"EntryTime": "2017-06-10T05:09:00-07:00", "NoiseLevelDB": 91, "DurationSecond": 17.1},
  {"EntryTime": "2017-06-10T05:09:00-07:00", "NoiseLevelDB": 115, "DurationSecond": 7.9},
  {"EntryTime": "2017-06-10T05:09:00-07:00", "NoiseLevelDB": 80, "DurationSecond": 28.3},
  {"EntryTime": "2017-06-10T05:10:00-07:00", "NoiseLevelDB": 55, "DurationSecond": 18.2},
  {"EntryTime": "2017-06-10T05:10:00-07:00", "NoiseLevelDB": 93, "DurationSecond": 25.8},
  {"EntryTime": "2017-06-10T05:11:00-07:00", "NoiseLevelDB": 83, "DurationSecond": 11.4},
  {"EntryTime": "2017-06-10T05:12:00-07:00", "NoiseLevelDB": 89, "DurationSecond": 7.9},
  {"EntryTime": "2017-06-10T05:15:00-07:00", "NoiseLevelDB": 112, "DurationSecond": 3.7},
  {"EntryTime": "2017-06-10T05:15:00-07:00", "NoiseLevelDB": 93, "DurationSecond": 9.7},
  {"EntryTime": "2017-06-10T05:18:00-07:00", "NoiseLevelDB": 96, "DurationSecond": 3.7},
  {"EntryTime": "2017-06-10T05:20:00-07:00", "NoiseLevelDB": 108, "DurationSecond": 0.99},
  {"EntryTime": "2017-06-10T05:20:00-07:00", "NoiseLevelDB": 113, "DurationSecond": 25.1},
  {"EntryTime": "2017-06-10T05:22:00-07:00", "NoiseLevelDB": 110, "DurationSecond": 5.3}
]