Implementujte uživatelsky definované agregáty v JavaScriptu v Azure Stream Analytics

Azure Stream Analytics podporuje uživatelsky definované agregáty (UDA) napsané v JavaScriptu, takže můžete implementovat složitou stavovou obchodní logiku. S UDA máte plnou kontrolu nad datovou strukturou stavu, akumulací stavů, deakumulací stavů a výpočtem agregovaných výsledků.

Použijte JavaScript UDA, když vestavěné agregační funkce nesplňují vaše potřeby a chcete agregovat události v oknech vlastním algoritmem.

Tento článek vám ukáže, jak vytvořit UDA a jak ji volat pomocí operací založených na oknech v dotazu Stream Analytics.

Předpoklady

Než začnete, ujistěte se, že máte:

Vyberte uživatelsky definovaný typ agregátu v JavaScriptu

Uživatelem definovaný agregát běží nad specifikací časového okna, agreguje události v tomto okně a vytváří jednu výslednou hodnotu. Stream Analytics podporuje dva typy rozhraní UDA: AccumulateOnly a AccumulateDeaccumulate. Oba typy fungují s tumblingem, skákáním, klouzáním a s okny pro sezení. Vyberte typ podle algoritmu, který používáte.

Agregáty AccumulateDeaccumulate fungují lépe než agregáty AccumulateOnly, když je používáte s přeskakujícími, posuvnými a relačními okny, protože Stream Analytics může odstranit události ze stavu místo jeho přepočítávání.

Kumulované agregace

Agregáty AccumulateOnly mohou do svého stavu pouze akumulovat nové události. Algoritmus neumožňuje deakumulaci hodnot. Tento typ zvolte, pokud nemůžete odstranit informace o události z hodnoty stavu. Následující kód je JavaScriptová šablona pro agregáty 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;
    }
}

Kumulované agregace

Agregace AccumulateDeaccumulate dekumuluje dříve akumulovanou hodnotu ze stavu. Například můžete odstranit pár klíč-hodnota ze seznamu hodnot událostí nebo odečíst hodnotu z agregátu součtu. Následující kód je JavaScriptová šablona pro agregáty 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;
    }
}

Pochopte deklaraci JavaScriptové funkce

Deklarace objektu Function definuje každou JavaScriptovou UDA. Následující seznam popisuje hlavní prvky definice UDA.

Alias funkce

Alias funkce je identifikátor UDA. Když voláte UDA v dotazu Stream Analytics, vždy používejte alias spolu s prefixem uda. .

Typ funkce

Pro UDA nastavte typ funkce na JavaScript UDA.

Typ výstupu

Nastavte výstupní typ na konkrétní typ, který podporuje úloha Stream Analytics, nebo na Any, pokud chcete tento typ zvládnout ve svém dotazu.

Název funkce

Název objektu Funkce. Název funkce musí odpovídat aliasu UDA.

Metoda: init()

Metoda init() inicializuje stav agregátu. Stream Analytics tuto metodu volá při spuštění okna.

Metoda: accumulate()

Metoda accumulate() vypočítává stav UDA na základě předchozího stavu a aktuálních hodnot událostí. Stream Analytics tuto metodu nazývá, když událost vstoupí do časového okna (TumblingWindow, HoppingWindow, SlidingWindow, nebo SessionWindow).

Metoda: deaccumulate()

Metoda deaccumulate() přepočítává stav na základě předchozího stavu a aktuálních hodnot událostí. Stream Analytics volá tuto metodu, když událost opustí SlidingWindow nebo SessionWindow.

Metoda: deaccumulateState()

Metoda deaccumulateState() přepočítává stav na základě předchozího stavu a stavu skoku. Stream Analytics volá tuto metodu, když sada událostí opustí HoppingWindow.

Metoda: computeResult()

Metoda computeResult() vrací agregovaný výsledek na základě aktuálního stavu. Analýza proudu tuto metodu označuje na konci časového okna (TumblingWindow, HoppingWindow, SlidingWindow, nebo SessionWindow).

Zkontrolujte podporované typy vstupních a výstupních dat

JavaScriptové uživatelsky definované agregáty používají stejné převody vstupních a výstupních typů jako uživatelsky definované funkce (UDF) v JavaScriptu. Pro úplné mapování mezi datovými typy Stream Analytics a JavaScript typy viz sekce Stream Analytics a konverze JavaScript typů v Integrate JavaScript UDFS.

Přidejte JavaScript UDA v portálu Azure

V této části vytvoříte UDA, která počítá časově vážený průměr. Pro vytvoření JavaScript UDA v existující práci Stream Analytics postupujte podle těchto kroků:

  1. Přihlaste se do portálu Azure a přejděte do své práce v Stream Analytics.

  2. V sekci Topologie práce vyberte Funkce.

  3. Vyberte Přidat a potom vyberte JavaScript UDA.

  4. Na stránce Nová funkce se v editoru objeví výchozí šablona UDA.

  5. Zadejte alias funkce a poté nahraďte TWA implementaci funkce následujícím kódem:

    // 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. Vyberte Uložit. Vaše UDA se objevuje v seznamu funkcí.

  7. Vyberte novou funkci TWA pro zhodnocení její definice.

Zavolejte JavaScript UDA v dotazu Stream Analytics

V Azure portálu otevřete svou práci a upravte dotaz. Zavolejte funkci TWA() s povinnou předponou uda.. Příklad:

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)

Otestujte dotaz pomocí UDA

Vytvořte lokální JSON soubor s následujícím obsahem, nahrajte ho jako ukázkový vstup do vaší úlohy Stream Analytics a poté otestujte předchozí dotaz:

[
  {"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}
]