Implemente agregados definidos pelo usuário em JavaScript no Azure Stream Analytics

O Azure Stream Analytics suporta agregados definidos pelo usuário (UDA) escritos em JavaScript para que você possa implementar lógica de negócio complexa e com estado. Com um UDA, você tem controle total sobre a estrutura de dados de estado, acumulação de estados, desacumulação de estados e cálculo agregado de resultados.

Use um UDA em JavaScript quando as funções agregadas embutidas não atenderem às suas necessidades e você quiser agregar eventos em janelas com seu próprio algoritmo.

Este artigo mostra como criar um UDA e como chamá-lo com operações baseadas em janelas em uma consulta de Análise de Fluxos.

Pré-requisitos

Antes de começar, verifique se você tem:

Escolha um tipo de agregado definido pelo usuário em JavaScript

Um agregado definido pelo usuário opera com base em uma especificação de janela de tempo para agregar os eventos nessa janela e produzir um único valor resultante. A Análise de Fluxos suporta dois tipos de interfaces UDA: AccumulateOnly e AccumulateDeaccumulate. Ambos os tipos funcionam com tombadas, pulos, deslizamentos e janelas de sessão. Escolha o tipo com base no algoritmo que você usa.

Os agregados AccumulateDeaccumulate têm desempenho melhor do que os agregados AccumulateOnly quando usados com janelas deslizantes, de salto e de sessão, porque o Stream Analytics pode remover eventos do estado em vez de recalculá-lo.

Agregações AccumulateOnly

Os agregados AccumulateOnly só podem acumular novos eventos ao seu estado. O algoritmo não permite a desacumulação de valores. Escolha este tipo quando não for possível remover as informações de um evento do valor do estado. O código a seguir é o modelo JavaScript para agregados 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;
    }
}

Agregações AccumulateDeaccumulate

AccumulateDeaccumulate desacomula um valor previamente acumulado do estado. Por exemplo, você pode remover um par-chave-valor de uma lista de valores de eventos ou subtrair um valor de um agregado total. O código a seguir é o modelo JavaScript para os 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;
    }
}

Entenda a declaração de função JavaScript

Uma declaração de objeto Function define cada UDA em JavaScript. A lista a seguir descreve os principais elementos em uma definição da UDA.

Alias da função

O alias da função é o identificador UDA. Ao chamar uma UDA em uma consulta do Stream Analytics, sempre use o alias junto com o prefixo uda..

Tipo de função

Para um UDA, defina o tipo de função para JavaScript UDA.

Tipo de saída

Defina o tipo de saída para um tipo específico que o trabalho de Stream Analytics suporta, ou para Qualquer se quiser lidar com o tipo da sua consulta.

Nome da função

O nome do objeto Function. O nome da função deve corresponder ao alias UDA.

Método: init()

O init() método inicializa o estado do agregado. O Stream Analytics chama esse método quando a janela começa.

Método: accumulate()

O accumulate() método calcula o estado UDA com base no estado anterior e nos valores atuais do evento. A Stream Analytics chama esse método quando um evento entra em uma janela de tempo (TumblingWindow, HoppingWindow, SlidingWindow, ou SessionWindow).

Método: deaccumulate()

O deaccumulate() método recalcula o estado com base no estado anterior e nos valores atuais do evento. A Análise de Fluxo chama esse método quando um evento sai de um SlidingWindow ou SessionWindow.

Método: deaccumulateState()

O método deaccumulateState() recalcula o estado com base no estado anterior e no estado de um salto. A Stream Analytics chama esse método quando um conjunto de eventos deixa um HoppingWindow.

Método: computeResult()

O computeResult() método retorna o resultado agregado com base no estado atual. A Análise de Fluxos chama esse método ao final de uma janela de tempo (TumblingWindow, HoppingWindow, SlidingWindow, ou SessionWindow).

Revise os tipos de dados de entrada e saída suportados

Agregados definidos pelo usuário em JavaScript usam as mesmas conversões de tipos de entrada e saída que as funções definidas pelo usuário (UDF) do JavaScript. Para ver o mapeamento completo entre os tipos de dados do Stream Analytics e os tipos de dados do JavaScript, consulte a seção Conversão de tipos do Stream Analytics e do JavaScript de Integrar UDFs do JavaScript.

Adicionar um UDA em JavaScript no portal do Azure

Nesta seção, você cria um UDA que calcula uma média ponderada pelo tempo. Para criar um UDA em JavaScript em um trabalho de Stream Analytics existente, siga estes passos:

  1. Faça login no portal do Azure e acesse seu emprego de Análise de Fluxos.

  2. Em Topologia de funções, selecione Funções.

  3. Selecione Adicionar e depois selecione JavaScript UDA.

  4. Na página de Nova função , um modelo padrão de UDA aparece no editor.

  5. Insira TWA como o alias da função e então substitua a implementação da função pelo seguinte 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. Clique em Salvar. Seu UDA aparece na lista de funções.

  7. Selecione a nova função TWA para revisar sua definição.

Chamar uma UDA JavaScript em uma consulta do Stream Analytics

No portal do Azure, abra seu trabalho e edite a consulta. Chame a TWA() função com o prefixo obrigatório uda. . Por exemplo:

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 a consulta com a UDA

Crie um arquivo JSON local com o seguinte conteúdo, faça o upload do arquivo como entrada de exemplo para seu trabalho de Stream Analytics e então teste a 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}
]