在 Azure 流分析 中实现 JavaScript 用户定义聚合

Azure 流分析 支持用 JavaScript 编写的用户定义聚合(UDA),以便实现复杂的有状态业务逻辑。 使用UDA时,你可以完全控制状态数据结构、状态累积、状态反累积以及聚合结果计算。

当内置聚合函数不满足你的需求,想用自己的算法聚合窗口事件时,可以使用JavaScript UDA。

本文将向您展示如何创建UDA以及如何在Stream Analytics查询中通过基于窗口的操作调用UDA。

先决条件

在开始之前,请确保具备:

选择一个JavaScript用户自定义聚合类型

用户自定义的聚合运行在时间窗口规范之上,对该窗口内的事件进行聚合,生成单一结果值。 Stream Analytics 支持两种类型的 UDA 接口:AccumulateOnly 和 AccumulateDeaccumulate。 这两种类型都支持滚动窗口、跳跃窗口、滑动窗口和会话窗口。 根据你使用的算法选择类型。

将 AccumulateDeaccumulate 聚合与跳跃窗口、滑动窗口和会话窗口配合使用时,其性能优于 AccumulateOnly 聚合,因为 Stream Analytics 可以从状态中删除事件,而不必重新计算状态。

仅聚合

AccumulateOnly 聚合只能将新事件积累到其状态。 算法不允许数值的去累积。 当无法从状态值中移除事件信息时,选择此类型。 以下代码是 AccumulateOnly 聚合的 JavaScript 模板:

// 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 聚合的 JavaScript 模板:

// 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函数声明

每个 JavaScript UDA 都由一个 Function 对象声明来定义。 以下列表描述了UDA定义中的主要要素。

函数别名

函数别名是 UDA 标识符。 在Stream Analytics查询中调用UDA时,务必使用别名和前 uda. 缀。

函数类型

对于 UDA,将函数类型设置为 JavaScript UDA

输出类型

将输出类型设置为Stream Analytics工作支持的特定类型,或者如果你想在查询中处理该类型,则设置为 任意

函数名称

函数对象的名称。 函数名称必须与UDA别名相匹配。

方法:init()

init() 方法初始化聚合的状态。 Stream Analytics 在窗口开始时调用此方法。

方法:累积()

accumulate() 方法根据之前的状态和当前事件值计算UDA状态。 当事件进入时间窗口(TumblingWindow, , HoppingWindowSlidingWindow, 或 SessionWindow)时,Stream Analytics 调用该方法。

方法:去累积()

deaccumulate() 方法根据之前的状态和当前事件值重新计算状态。 当事件离开 SlidingWindowSessionWindow 时,流分析会调用此方法。

方法:deaccumulateState()

deaccumulateState() 方法根据先前的状态和一个 hop 的状态重新计算状态。 当一组事件离开 HoppingWindow 时,流分析会调用此方法。

方法:computeResult()

computeResult() 方法返回基于当前状态的汇总结果。 Stream Analytics 在时间窗口结束时调用该方法(TumblingWindowHoppingWindowSlidingWindow, , 或 SessionWindow)。

回顾支持的输入和输出数据类型

JavaScript 用户自定义聚合使用与 JavaScript 用户自定义函数(UDF)相同的输入和输出类型转换。 有关流分析数据类型与JavaScript数据类型的完整映射,请参见“集成JavaScript UDFs”中的流分析与JavaScript类型转换部分。

在 Azure 门户中添加 JavaScript UDA

在本节中,你创建一个计算时间加权平均值的UDA文件。 要在现有的Stream Analytics作业中创建JavaScript UDA,请遵循以下步骤:

  1. 登录Azure门户,进入你的Stream Analytics职位。

  2. 作业拓扑中,选择 函数

  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 函数以查看其定义。

在 Stream Analytics 查询中调用 JavaScript UDA

在 Azure 门户中,打开你的工作并编辑查询。 调用带有强制uda.前缀的TWA()函数。 例如:

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}
]