Implementa aggregati definiti dall'utente JavaScript in Analisi di flusso di Azure

Analisi di flusso di Azure supporta aggregati definiti dall'utente (UDA) scritti in JavaScript, così da poter implementare una logica aziendale complessa e con stato. Con un UDA, hai il pieno controllo sulla struttura dati di stato, sull'accumulo di stati, sulla deaccumulazione di stati e sul calcolo aggregato dei risultati.

Usa un UDA JavaScript quando le funzioni aggregate integrate non soddisfano le tue esigenze e vuoi aggregare eventi a finestra con il tuo algoritmo.

Questo articolo ti mostra come creare un UDA e come chiamarlo con operazioni basate su finestre in una query Stream Analytics.

Prerequisiti

Prima di iniziare, assicurarsi di avere:

Scegli un tipo di aggregato definito dall'utente in JavaScript

Un aggregato definito dall'utente viene eseguito su una specifica di una finestra temporale per aggregare gli eventi all'interno di tale finestra e produrre un singolo valore di risultato. Stream Analytics supporta due tipi di interfacce UDA: AccumulateOnly e AccumulateDeaccumulate. Entrambi i tipi funzionano con tumbling, hopping, slipping e sessioni di finestre. Scegli il tipo in base all'algoritmo che usi.

Gli aggregati AccumulateDeaccumulate funzionano meglio degli aggregati AccumulateOnly quando li usi con finestre di salto e scorrimento, perché Stream Analytics può rimuovere eventi dallo stato invece di ricalcolarli.

Aggregazioni di tipo AccumulateOnly

Gli aggregati AccumulateOnly possono solo accumulare nuovi eventi nel loro stato. L'algoritmo non permette la deaccumulazione dei valori. Scegli questo tipo quando non puoi rimuovere le informazioni di un evento dal valore di stato. Il seguente codice è il template JavaScript per gli aggregati 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;
    }
}

Aggregazioni di tipo AccumulateDeaccumulate

Accumulare Gli aggregati Deaccumula deaccumulano un valore precedentemente accumulato dallo stato. Ad esempio, puoi rimuovere una coppia chiave-valore da una lista di valori di eventi o sottrarre un valore da un aggregato sommato. Il seguente codice è il template JavaScript per gli aggregati 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;
    }
}

Comprendere la dichiarazione di funzione JavaScript

Una dichiarazione di oggetto Funzione definisce ogni UDA JavaScript. Il seguente elenco descrive gli elementi principali in una definizione UDA.

Alias di funzione

L'alias funzionale è l'identificatore UDA. Quando chiami un UDA in una query di Stream Analytics, usa sempre l'alias insieme a un uda. prefisso.

Tipo di funzione

Per un UDA, imposta il tipo di funzione su JavaScript UDA.

Tipo di output

Imposta il tipo di output a un tipo specifico supportato dal lavoro Stream Analytics, oppure a Any se vuoi gestire il tipo nella tua query.

Nome della funzione

Il nome dell'oggetto Function. Il nome della funzione deve corrispondere all'alias UDA.

Metodo: init()

Il init() metodo inizializza lo stato dell'aggregato. Stream Analytics chiama questo metodo quando la finestra inizia.

Metodo: accumula()

Il accumulate() metodo calcola lo stato UDA basandosi sullo stato precedente e sui valori attuali degli eventi. Stream Analytics chiama questo metodo quando un evento entra in una finestra temporale (TumblingWindow, HoppingWindow, SlidingWindow, o SessionWindow).

Metodo: deaccumulate()

Il deaccumulate() metodo ricalcola lo stato basandosi sullo stato precedente e sui valori attuali dell'evento. Stream Analytics chiama questo metodo quando un evento lascia un SlidingWindow o SessionWindow.

Metodo: deaccumulateState()

Il deaccumulateState() metodo ricalcola lo stato basandosi sullo stato precedente e sullo stato di un salto. Stream Analytics chiama questo metodo quando un insieme di eventi lascia un HoppingWindow.

Metodo: computeResult()

Il computeResult() metodo restituisce il risultato aggregato basato sullo stato attuale. Stream Analytics chiama questo metodo alla fine di una finestra temporale (TumblingWindow, HoppingWindow, SlidingWindow, o SessionWindow).

Rivedi i tipi di dati di input e output supportati

Gli aggregati definiti dall'utente di JavaScript utilizzano le stesse conversioni di tipi di input e output delle funzioni definite dall'utente (UDF) di JavaScript. Per la mappatura completa tra i tipi di dati di Stream Analytics e quelli di JavaScript, consulta la sezione Conversione dei tipi tra Stream Analytics e JavaScript di Integrare le UDF JavaScript.

Aggiungi un UDA JavaScript nel portale Azure

In questa sezione, crei un UDA che calcola una media ponderata nel tempo. Per creare un UDA JavaScript in un lavoro Stream Analytics esistente, segui questi passaggi:

  1. Accedi al portale Azure e vai al tuo lavoro di Stream Analytics.

  2. Sotto Topologia del lavoro, seleziona Funzioni.

  3. Seleziona Aggiungi, poi seleziona JavaScript UDA.

  4. Nella pagina Nuova funzione , appare un template UDA predefinito nell'editor.

  5. Inserisci TWA come alias di funzione, e poi sostituisci l'implementazione della funzione con il seguente codice:

    // 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. Seleziona Salva. Il tuo UDA appare nella lista delle funzioni.

  7. Seleziona la nuova funzione TWA per rivedere la sua definizione.

Chiamare una UDA JavaScript in una query di Stream Analytics

Nel portale Azure, apri il tuo lavoro e modifica la query. Chiama la TWA() funzione con il prefisso obbligatorio uda. . Ad esempio:

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)

Testa la query con l'UDA

Crea un file JSON locale con il seguente contenuto, carica il file come input di esempio al tuo job di Stream Analytics e poi testa la query precedente:

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