Veri akışı grafiklerinde pencere dönüşümleriyle veri topla

Bir pencere dönüşümü, gelen mesajları gruplar ve pencere kapandığında toplu değerlerle tek bir çıkış mesajı üretir. Her okumayı tek tek iletmek yerine, ortalamalar, minimumlar veya sayılar gibi istatistikleri hesaplayıp tek bir konsolide sonucu aşağıya gönderebilirsiniz.

Şu anda, bir pencere süre, sayı, bellek veya tetikleme koşullarına göre kapanabilir.

Uyarı

Süreye dayalı olmayan pencereleme, azureiotoperations/graph-dataflow-window:1.1.0 veya sonraki sürümlerini gerektirir.

Dönüşümler, değerleri, test koşullarını ve referans alanlarını hesaplamak için bir ifade dili kullanır. İfadeler, girdilere isimle değil, konuma göre atıfta bulunur: listedeki inputs ilk girdi $1, ikinci girdi $2, ve benzeri. cToF gibi yerleşik işlevler, bu değerleri dönüştürür ve işler.

Operatörler, fonksiyonlar, veri tipleri ve metaveri alanlarının tam listesi için İfadeler referansına bakınız.

Pencere dönüşümleri, yalnızca birikim kurallarında kullanılabilen average, min ve max gibi toplulaştırma işlevleri ekler. Tam liste için bkz. Toplama fonksiyonları.

Dönüşümler, değerleri, test koşullarını ve referans alanlarını hesaplamak için bir ifade dili kullanır. İfadeler, girdilere isimle değil, konuma göre atıfta bulunur: listedeki inputs ilk girdi $1, ikinci girdi $2, ve benzeri. cToF gibi yerleşik işlevler, bu değerleri dönüştürür ve işler.

Operatörler, fonksiyonlar, veri tipleri ve metaveri alanlarının tam listesi için İfadeler referansına bakınız.

Pencere dönüşümleri, yalnızca birikim kurallarında kullanılabilen average, min ve max gibi toplulaştırma işlevleri ekler. Tam liste için bkz. Toplama fonksiyonları.

Dönüşümler, değerleri, test koşullarını ve referans alanlarını hesaplamak için bir ifade dili kullanır. İfadeler, girdilere isimle değil, konuma göre atıfta bulunur: listedeki inputs ilk girdi $1, ikinci girdi $2, ve benzeri. cToF gibi yerleşik işlevler, bu değerleri dönüştürür ve işler.

Operatörler, fonksiyonlar, veri tipleri ve metaveri alanlarının tam listesi için İfadeler referansına bakınız.

Pencere dönüşümleri, yalnızca birikim kurallarında kullanılabilen average, min ve max gibi toplulaştırma işlevleri ekler. Tam liste için bkz. Toplama fonksiyonları.

Önkoşullar

  • Dağıtım sırasında default işaret eden mcr.microsoft.com adlı varsayılan kayıt defteri uç noktası otomatik olarak oluşturulur.

Bu makaledeki Azure CLI örnek, her değeri bir kez ayarlayabilmeniz ve ardından komutları kopyalayıp yapıştırmanız için ortam değişkenlerini kullanır as-is. Eğer hızlı başlangıçta Azure IoT İşlemleri Codespaces ortamını kullanıyorsanız, bu değişkenler zaten sizin için ayarlanmış ve bu adımı atlayabilirsiniz. Aksi takdirde, komutları çalıştırmadan önce shell'inizde aşağıdaki ortam değişkenlerini ayarlayın.

Aşağıdaki betikler en sık kullanılan ortam değişkenlerini ayarlar:

