Implementieren von benutzerdefinierten JavaScript-Aggregaten in Azure Stream Analytics

Azure Stream Analytics unterstützt benutzerdefinierte Aggregate (UDA), die in JavaScript geschrieben sind, sodass Sie komplexe stateful Business Logic implementieren können. Mit einer UDA haben Sie die volle Kontrolle über die Zustandsdatenstruktur, Zustandsakkumulation, Zustandsdeakkumulation und aggregierte Ergebnisberechnung.

Verwenden Sie ein JavaScript-UDA, wenn die integrierten Aggregatfunktionen Ihren Anforderungen nicht entsprechen und Sie fensterbasierte Ereignisse mit Ihrem eigenen Algorithmus aggregieren möchten.

Dieser Artikel zeigt Ihnen, wie Sie ein UDA erstellen und wie Sie es mit fensterbasierten Operationen in einer Stream Analytics-Abfrage aufrufen.

Voraussetzungen

Bevor Sie beginnen, stellen Sie sicher, dass Sie folgendes haben:

Wählen Sie einen JavaScript-benutzerdefinierten Aggregate-Typ

Ein benutzerdefiniertes Aggregat läuft auf einer Zeitfensterspezifikation, um die Ereignisse in diesem Fenster zu aggregieren und einen einzigen Ergebniswert zu erzeugen. Stream Analytics unterstützt zwei Arten von UDA-Schnittstellen: AccumulateOnly und AccumulateDeaccumulate. Beide Typen funktionieren mit Dreh-, Hüpf-, Gleit- und Sitzungsfenstern. Wählen Sie den Typ basierend auf dem verwendeten Algorithmus.

AccumulateDeaccumulate-Aggregate bieten eine bessere Leistung als AccumulateOnly-Aggregate, wenn sie mit Hopping-, Sliding- und Session-Fenstern verwendet werden, da Stream Analytics Ereignisse aus dem Zustand entfernen kann, anstatt den Zustand neu zu berechnen.

AccumulateOnly-Aggregate

AkkumulierenNur Aggregate können neue Ereignisse in ihren Zustand akkumulieren. Der Algorithmus erlaubt keine Deakkumulierung von Werten. Wählen Sie diesen Typ, wenn Sie die Informationen eines Ereignisses nicht aus dem Zustandswert entfernen können. Der folgende Code ist die JavaScript-Vorlage für AccumulateOnly-Aggregate:

// 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-Aggregate

AkkumulierenDeakkumulierte Aggregate deakkumulieren einen zuvor akkumulierten Wert aus dem Bundesstaat. Zum Beispiel kann man ein Schlüssel-Wert-Paar aus einer Liste von Ereigniswerten entfernen oder einen Wert von einem Summenaggregat abziehen. Der folgende Code ist die JavaScript-Vorlage für AccumulateDeaccumulate-Aggregate:

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

Verstehen Sie die JavaScript-Funktionsdeklaration

Eine Funktionsobjekt-Deklaration definiert jede JavaScript-UDA. Die folgende Liste beschreibt die Hauptelemente einer UDA-Definition.

Funktionsalias

Das Funktionsalias ist die UDA-Kennung. Wenn Sie eine UDA in einer Stream Analytics-Abfrage aufrufen, verwenden Sie immer das Alias zusammen mit einem uda. Präfix.

Funktionstyp

Für eine UDA setzen Sie den Funktionstyp auf JavaScript UDA.

Ausgabetyp

Setze den Ausgabe-Typ auf einen bestimmten Typ, den der Stream Analytics-Job unterstützt, oder auf Any, wenn du den Typ in deiner Abfrage bearbeiten möchtest.

Funktionsname

Der Name des Funktionsobjekts. Der Funktionsname muss mit dem UDA-Alias übereinstimmen.

Methode: init()

Die Methode init() initialisiert den Zustand des Aggregats. Stream Analytics ruft diese Methode auf, wenn das Fenster beginnt.

Methode: akkumulieren()

Die Methode accumulate() berechnet den UDA-Zustand basierend auf dem vorherigen Zustand und den aktuellen Ereigniswerten. Stream Analytics ruft diese Methode auf, wenn ein Ereignis in ein Zeitfenster eintritt (TumblingWindow, HoppingWindow, SlidingWindow, oder SessionWindow).

Methode: deaccumulate()

Die Methode deaccumulate() berechnet den Zustand basierend auf dem vorherigen Zustand und den aktuellen Ereigniswerten neu. Stream Analytics ruft diese Methode auf, wenn ein Ereignis ein SessionWindow oder SlidingWindow verlässt.

Methode: deaccumulateState()

Die Methode deaccumulateState() berechnet den Zustand basierend auf dem vorherigen Zustand und dem Zustand eines Hopps neu. Stream Analytics ruft diese Methode auf, wenn eine Reihe von Ereignissen ein HoppingWindow verlässt.

Methode: computeResult()

Die Methode computeResult() liefert das aggregierte Ergebnis basierend auf dem aktuellen Zustand. Stream Analytics ruft diese Methode am Ende eines Zeitfensters auf (TumblingWindow, HoppingWindow, SlidingWindow, oder SessionWindow).

Überprüfen Sie unterstützte Eingabe- und Ausgabedatentypen

JavaScript-benutzerdefinierte Aggregate verwenden dieselben Eingabe- und Ausgabetyp-Konvertierungen wie JavaScript-benutzerdefinierte Funktionen (UDF). Für die vollständige Zuordnung zwischen Stream Analytics Datentypen und JavaScript-Datentypen siehe den Abschnitt Stream Analytics und JavaScript-Typkonvertierung in Integrate JavaScript UDFs.

Hinzufügen einer JavaScript-UDA im Azure-Portal

In diesem Abschnitt erstellen Sie eine UDA, die einen zeitgewichteten Durchschnitt berechnet. Um ein JavaScript-UDA in einem bestehenden Stream Analytics-Job zu erstellen, folgen Sie diesen Schritten:

  1. Melden Sie sich im Azure-Portal an und gehen Sie zu Ihrem Stream Analytics-Job.

  2. Wählen Sie unter JobtopologieFunktionen aus.

  3. Wählen Sie Hinzufügen und dann JavaScript UDA.

  4. Auf der Seite Neue Funktionen erscheint im Editor eine Standard-UDA-Vorlage.

  5. Geben Sie als Funktionsalias ein TWA und ersetzen Sie dann die Funktionsimplementierung durch folgenden Code:

    // 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. Wählen Sie Speichern aus. Deine UDA erscheint in der Funktionsliste.

  7. Wählen Sie die neue TWA-Funktion aus, um deren Definition zu überprüfen.

Rufen Sie eine JavaScript-UDA in einer Stream Analytics-Abfrage auf

Im Azure-Portal öffnen Sie Ihren Job und bearbeiten Sie die Anfrage. Ruf die TWA() Funktion mit dem obligatorischen uda. Präfix auf. Beispiel:

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)

Teste die Abfrage mit der UDA

Erstellen Sie eine lokale JSON-Datei mit folgendem Inhalt, laden Sie die Datei als Beispieleingabe in Ihren Stream Analytics-Job hoch und testen Sie dann die vorherige Abfrage:

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