Implementar agregados definidos por el usuario en JavaScript en Azure Stream Analytics

Azure Stream Analytics soporta agregados definidos por el usuario (UDA) escritos en JavaScript para que puedas implementar lógica de negocio con estado complejo. Con una UDA, tienes control total sobre la estructura de datos de estados, acumulación de estados, desacumulación de estados y cálculo agregado de resultados.

Usa un UDA en JavaScript cuando las funciones agregadas integradas no cumplan tus necesidades y quieras agregar eventos en ventana con tu propio algoritmo.

Este artículo te muestra cómo crear un UDA y cómo llamarlo mediante operaciones basadas en ventanas en una consulta de Stream Analytics.

Prerequisites

Antes de comenzar, asegúrese de que tiene:

Elige un tipo de agregado definido por el usuario en JavaScript

Un agregado definido por el usuario se ejecuta sobre una especificación de ventana temporal para agregar sobre los eventos de esa ventana y producir un único valor de resultado. Stream Analytics soporta dos tipos de interfaces UDA: AccumulateOnly y AccumulateDeaccumulate. Ambos tipos son compatibles con ventanas contiguas, ventanas con salto, ventanas deslizantes y ventanas de sesión. Elige el tipo según el algoritmo que utilices.

Los agregados AccumulateDeaccumulate rinden mejor que los agregados AccumulateOnly cuando los usas con ventanas de salto, deslizamiento y sesión, porque Stream Analytics puede eliminar eventos del estado en lugar de recalcularlos.

Agregados de AccumulateOnly

Los agregados AccumulateOnly solo pueden acumular nuevos eventos en su estado. El algoritmo no permite la desacumulación de valores. Elige este tipo cuando no puedas eliminar la información de un evento del valor del estado. El siguiente código es la plantilla de JavaScript para los agregados de 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;
    }
}

Agregados de AccumulateDeaccumulate

AcumularDesacumular agregados desacumulan un valor previamente acumulado del estado. Por ejemplo, puedes eliminar un par clave-valor de una lista de valores de eventos o restar un valor de un agregado sumado. El siguiente código es la plantilla JavaScript para los agregados 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;
    }
}

Entiende la declaración de función JavaScript

Una declaración de objeto Function define cada UDA de JavaScript. La siguiente lista describe los elementos principales en una definición de UDA.

Alias de función

El alias de la función es el identificador UDA. Cuando llames a una UDA en una consulta de Stream Analytics, usa siempre el alias precedido del prefijo uda..

Tipo de función

Para un UDA, se establece el tipo de función en JavaScript UDA.

Tipo de salida

Configura el tipo de salida a un tipo específico que soporte el trabajo de Stream Analytics, o a Cualquiera si quieres manejar el tipo de tu consulta.

Nombre de función

El nombre del objeto Función. El nombre de la función debe coincidir con el alias de UDA.

Método: init()

El init() método inicializa el estado del agregado. Stream Analytics llama a este método cuando comienza la ventana.

Método: acumular()

El accumulate() método calcula el estado UDA basándose en el estado anterior y los valores actuales del evento. Stream Analytics llama a este método cuando un evento entra en una ventana temporal (TumblingWindow, HoppingWindow, SlidingWindow, o SessionWindow).

Método: deaccumulate()

El deaccumulate() método recalcula el estado basándose en el estado anterior y en los valores actuales del evento. Stream Analytics llama a este método cuando un evento abandona un SlidingWindow o SessionWindow.

Método: deaccumulateState()

El deaccumulateState() método recalcula el estado en función del estado anterior y del estado de un salto. Stream Analytics llama a este método cuando un conjunto de eventos deja un HoppingWindow.

Método: computeResult()

El computeResult() método devuelve el resultado agregado basado en el estado actual. Stream Analytics llama a este método al final de una ventana temporal (TumblingWindow, HoppingWindow, SlidingWindow, o SessionWindow).

Revisar los tipos de datos de entrada y salida compatibles

Los agregados definidos por el usuario de JavaScript utilizan las mismas conversiones de tipos de entrada y salida que las funciones definidas por el usuario (UDF) de JavaScript. Para ver la correspondencia completa entre los tipos de datos de Stream Analytics y los tipos de datos de JavaScript, consulte la sección Stream Analytics y conversión de tipos de JavaScript de Integrar UDF de JavaScript.

Añadir un UDA en JavaScript en el portal de Azure

En esta sección, creas una UDA que calcula un promedio ponderado en el tiempo. Para crear un UDA en JavaScript en un trabajo existente de Análisis de Flujos, sigue estos pasos:

  1. Inicia sesión en el portal de Azure y accede a tu trabajo de Stream Analytics.

  2. En Topología de trabajos, selecciona Funciones.

  3. Selecciona Añadir y luego selecciona JavaScript UDA.

  4. En la página de Nueva función , aparece una plantilla UDA predeterminada en el editor.

  5. Introduce TWA como alias de función y luego reemplaza la implementación de la función por el siguiente código:

    // 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. Haga clic en Guardar. Tu UDA aparece en la lista de funciones.

  7. Selecciona la nueva función TWA para revisar su definición.

Invocar una UDA de JavaScript desde una consulta de Stream Analytics

En el portal de Azure, abre tu trabajo y edita la consulta. Llame a la función TWA() con el prefijo uda. obligatorio. Por ejemplo:

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)

Prueba la consulta con la UDA

Crea un archivo JSON local con el siguiente contenido, sube el archivo como entrada de muestra a tu trabajo de Stream Analytics y luego prueba la consulta anterior:

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