Azure Stream Analytics'te JavaScript kullanıcı tanımlı toplamaları uygulama

Azure Stream Analytics, JavaScript'le yazılmış kullanıcı tanımlı agregasyonları (UDA) destekler, böylece karmaşık durumlu iş mantığı uygulayabilirsiniz. UDA ile durum veri yapısı, durum birikimi, durum dekomülasyonu ve toplam sonuç hesaplaması üzerinde tam kontrole sahip olursunuz.

Yerleşik toplu fonksiyonlar ihtiyaçlarınızı karşılamadığında ve kendi algoritmanızla pencereli olayları toplamak istiyorsanız JavaScript UDA kullanın.

Bu makale, bir UDA nasıl oluşturulacağını ve Stream Analytics sorgusunda pencere tabanlı işlemlerle nasıl çağrılacağını gösteriyor.

Prerequisites

Başlamadan önce şunları yaptığınızdan emin olun:

JavaScript kullanıcı tanımlı bir aggregate türü seçin

Kullanıcı tanımlı bir agregat, bir zaman penceresi spesifikasyonunun üzerinde çalışır ve o penceredeki olayları toplar ve tek bir sonuç değeri oluşturur. Stream Analytics, iki tür UDA arayüzü destekler: AccumulateOnly ve AccumulateDeaccumulate. Her iki tip de tumbling, hopping, kayma ve oturum pencereleriyle çalışır. Kullandığınız algoritmaya göre tipi seçin.

AccumulateDeaccumulate toplamaları, sıçramalı, kayan ve oturum pencereleriyle kullanıldığında, AccumulateOnly toplamalarına göre daha iyi performans gösterir; çünkü Stream Analytics, durumu yeniden hesaplamak yerine olayları durum bilgisinden kaldırabilir.

AccumulateOnly toplamları

AccumulateOnly agregaları durumlarına yalnızca yeni olaylar ekleyebilir. Algoritma değerlerin debirikimine izin vermiyor. Bir olayın bilgisini durum değerinden kaldıramadığınızda bu türü seçin. Aşağıdaki kod, AccumulateOnly toplamları için JavaScript şablonudur:

// 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 işlevi toplamları birleştirir ve ayırır

AccumulateDeaccumulate, daha önce biriktirilmiş bir değeri durumdan geri çıkarır. Örneğin, bir olay değerleri listesinden bir anahtar-değer çiftini çıkarabilir veya bir toplam agrematından bir değeri çıkarabilirsiniz. Aşağıdaki kod, AccumulateDeaccumulation agregaları için JavaScript şablonudur:

// 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 fonksiyon bildirgesini anlayın

Her JavaScript UDA'sını bir Fonksiyon nesne bildirmesi tanımlar. Aşağıdaki liste, UDA tanımındaki ana unsurları tanımlamaktadır.

Fonksiyon takma adı

Fonksiyon takma adı UDA tanımlayıcısıdır. Bir Stream Analytics sorgusunda UDA çağırdığınızda, her zaman takma adı uda. ön ekiyle birlikte kullanın.

İşlev türü

UDA için fonksiyon tipini JavaScript UDA olarak ayarlayın.

Çıkış türü

Çıkış tipini Stream Analytics işinin desteklediği belirli bir tipe veya sorgudaki tipi işlemek istiyorsanız Herhangi bir tipe ayarlayın.

İşlev adı

Function nesnesinin adı. Fonksiyon adı UDA takma adıyla eşleşmelidir.

Yöntem: init()

Yöntem, init() agreganın durumunu başlatır. Stream Analytics bu yöntemi pencere başladığında çağırır.

Yöntem: biriktir()

Yöntem, accumulate() UDA durumunu önceki durum ve mevcut olay değerleri üzerinden hesaplar. Stream Analytics, bir olay bir zaman penceresine (TumblingWindow, HoppingWindow, SlidingWindow veya SessionWindow) girdiğinde bu yöntemi çağırır.

Yöntem: deaccumulate()

Yöntem, deaccumulate() durumu önceki durum ve mevcut olay değerleri üzerinden yeniden hesaplar. Stream Analytics, bir olay bir SlidingWindow veya SessionWindow'dan ayrıldığında bu yöntemi çağırır.

Yöntem: deaccumulateState()

deaccumulateState() yöntemi, durumu önceki duruma ve bir atlamanın durumuna göre yeniden hesaplar. Stream Analytics, bir olay kümesi HoppingWindow öğesinden çıktığında bu yöntemi çağırır.

Yöntem: computeResult()

Yöntem, computeResult() mevcut duruma göre toplam sonucu döndürür. Stream Analytics bu yöntemi bir zaman penceresinin sonunda (TumblingWindow, HoppingWindow, SlidingWindow veya SessionWindow) çağırır.

Desteklenen giriş ve çıkış veri türlerini gözden geçirin

JavaScript kullanıcı tanımlı agregalar, JavaScript kullanıcı tanımlı fonksiyonlar (UDF) ile aynı giriş ve çıktı türü dönüşümlerini kullanır. Stream Analytics veri türleri ile JavaScript veri türleri arasındaki tam eşleme için, Stream Analytics ve JavaScript tür dönüştürme bölümüne bakın (JavaScript UDF'lerini tümleştirme).

Azure portalında bir JavaScript UDA ekleyin

Bu bölümde, zaman ağırlıklı bir ortalama hesaplayan bir UDA oluşturursunuz. Mevcut bir Stream Analytics işinde JavaScript UDA oluşturmak için şu adımları izleyin:

  1. Azure portalına giriş yapın ve Stream Analytics işinize gidin.

  2. İş topolojisi altından Fonksiyonlar'ı seçin.

  3. Ekle seçeneğini seçin, ardından JavaScript UDA'yı seçin.

  4. Yeni fonksiyon sayfasında, düzenleyicide varsayılan bir UDA şablonu görünür.

  5. Fonksiyon takma adını girin TWA ve ardından fonksiyon uygulamasını aşağıdaki kodla değiştirin:

    // 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. Kaydetseçeneğini seçin. UDA'nız fonksiyon listesinde görünüyor.

  7. Tanımını gözden geçirmek için yeni TWA fonksiyonunu seçin.

Stream Analytics sorgusunda JavaScript UDA'yı çağırın

Azure portalında işinizi açın ve sorguyu düzenleyin. TWA() zorunlu ön ekiyle uda. işlevini çağırın. Örneğin:

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)

Sorguyu UDA ile test edin

Aşağıdaki içerikle yerel bir JSON dosyası oluşturun, dosyayı örnek giriş olarak Stream Analytics işinize yükleyin ve ardından önceki sorguyu test edin:

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