JavaScript felhasználó által definiált aggregátumok implementálása az Azure Stream Analyticsben

Az Azure Stream Analytics támogatja a JavaScript-ben írt felhasználói definiált aggregátumokat (UDA), így összetett állapotos üzleti logikát tudsz megvalósítani. UDA-val teljes irányításod van az állapot-adatszerkezet, az állapot felhalmozódása, az állapot derakkulációja és az aggregált eredményszámítás felett.

Használj JavaScript UDA-t, ha a beépített aggregált függvények nem felelnek meg az igényeidnek, és saját algoritmusoddal szeretnél ablakos eseményeket aggregálni.

Ez a cikk bemutatja, hogyan lehet UDA-t létrehozni, és hogyan hívja ablakalapú műveleteket egy Stream Analytics lekérdezésben.

Prerequisites

Mielőtt hozzákezdene, győződjön meg arról, hogy rendelkezik az alábbiakval:

Válassz JavaScript felhasználó által definiált aggregált típust

Egy felhasználó által definiált aggregált rendszer egy időablak specifikáción fut, hogy az adott ablakban lévő eseményeket aggregálja, és egyetlen eredményértéket hozzon létre. A Stream Analytics kétféle UDA interfészt támogat: AccumulateOnly és AccumulateDeaccumulate. Mindkét típus működik dobó, ugró, csúszó és játékablakokkal. Válaszd ki a típust az általad használt algoritmus alapján.

Az AccumulateDeaccumulation aggregátumok jobban teljesítenek, mint az AccumulateOnly aggregációk, ha ugrás, csúszás és munkamenet ablakokkal használod őket, mert a Stream Analytics képes eltávolítani az eseményeket az állapotból ahelyett, hogy újraszámításra helyezné.

Halmozott összesítések

AccumulateOnly aggregátumok csak új eseményeket tudnak felhalmozni az állapotukba. Az algoritmus nem engedi az értékek dehalmozódását. Ezt a típust akkor válaszd, ha nem tudod eltávolítani egy esemény információját az állapotértékből. A következő kód a JavaScript sablonja az AccumulateOnly aggregációkhoz:

// 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;
    }
}

AccumulateDeaccumulate összesítések

Az AccumulateDeaccumulate eltávolít egy korábban felhalmozott értéket az állapotból. Például eltávolíthatsz egy kulcs-érték párt az eseményértékek listájáról, vagy levonhatsz egy értéket egy összegaggregátumból. Az alábbi kód a JavaScript sablonja az AccumulateDeaccumulation aggregátumokhoz:

// 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;
    }
}

Értsd meg a JavaScript függvény deklarációját

Egy függvényobjektum deklaráció határozza meg minden JavaScript UDA-t. Az alábbi lista az UDA definíciójának főbb elemeit írja le.

Függvény aliasa

A függvényalias az UDA azonosítója. Amikor UDA-t hívsz egy Stream Analytics lekérdezésben, mindig használd az aliast előtaggal uda. együtt.

Függvény típusa

UDA esetén állítsuk be a függvénytípust JavaScript UDA-nak.

Kimeneti típus

Állítsd be a kimeneti típust egy olyan típusra, amit a Stream Analytics munka támogat, vagy Any-ra , ha kezelni akarod a lekérdezésedben lévő típust.

Függvénynév

A függvényobjektum neve. A függvény nevének egyeznie kell az UDA álnévvel.

Módszer: init()

A init() módszer inicializálja az aggregált állapotot. A Stream Analytics ezt a módszert akkor hívja, amikor az ablak elkezdődik.

Módszer: accumulate()

A accumulate() módszer az UDA állapotot az előző állapot és a jelenlegi eseményértékek alapján számolja ki. A Stream Analytics ezt a módszert akkor nevezi el, amikor egy esemény belép egy időablakba (TumblingWindow, HoppingWindow, SlidingWindow, vagy SessionWindow).

Módszer: deaccumulate()

A deaccumulate() módszer az állapotot az előző állapot és a jelenlegi eseményértékek alapján újraszámolja. A Stream Analytics akkor hívja meg ezt a metódust, amikor egy esemény elhagy egy SlidingWindow vagy SessionWindow elemet.

Módszer: deaccumulateState()

A deaccumulateState() módszer az állapotot az előző állapot és egy hop állapota alapján újraszámítja. A Stream Analytics ezt a metódust akkor hívja meg, amikor események egy halmaza elhagy egy HoppingWindow elemet.

Metódus: computeResult()

A computeResult() módszer az aktuális állapot alapján adja vissza az összesített eredményt. A Stream Analytics ezt a módszert egy időablak végén nevezi (TumblingWindow, HoppingWindow, SlidingWindow, vagy SessionWindow).

Felülvizsgálja a támogatott bemeneti és kimeneti adattípusokat

A JavaScript felhasználódefiniált aggregátumok ugyanazokat a bemeneti és kimeneti típus-konverziókat használják, mint a JavaScript felhasználódefiniált függvények (UDF). A Stream Analytics adattípusok és JavaScript adattípusok közötti teljes leképezésért lásd az Integrál JavaScript UDF-ekStream Analytics és JavaScript típus-átalakítási szakaszát.

Hozzáadj egy JavaScript UDA-t az Azure portálba

Ebben a részben létrehozol egy UDA-t, amely idősúlyozott átlagot számol ki. JavaScript UDA létrehozásához egy meglévő Stream Analytics feladatban kövesse ezeket a lépéseket:

  1. Jelentkezz be az Azure portálra, és menj a Stream Analytics feladatodhoz.

  2. A Jobtopológia alatt válassza a Függvények lehetőséget.

  3. Válaszd az Add elemet, majd válaszd a JavaScript UDA elemet.

  4. Az Új funkció oldalon egy alapértelmezett UDA sablon jelenik meg a szerkesztőben.

  5. Adja meg a(z) TWA értéket a függvény aliasaként, majd cserélje le a függvény implementációját a következő kódra:

    // 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. Válassza az Mentésgombot. Az UDA-d megjelenik a funkciólistában.

  7. Válaszd ki az új TWA függvényt, hogy áttekintsd annak definícióját.

Hívjunk egy JavaScript UDA-t egy Stream Analytics lekérdezésben

Az Azure portálban nyisd meg a munkádat, és szerkessze a lekérdezést. Hívd meg a TWA() függvényt a kötelező uda. előtaggal. Példa:

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)

Teszteld a lekérdezést az UDA-val

Hozz létre egy helyi JSON fájlt a következő tartalommal, töltsd fel a fájlt mintabemenetként a Stream Analytics feladatodba, majd teszteld az előző lekérdezést:

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