在 Azure 串流分析 中實作 JavaScript 使用者定義彙總

Azure 串流分析 支援以 JavaScript 撰寫的使用者定義聚合(UDA),以便實作複雜的有狀態商業邏輯。 使用 UDA 時,你可以完全控制狀態資料結構、狀態累積、狀態去累積以及結果彙總計算。

當內建的聚合函式不符合需求,想用自己的演算法聚合視窗事件時,可以使用 JavaScript UDA。

本文將教你如何建立 UDA,以及如何在 Stream Analytics 查詢中用視窗運算呼叫它。

先決條件

在開始之前,請確保你具備:

選擇一個 JavaScript 使用者定義的聚合型別

使用者自訂的聚合系統會建立在時間視窗規範之上,對該視窗內的事件進行聚合,產生單一結果值。 串流分析支援兩種 UDA 介面:AccumulateOnly 與 AccumulateDeaccumulate。 兩種類型皆適用於輪轉視窗、跳躍視窗、滑動視窗和工作階段視窗。 根據你使用的演算法選擇類型。

AccumulateDeaccumulate 彙總與跳躍視窗、滑動視窗和工作階段視窗搭配使用時,其效能優於 AccumulateOnly 彙總,因為串流分析可以從狀態中移除事件,而不需重新計算。

AccumulateOnly 彙總

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 彙總

AccumulateDeaccumulate 彙總會從狀態中取消累積先前已累積的值。 例如,你可以從事件值清單中移除鍵值對,或從總和彙總中減去一個值。 以下程式碼為 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 都由函式物件宣告定義。 以下列表描述了 UDA 定義中的主要元素。

函式別名

函式別名是 UDA 識別碼。 當你在 Stream Analytics 查詢中呼叫 UDA 時,務必使用別名和 uda. 前綴。

函式類型

對於 UDA,將函式類型設為 JavaScript UDA

輸出類型

輸出類型設為 Stream Analytics 工作支援的特定類型,或者如果你想處理查詢中的該類型,則設定 為任意

函式名稱

函式物件的名稱。 函式名稱必須與 UDA 別名相符。

方法:init()

init() 方法初始化集合的狀態。 Stream Analytics 在視窗開始時會呼叫此方法。

方法:累積()

accumulate() 方法根據先前狀態與當前事件值計算 UDA 狀態。 當事件進入時間窗(TumblingWindow、、HoppingWindowSlidingWindow、或SessionWindow)時,Stream Analytics 會呼叫此方法。

方法:deaccumulate()

deaccumulate() 方法會根據先前的狀態與當前事件值重新計算狀態。 當事件離開 SlidingWindowSessionWindow 時,串流分析會呼叫此方法。

方法:deaccumulateState()

deaccumulateState() 方法根據前一狀態及跳躍狀態重新計算狀態。 當一組事件離開 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 入口網站,打開你的工作並編輯查詢。 以必要的 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}
]