Ortam değişkeni Açıklama
SUBSCRIPTION_ID Azure IoT İşlemleri örneğini içeren aboneliğin kimliği.
RESOURCE_GROUP Azure IoT İşlemleri örneğini içeren kaynak grubunun adı.
AIO_INSTANCE_NAME Azure IoT İşlemleri örneğinizin adı. Örneklerinizi listelemek için çalıştırın az iot ops list -o table.
CLUSTER_NAME Azure Arc özellikli Kubernetes kümesinin adı, örneğini barındırıyor.
LOCATION Yeni kaynaklar için kullanılacak Azure bölgesi, örneğin eastus.
SUBSCRIPTION_ID=<subscription-id>
RESOURCE_GROUP=<resource-group-name>
AIO_INSTANCE_NAME=<instance-name>
CLUSTER_NAME=<cluster-name>
LOCATION=<region>

Sadece bu makalenin kullandığı değişkenleri ayarlamanız yeterlidir. Bu makale, seçtiğiniz kaynak adları için ek ortam değişkenleri kullanabilir. Makale, tanıtıldıkları yere nasıl yerleştirileceğini açıklıyor.

Durum bilgisi olan grafikler için ölçeklendirme sınırlaması

Önemli

Pencere dönüşümü içeren veri akışı grafikleri durum bilgisi içerir; her örnek iletileri bağımsız olarak biriktir. Veri akışı profilinde örnek sayısı birden fazla olduğunda, iletiler paylaşılan abonelikler aracılığıyla örnekler arasında dağıtılır. Her örnek kendi toplama durumunu koruduğundan ve örnekler birbirleriyle durum paylaşmadığından, her örnek iletilerin yalnızca bir alt kümesini görür. Bu, ortalamalar, toplamlar ve sayılar gibi toplama sonuçlarının kısmi bir veri kümesi üzerinden hesaplandığı ve yanlış olduğu anlamına gelir.

Doğru toplama sonuçlarından emin olmak için, pencere dönüşümü kullanan veri akışı grafı için veri akışı profili örnek sayısını1 olarak ayarlayın.

Pencere dönüşümü ne zaman kullanılır?

Yüksek frekanslı algılayıcı verileri aldığınızda ve aşağı akış göndermeden önce hacmi azaltmak istediğinizde pencere dönüşümünü kullanın. Yaygın senaryolar şunlardır:

  • İşlem ortalamaları: Sıcaklık algılayıcısı her saniye yayımlar, ancak bulut uygulamanızın yalnızca 30 saniyelik ortalamaya ihtiyacı vardır.
  • Aşırı değerleri izleme: Her bir dakikalık aralıkta en düşük ve en yüksek basınç okumalarını istiyorsunuz.
  • Olayları say: Son beş dakika içinde kaç tane kapı açık olay gerçekleştiğini bilmeniz gerekir.
  • Üretim partileri oluşturun: Her sabit boyutlu parti için istatistikleri hesaplamak istersiniz, örneğin bir dolum hattından çıkan her 100 paket gibi.
  • Durum değişikliklerine yanıt verin: Bir çalışma sinyalinin ne zaman değiştiğini bilmek isteyebilirsiniz; örneğin, bir mikser running durumundan draining durumuna geçtiğinde.

Pencere dönüştürme nasıl çalışır?

Pencere dönüşümünde sıralı olarak bağlı iki iç adım vardır:

  1. Pencere: Mesajları, yapılandırılmış kapanış koşullarından biri devreye girene kadar arabellekte tutar.
  2. Birikme: Pencere kapandığında toplama kurallarınızı uygular. Penceredeki tüm iletiler tek bir çıkış iletisine indirilir.

Uyarı

Bir pencere dönüşümü en az bir kapanış koşulu yapılandırmalıdır: delay, count, memory, veya triggers.

Pencere kapanma koşullarını yapılandırın

Sürümden 1.1.0itibaren pencere grafiği, mevcut delay anahtarın yanına üç yeni eş yapılandırma anahtarı ekler:

Yapılandırma anahtarı Pencere tipi Purpose
delay Süreye dayalı pencere Pencereyi belirli bir süre sonra kapatın.
count Sayıya dayalı pencere Sabit sayıda mesaj aldıktan sonra pencereyi kapatın.
memory Bellek tabanlı pencere Tamponlu yük boyutu bir sınıra ulaştığında pencereyi kapatın.
triggers Tetik tabanlı pencere Özel bir ifade true olarak değerlendirilirse pencereyi kapatın.

