Реализовать пользовательские агрегатные функции на JavaScript в Azure Stream Analytics

Azure Stream Analytics поддерживает пользовательские агрегатные функции (UDA), написанные на JavaScript, что позволяет реализовать сложную бизнес-логику с отслеживанием состояния. С UDA вы полностью контролируете структуру данных состояний, накопление состояний, деаккуляцию состояния и вычисление агрегированных результатов.

Используйте JavaScript-UDA, если встроенные агрегатные функции не соответствуют вашим потребностям, и вы хотите агрегировать окнистые события с помощью собственного алгоритма.

В этой статье показано, как создать UDA и как вызвать её с помощью оконных операций в запросе Stream Analytics.

Prerequisites

Прежде чем начать, убедитесь, что у вас есть:

Выберите пользовательский агрегатный тип JavaScript

Пользовательский агрегат работает поверх спецификации временного окна, чтобы агрегировать события в этом окне и получить одно значение результата. Stream Analytics поддерживает два типа интерфейсов UDA: AccumulateOnly и AccumulateDeaccumulate. Оба типа работают с фиксированными, скачкообразными, скользящими и сеансовыми окнами. Выбирайте тип, исходя из используемого вами алгоритма.

Агрегатные функции AccumulateDeaccumulate работают лучше, чем агрегаты AccumulateOnly, при использовании со скачущими, скользящими и сеансовыми окнами, потому что Stream Analytics может удалять события из состояния вместо его повторного вычисления.

Статистические выражения AccumulateOnly

Агрегаты AccumulateOnly могут только накапливать новые события в своё состояние. Алгоритм не позволяет деаккумуляции значений. Выбирайте этот тип, если не можете удалить информацию о событии из значения состояния. Следующий код является JavaScript-шаблоном для агрегатов 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;
    }
}

Статистические выражения AccumulateDeaccumulate

AccumulateDeaccumulate деаккумулирует из состояния ранее накопленное значение. Например, вы можете удалить пару ключ-значение из списка значений событий или вычесть значение из суммарного агрегата. Следующий код является JavaScript-шаблоном для агрегатов 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;
    }
}

Понимание объявления функций JavaScript

Объявление объекта Function используется для определения каждого UDA JavaScript. Следующий список описывает основные элементы определения UDA.

Псевдоним функции

Псевдоним функции является идентификатором UDA. Когда вы вызываете UDA в запросе Stream Analytics, всегда используйте псевдоним вместе с uda. префиксом.

Тип функции *

Для UDA установите тип функции на JavaScript UDA.

Тип выходных данных

Установите тип вывода на конкретный тип, который поддерживает работа Stream Analytics, или на Any , если хотите обработать тип запроса.

Имя функции

Имя объекта Function. Название функции должно совпадать с псевдонимом UDA.

Метод: init()

Метод init() инициализирует состояние агрегата. Stream Analytics вызывает этот метод при начале окна.

Метод: накопление()

Метод вычисляет состояние accumulate() UDA на основе предыдущего состояния и текущих значений событий. Stream Analytics вызывает этот метод, когда событие входит в окно времени (TumblingWindow, HoppingWindow, SlidingWindow, или SessionWindow).

Метод: деаккуляция()

deaccumulate() Метод пересчитывает состояние на основе предыдущего состояния и текущих значений событий. Stream Analytics вызывает этот метод, когда событие покидает SlidingWindow или SessionWindow.

Метод: deaccumulateState()

Метод deaccumulateState() пересчитывает состояние на основе предыдущего состояния и состояния перехода. Stream Analytics вызывает этот метод, когда набор событий выходит из HoppingWindow.

Метод: computeResult()

Метод computeResult() возвращает совокупный результат на основе текущего состояния. Stream Analytics вызывает этот метод в конце временного окна (TumblingWindow, HoppingWindow, SlidingWindow, или SessionWindow).

Обзор поддерживаемых типов входных и выходных данных

Определяемые пользователем агрегатные функции JavaScript используют те же преобразования типов входных и выходных данных, что и определяемые пользователем функции JavaScript (UDF). Для полного сопоставления между типами данных Stream Analytics и типами данных JavaScript см. раздел Stream Analytics и конвертация типов JavaScript в разделе Интеграция JavaScript UDFs.

Добавить JavaScript UDA в портале Azure

В этом разделе вы создаёте UDA, который вычисляет взвешенное по времени среднее. Чтобы создать JavaScript-UDA в существующей работе Stream Analytics, следуйте следующим шагам:

  1. Войдите в портал Azure и перейдите на работу в Stream Analytics.

  2. В разделе Job topology выберите Functions.

  3. Выберите Добавить, а затем JavaScript UDA.

  4. На странице новой функции в редакторе отображается стандартный шаблон UDA.

  5. Введите TWA псевдоним функции, а затем замените реализацию функции следующим кодом:

    // 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. Нажмите Сохранить. Ваш UDA отображается в списке функций.

  7. Выберите новую функцию TWA , чтобы проверить её определение.

Вызовите JavaScript-UDA в запросе Stream Analytics

В портале Azure откройте задание и отредактируйте запрос. Вызовите TWA() функцию с обязательным uda. префиксом. Например:

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)

Проверьте запрос с помощью UDA

Создайте локальный JSON-файл с следующим содержанием, загрузите его в качестве примера входных данных в вашу задачу Stream Analytics и затем протестируйте предыдущий запрос:

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