Azure Stream Analytics에서 JavaScript 사용자 정의 집계를 구현하세요

Azure Stream Analytics는 JavaScript로 작성된 사용자 정의 집계(UDA)를 지원하여 복잡한 상태 기반 비즈니스 로직을 구현할 수 있습니다. UDA를 사용하면 상태 데이터 구조, 상태 축적, 상태 축적 해제, 그리고 집계 결과 계산을 완전히 제어할 수 있습니다.

내장 집계 함수가 필요를 충족하지 못하고 자체 알고리즘으로 창 이벤트를 집계하고 싶을 때는 JavaScript UDA를 사용하세요.

이 글에서는 UDA를 만드는 방법과 Stream Analytics 쿼리에서 윈도우 기반 연산으로 UDA를 호출하는 방법을 설명합니다.

Prerequisites

시작하기 전에 다음이 있는지 확인합니다.

JavaScript 사용자 정의 집합 타입을 선택하세요

사용자 정의 집계는 타임 윈도우 명세 위에 실행되어 해당 윈도우 내 이벤트들을 집계하여 단일 결과 값을 생성합니다. 스트림 애널리틱스는 두 가지 유형의 UDA 인터페이스를 지원합니다: AccumulateOnly와 AccumulateDeaccumulate. 두 유형 모두 텀블링, 홉핑, 슬라이딩, 세션 창을 지원합니다. 사용하는 알고리즘에 따라 유형을 선택하세요.

AccumulateDeaccumulate 집계 함수는 홉핑, 슬라이딩 및 세션 윈도우와 함께 사용할 때 AccumulateOnly 집계 함수보다 성능이 더 좋습니다. 이는 Stream Analytics가 상태를 다시 계산하는 대신 상태에서 이벤트를 제거할 수 있기 때문입니다.

AccumulateOnly 집계 기능

AccumulateOnly 애그리게이트는 새로운 이벤트만 자신의 상태에 누적할 수 있습니다. 알고리즘은 값의 축적 해제를 허용하지 않습니다. 이벤트 정보를 상태 값에서 제거할 수 없을 때 이 유형을 선택하세요. 다음 코드는 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는 상태에서 이전에 누적된 값의 누적을 해제합니다. 예를 들어, 이벤트 값 목록에서 키-값 쌍을 제거하거나 합산 집계에서 값을 뺄 수 있습니다. 다음 코드는 Accumulate Deaccumulate 집계를 위한 자바스크립트 템플릿입니다:

// 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 함수 선언을 이해하세요

함수 객체 선언은 각 자바스크립트 UDA를 정의합니다. 다음 목록은 UDA 정의의 주요 요소들을 설명합니다.

함수 별칭

함수 별칭은 UDA 식별자입니다. Stream Analytics 쿼리에서 UDA를 호출할 때는 항상 별칭과 uda. 접두사를 함께 사용하세요.

함수 유형

UDA의 경우 함수 타입을 JavaScript UDA로 설정하세요.

출력 형식

출력 타입을 Stream Analytics 직무가 지원하는 특정 타입으로 설정하거나, 쿼리에서 해당 타입을 처리하고 싶다면 Any 로 설정하세요.

함수 이름

함수 객체의 이름입니다. 함수 이름은 UDA 별칭과 일치해야 합니다.

메서드: init()

이 메서드는 init() 집합체의 상태를 초기화합니다. 스트림 애널리틱스는 창이 시작될 때 이 방식을 호출합니다.

방법: 누적()

accumulate() 방법은 이전 상태와 현재 이벤트 값을 기반으로 UDA 상태를 계산합니다. 스트림 애널리틱스는 이벤트가 시간 창TumblingWindow(, , HoppingWindowSlidingWindow, , 또는 SessionWindow)에 들어갈 때 이 방법을 호출합니다.

메서드: deaccumulate()

이 메서드는 deaccumulate() 이전 상태와 현재 이벤트 값을 기반으로 상태를 재계산합니다. Stream Analytics는 이벤트가 SlidingWindow 또는 SessionWindow를 벗어날 때 이 메서드를 호출합니다.

방법: deaccumulateState()

deaccumulateState() 방법은 이전 상태와 홉의 상태를 바탕으로 상태를 재계산합니다. 스트림 애널리틱스는 일련의 이벤트가 HoppingWindow를 벗어날 때 이 메서드를 호출합니다.

메서드: computeResult()

이 메서드는 computeResult() 현재 상태를 기반으로 한 집계 결과를 반환합니다. 스트림 애널리틱스는 시간 창이 끝날 때 이 방법을 호출합니다(TumblingWindow, , HoppingWindowSlidingWindow, , ).SessionWindow

지원되는 입력 및 출력 데이터 타입 검토

JavaScript 사용자 정의 집합체는 JavaScript 사용자 정의 함수(UDF)와 동일한 입력 및 출력 타입 변환을 사용합니다. 스트림 애널리틱스 데이터 타입과 자바스크립트 데이터 타입 간의 전체 매핑은 '자바스크립트 UDF를 통합하기'의 스트림 애널리틱스 및 자바스크립트 타입 변환 섹션을 참조하세요.

Azure Portal에서 JavaScript UDA 추가

이 섹션에서는 시간 가중 평균을 계산하는 UDA를 만듭니다. 기존 스트림 분석 작업에서 JavaScript UDA를 생성하려면 다음 단계를 따르세요:

  1. Azure 포털에 로그인한 후 Stream Analytics 직무로 이동하세요.

  2. 작업 토폴로지에서 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를 호출합니다

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