Süreye dayalı pencere

Konfigürasyon, delay her yuvarlanma penceresinin ne kadar süreceğini kontrol eder.

Uyarı

Gecikme adımı, ileti zaman damgalarını pencere sınırlarına hizalar. Bir mesaj 10 saniyelik bir pencerenin 7 saniyesinde gelirse, 10 saniyelik sınıra bağlıdır.

Uyarı

delay belirtmezseniz, güvenlik önlemi olarak 60 saniyelik varsayılan bir pencere zaman aşımı kullanılır.

Pencere dönüştürme yapılandırmasında Pencere süresini saniye cinsinden ayarlayın. Örneğin, 30 saniyelik düşen pencere için ayarını 30 olarak yapın.

Mülkiyet Türü Açıklama
type String olmalıdır "duration".
delaySeconds uint64 Pencerenin kapanmasından önceki saniye sayısı. 0'dan büyük olmalıdır.

Sayıya dayalı pencere

Sabit sayıda mesaj sonrası pencereyi kapatmak için count yapılandırmayı kullanın.

Pencere dönüşümü yapılandırmasında, Mesaj sayısını ayar 5 ve sınır mesaj davranışını messageInCurrent olarak ayarlayın.

Mülkiyet Türü Açıklama
type String olmalıdır "messageCount".
maxMessageCount uint64 Pencere kapanmadan önce tampon edilmesi gereken mesaj sayısı. 0'dan büyük olmalıdır.
boundaryMessage String Pencereyi kapatan mesajın mevcut pencerede mi kalacağı (messageInCurrent) yoksa sonraki pencereyi mi başlatacağı (messageInNext).

Bellek tabanlı pencere

Bufferli yük boyutu bir sınıra ulaştığında pencereyi kapatmak için yapılandırmayı memory kullanın.

Pencere dönüşüm yapılandırmasında, Buffer boyutunu bayt 1048576 olarak ayarlayın ve sınır mesaj davranışını messageInNext olarak ayarlayın.

Mülkiyet Türü Açıklama
type String olmalıdır "bufferSize".
maxBufferBytes uint64 Pencere kapanmadan önce maksimum birikimli yük baytları. 0'dan büyük olmalıdır.
boundaryMessage String Pencereyi kapatan mesajın mevcut pencerede mi kalacağı (messageInCurrent) yoksa sonraki pencereyi mi başlatacağı (messageInNext).

Tetik tabanlı pencere

Pencerenin kapanması gereken zaman, mesaj içeriğine veya mevcut penceredeki çalışma durumuna göre yapılandırmayı kullanın triggers .

Pencere dönüşümü yapılandırmasında, giriş alanı temperature, ifade running_sum($1) + $1 > 100ve sınır mesaj davranışı messageInCurrent ile bir tetikleme kuralı ekleyin.

Mülkiyet Zorunlu Açıklama
type Evet olmalıdır "expression".
rules Evet Tetikleyici kurallar dizisi. Kurallar mesaj başına sıralı olarak değerlendirilir; İlk eşleşme kuralı pencereyi kapatır.
datasets Hayır Durum deposuna referans veren isteğe bağlı durum-depo veri setleri.

Her tetikleme kuralı şu alanları destekler:

Mülkiyet Zorunlu Açıklama
inputs Evet Giriş alanı referanslarının dizisi. İfade $1, $2, ve benzeri şekilde bağlanır.
trigger Evet Değeri true olduğunda pencereyi kapatan Boole ifadesi.
boundaryMessage Evet Pencereyi kapatan mesajın mevcut pencerede mi kalacağı (messageInCurrent) yoksa sonraki pencereyi mi başlatacağı (messageInNext).

inputs alanı, veri akışı grafiklerinin başka yerlerinde kullanılan aynı giriş söz dizimini destekler; buna düz alanlar, ?? varsayılanları, ? $last, $context(key).field ve $metadata.* dahildir. Kullanımı $context(key)hakkında daha fazla bilgi için, Dış verilerle zenginleştirme bölümünü inceleyebilirsiniz.

