Azure Stream AnalyticsでJavaScriptのユーザー定義アグリゲートを実装してください

Azure Stream AnalyticsはJavaScriptで書かれたユーザー定義集約(UDA)をサポートしており、複雑なステートフルなビジネスロジックを実装できます。 UDAを使えば、状態データ構造、状態の蓄積、状態の非蓄積、集計結果の計算を完全に制御できます。

組み込みの集約関数がニーズに合わず、独自のアルゴリズムでウィンドウイベントを集約したい場合はJavaScript UDAを使いましょう。

この記事では、UDAの作成方法と、Stream Analyticsクエリでウィンドウベースの操作でUDAを呼び出す方法を紹介します。

前提条件

開始する前に、次の内容があることを確認します。

JavaScriptのユーザー定義アグリゲートタイプを選択する

ユーザー定義の集約は、タイムウィンドウ仕様の上に動作し、そのウィンドウ内のイベントを集約して単一の結果値を生成します。 Stream AnalyticsはCollectOnly とAccumulate Deaccumulateの2種類のUDAインターフェースをサポートしています。 どちらのタイプもタンブリング、ホッピング、スライド、セッションウィンドウに対応します。 使用するアルゴリズムに基づいてタイプを選びましょう。

AccumulateDeaccumulateの集約は、ホッピング、スライディング、セッションウィンドウで使用した場合、Stream Analyticsがイベントを再計算せずに状態から除去できるため、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()

accumulate()メソッドは、前の状態と現在のイベント値に基づいてUDAの状態を計算します。 Stream Analyticsは、イベントがタイムウィンドウ(TumblingWindowHoppingWindowSlidingWindow、または SessionWindow)に入るときにこの手法を呼びます。

メソッド: deaccumulate()

deaccumulate()メソッドは前の状態と現在のイベント値に基づいて状態を再計算します。 Stream Analyticsは、イベントが SlidingWindowSessionWindowを離れた場合にこのメソッドを呼び出します。

メソッド:deaccumulateState()

deaccumulateState()メソッドは、前の状態とホップの状態に基づいて状態を再計算します。 ストリームアナリティクスは、イベントの集合が一つの HoppingWindowを離れた場合にこのメソッドを呼びます。

メソッド:computeResult()

computeResult()メソッドは現在の状態に基づく集計結果を返します。 Stream Analyticsはこの手法を、時間の終わり(TumblingWindowHoppingWindowSlidingWindow、または SessionWindow)に呼びます。

サポートされている入力および出力データ型のレビュー

JavaScriptのユーザー定義集約は、JavaScriptのユーザー定義関数(UDF)と同じ入出力タイプの変換を使用します。 Stream Analyticsのデータ型とJavaScriptデータ型の完全なマッピングについては、「JavaScript UDFsの統合」のStream AnalyticsおよびJavaScript型変換セクションをご覧ください。

Azure portal で 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}
]