Tetikleyici ifadeler, düzenli grafik ifade fonksiyonlarını ve pencere kapandığında sıfırlanan aşağıdaki çalışma durumu fonksiyonlarını kullanabilir:

Function Açıklama
running_sum($1) Mevcut penceredeki önceki mesajlardaki $1 değerlerinin kümülatif toplamı.
running_avg($1) Önceki mesajlardaki $1 değerinin kümülatif ortalaması.
running_min($1) Önceki mesajlarda görülen $1 için minimum değer. İlk mesajda geri $1 gelir (bir elemanın minimumu kendisidir).
running_max($1) Önceki mesajlarda görülen $1 maksimum değeri. İlk iletide $1 döndürür (tek bir öğenin maksimumu kendisidir).
running_count($1) Orada $1 bulunan mesaj sayısı.
running_count() Toplam mesaj sayısı (alan filtresi yok).
first($1) Geçerli penceredeki $1 öğesinin ilk boş olmayan değeri. İlk iletide $1 döndürür.
changed($1) true if $1 , önceki mesajdaki değerinden farklıdır. false pencerenin ilk mesajında (karşılaştırılacak önceki değer yok).
prev($1) Geçerli penceredeki önceki bir mesajdan alınan $1 değerinin boş olmayan en son değeri. $1 değerinin boş olduğu mesajlar atlanır (depolanan değerin üzerine yazılmaz). Bir pencerenin ilk iletisinde $1 döndürür.

Uyarı

running_sum($1) benzer fonksiyonlar daha önce işlenmiş mesajlardan değerler döndürür. Geçerli ileti için $1 kullanın.

Tetikleme kuralı örnekleri

Tam bir tetikleyici yapılandırma nesnesinde yaygın inputs ve trigger desenleri görmek için bu örnekleri kullanın:

  • Bu örnek, normal bir tetikleme ifadesini gösterir. Pencere, akım temperature 80'in üzerinde olduğunda kapanır.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature"],
      "trigger": "$1 > 80",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Bu örnek, mevcut penceredeki önceki mesajları mevcut mesajla birleştiren ve eşik aşıldığında pencereyi kapatan bir tetikleyici ifadesini running_sum($1) + $1 gösterir.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature"],
      "trigger": "running_sum($1) + $1 > 100",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Bu örnek, temperature ?? 0 ile null güvenli girdi işlemeyi ve sınır mesajını sonraki pencereye yerleştirmek için messageInNext kullanımını gösterir.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature ?? 0"],
      "trigger": "running_avg($1) > 80",
      "boundaryMessage": "messageInNext"
    }
  ]
}
  • Bu örnek, belirli bir konu değeri için pencerenin kapandığı metaveri tabanlı bir tetikleyiciyi göstermektedir $metadata.topic.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["$metadata.topic"],
      "trigger": "$1 == \"telemetry/high-priority\"",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Bu örnek, veri kümesi zenginleştirmesi kullanılarak tetikleme kurallarını gösterir: iletiyi factoryId durum deposundaki bir satırla eşleştirir, shiftId içinden $context(factory).shiftId okur ve bu shift değeri değiştiğinde pencereyi kapatır (changed($1)).
{
  "type": "expression",
  "datasets": [
    {
      "key": "factory",
      "inputs": ["$source.factoryId", "$context.factoryId"],
      "expression": "$1 == $2"
    }
  ],
  "rules": [
    {
      "inputs": ["$context(factory).shiftId"],
      "trigger": "changed($1)",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}

Bu örnekte, factory ile temsil edilen durum deposu veri kümesinin factoryId ve shiftId gibi alanlar içermesi beklenir.

Sınır davranışı

Ayar, boundaryMessagesayım, bellek veya tetik tabanlı pencerenin kapanmasına neden olan mesajın ne olduğunu kontrol eder:

  • messageInCurrent: sınır mesajını kapanma penceresine ekleyin.
  • messageInNext: önce mevcut pencereyi kapatın, sonra sınır mesajıyla bir sonraki pencereyi başlatın.

messageInNext, yeni bir penceredeki ilk mesajda tetiklenirse, boş bir pencere üretilmemesi için kapatma engellenir.

Uyarı

Süre tabanlı pencere boundaryMessage kullanmaz. Süre sınırları zamana bağlıdır, mesaj tabanlı değildir, bu yüzden mevcut veya sonraki pencereye koyulacak bir sınır mesajı yoktur.

Kapanış koşullarını birleştirin

Süre, sayı, bellek ve tetikleme koşullarını aynı grafikte birleştirebilirsiniz.

  • Süre zamana bağlıdır ve zamanlayıcı tarafından değerlendirilir.
  • Her gelen mesaj için, mesaj odaklı koşullar şu sırayla değerlendirilir: Memory > Count > Trigger.
  • triggers.rules içinde kurallar sırayla değerlendirilir ve eşleşen ilk kural geçerli olur.

Mesaj odaklı koşullarda değerlendirme sırası, bir mesajın aynı anda birden fazla koşulu tatmin ettiğinde birikim sonucu açısından önemlidir. Örneğin, memorymessageInCurrent kullanıyorsa ve countmessageInNext kullanıyorsa, her iki koşulu da karşılayan bir ileti bellek yapılandırmasının ardından gelir. Mesaj mevcut pencerede kalır ve o pencerenin birikim çıktısına katkıda bulunur.

Birikmesi kurallarını tanımlama

Her birikmesi kuralı, bir ileti penceresinin tek bir çıkış değerine nasıl küçültüldüğünü belirtir. Yapılandırma anahtarı şeklindedir rules.

Pencere dönüşümü yapılandırmasında, girdi temperature, çıktı avgTemperatureve toplama fonksiyonu average($1)ile birikim kuralı ekleyin.

Mülkiyet Zorunlu Açıklama
inputs Evet Gelen her iletiden okunacak olan alan yollarının listesi.
output Evet Toplu sonuç için alan yolu. Her kuralın benzersiz bir çıktısı olmalıdır.
expression Evet Pencere genelindeki giriş değerlerini tek bir skaler olarak azaltan formül. En az bir toplama işlevi içermelidir.
description Hayır İnsan tarafından okunabilen bir açıklama.

Harita kurallarının aksine, expression her birikmesi kuralı için gereklidir . $1'ün tek başına kullanımı geçerli değildir çünkü tek bir skaler yerine bir değer koleksiyonuna atıfta bulunur. Bunu average($1) gibi bir toplama işlevine sarmalamanız gerekir.

Toplama işlevleri

Function İadeler Boş pencere davranışı
average Sayısal değerlerin ortalaması Error
sum Sayısal değerlerin toplamı 0,0
min En düşük sayısal değer Error
max En büyük sayısal değer Error
count Alanın bulunduğu iletilerin sayısı 0
first Penceredeki ilk değer Error
last Penceredeki son değer Error

Her işlev bağımsız değişkeni olarak tek bir konumsal değişken alır ($1 ilk giriş için, $2 ikinci için vb.).

Sayısal olmayan değerler: average, sum, minve max işlevleri sayısal olmayan değerleri sessizce atlar.

Varlık tabanlı işlevler: count, first, ve last değer türünden bağımsız olarak alan varlığı üzerinde çalışır.

Toplamaları birleştirme

Birden çok toplama işlevini tek bir ifadede birleştirebilirsiniz:

temperature ve humidity girişlerini ve average($1) + max($2) ifadesini içeren bir kural ekleyin.

Toplanan bir değeri dönüştürmek için dönüştürme işlevini toplamanın dışına uygulayın. Örneğin, cToF(average($1)) ortalama sıcaklığı Fahrenheit'e dönüştürür.

Her toplama işlevinin doğrudan tek bir konumsal değişkene başvurması gerekir. average($1) + max($2) geçerli, ancak average($1 + $2) geçerli değil.

Harita kurallarından farklar

Kapasite Harita kuralları Birikim Kuralları
İfade gerekli Hayır Evet
Joker karakter girişleri Destekleniyor Desteklenmiyor
$metadata Erişim Destekleniyor Desteklenmiyor
$context Zenginleştirme Destekleniyor Desteklenmiyor
? $last Talimat Destekleniyor Desteklenmiyor
Çıkış içerik türü Girişle eşleşir Her zaman application/json

Tam yapılandırma örneği

30 saniye, 5 mesaj, 1048576 arabelleğe alınmış bayt sonrasında veya running_sum($1) + $1 > 100 olduğunda pencereyi kapatan, son üç koşul için boundaryMessage değeri messageInCurrent olarak ayarlanmış ve ardından sıcaklık istatistiklerini hesaplayan eksiksiz bir pencere yapılandırması.

Pencereyi hangi koşulla kapatacağı, mesaj zamanlamasına, sayına, yük boyutuna ve içeriğe bağlıdır. Aşağıdaki örnekler, her kapanış koşulu için çıkan çıktıyı göstermektedir.

Süre sona erer

Eğer başka bir koşul önce ateşlenmezse ve pencere bu üç mesajı aldıktan sonra 30 saniyeye ulaşırsa:

{ "temperature": 21.5 }
{ "temperature": 23.0 }
{ "temperature": 19.8 }

Çıkış iletisi:

{
  "avgTemperature": 21.433333333333334,
  "minTemperature": 19.8,
  "maxTemperature": 23.0,
  "readingCount": 3,
  "tempRange": 3.2
}

Sayım kapanıyor

Pencere, başka herhangi bir koşul tetiklenmeden önce bu beş mesajı alırsa:

{ "temperature": 20.0 }
{ "temperature": 22.0 }
{ "temperature": 21.0 }
{ "temperature": 24.0 }
{ "temperature": 23.0 }

Çıkış iletisi:

{
  "avgTemperature": 22.0,
  "minTemperature": 20.0,
  "maxTemperature": 24.0,
  "readingCount": 5,
  "tempRange": 4.0
}

Hafıza kapanıyor

Arabelleğe alınan yükün boyutu, başka herhangi bir koşul tetiklenmeden önce 1.048.576 bayta ulaşırsa, örneğin şu iki büyük iletiden sonra:

{ "temperature": 21.0, "payloadPad": "<large string>" }
{ "temperature": 22.5, "payloadPad": "<large string>" }

Çıkış iletisi:

{
  "avgTemperature": 21.75,
  "minTemperature": 21.0,
  "maxTemperature": 22.5,
  "readingCount": 2,
  "tempRange": 1.5
}

Tetik kapanıyor

Tetikleyici ifadesi running_sum($1) + $1 > 100 başka herhangi bir koşuldan önce ateşlenirse, örneğin bu üç mesajdan sonra:

{ "temperature": 40.0 }
{ "temperature": 35.0 }
{ "temperature": 30.0 }

Çıkış iletisi:

{
  "avgTemperature": 35.0,
  "minTemperature": 30.0,
  "maxTemperature": 40.0,
  "readingCount": 3,
  "tempRange": 10.0
}

İşlemler deneyiminde, pencere dönüşümüyle bir veri akışı grafiği oluşturun:

  1. kaynağından okuyan bir telemetry/temperature ekleyin.
  2. Pencere dönüşümü ekleyin. 30 saniyelik bir pencere, 5 mesaj sınırı, 1.048.576 baytlık arabellek boyutu sınırı ve temperature üzerinde running_sum($1) + $1 > 100 ifadesine sahip bir tetikleme kuralı yapılandırın. Sayım, bellek ve tetikleme koşulları için, sınır iletisi davranışını messageInCurrent olarak ayarlayın. Sahada temperature ortalama, minimum, max, sayı ve menzil için birikim kuralları ekleyin.
  3. Varış noktası ekleyin ki telemetry/aggregated adresine gönderebilsin.

Sonraki Adımlar