Użyj transformacji WebAssembly (WASM) w grafach przepływu danych

Operacje Azure IoT data flow graphs obejmują wbudowane przekształcenia dla typowych zadań przetwarzania, takich jak mapowanie, filtrowanie i agregacja. Jeśli potrzebujesz niestandardowej logiki poza wbudowanymi przekształceniami, możesz wdrożyć moduły WebAssembly (WASM) jako niestandardowe przekształcenia w potokach grafu przepływu danych.

Ważne

Obecnie interfejs webowy doświadczenia operacyjnego obsługuje jedynie tworzenie i przeglądanie artefaktów grafów przepływu danych pochodzących z Azure Container Registry (ACR), a dla wbudowanych transformacji – mcr.microsoft.com. Aby dowiedzieć się więcej, zobacz Internetowy interfejs użytkownika środowiska obsługi wyświetla tylko artefakty grafu przepływu danych pochodzące z usługi Azure Container Registry (ACR) i witryny mcr.microsoft.com.

Wymagania wstępne

  • Instancja usługi Operacje Azure IoT wdrożona w klastrze Kubernetes. Aby uzyskać więcej informacji, zobacz Deploy Operacje Azure IoT.

Przykłady interfejsu wiersza polecenia platformy Azure (Azure CLI) w tym artykule wykorzystują zmienne środowiskowe, dzięki czemu można ustawić każdą wartość tylko raz, a następnie kopiować i wklejać polecenia bez zmian. Jeśli korzystasz ze środowiska Operacje Azure IoT Codespaces z quickstartu, te zmienne są już ustawione i możesz pominąć ten krok. W przeciwnym razie ustaw następujące zmienne środowiskowe w swojej poskoczce przed uruchomieniem poleceń.

Poniższe skrypty określają najczęściej używane zmienne środowiskowe:

Zmienna środowiskowa Description
SUBSCRIPTION_ID ID subskrypcji zawierającej Twoją instancję Operacje Azure IoT.
RESOURCE_GROUP Nazwa grupy zasobów zawierającej Twoją instancję Operacje Azure IoT.
AIO_INSTANCE_NAME Nazwa Twojej instancji Operacje Azure IoT. Aby wyświetlić swoje instancje, uruchom az iot ops list -o table.
CLUSTER_NAME Nazwa klastra Kubernetes z włączonym Azure Arc, który hostuje Twoją instancję.
LOCATION Region platformy Azure, który ma być używany dla nowych zasobów, na przykład eastus.
SUBSCRIPTION_ID=<subscription-id>
RESOURCE_GROUP=<resource-group-name>
AIO_INSTANCE_NAME=<instance-name>
CLUSTER_NAME=<cluster-name>
LOCATION=<region>

Wystarczy ustawić zmienne, których używa ten artykuł. W tym artykule może być użytych dodatkowych zmiennych środowiskowych do wyboru nazw zasobów. Artykuł wyjaśnia, jak ustawić je tam, gdzie są wprowadzane.

W tym artykule używasz także następujących zmiennych środowiskowych dla wybranych przez Ciebie wartości: REGISTRY_ENDPOINT (nazwa endpointu rejestru), GRAPH_NAME (nazwa grafu przepływu danych), PROFILE (nazwa profilu przepływu danych) oraz REGISTRY_HOST (nazwa hosta rejestru kontenera). Ustaw każdą z nich przed uruchomieniem odpowiednich poleceń.

Przegląd

Korzystając z modułów WebAssembly (WASM) w wykresach przepływu danych Operacje Azure IoT, można przetwarzać dane na brzegu z wysoką wydajnością i zabezpieczeniami. WASM działa w środowisku piaskownicy i obsługuje języki Rust i Python.

Przepływ danych to potok, który przenosi i przekształca dane między punktami końcowymi za pomocą wbudowanych transformacji. Wykres przepływu danych rozszerza przepływy danych za pomocą kroków przetwarzania komponowalnego. Operacje Azure IoT udostępnia wykresy przepływu danych wbudowane dla typowych operacji, takich jak mapowanie, filtrowanie, rozgałęzianie i agregacja. W przypadku niestandardowej logiki przetwarzania zaimplementuj moduły WebAssembly, jak opisuje ten artykuł. Wykresy przepływu danych używają definicji grafu YAML, które określają sposób łączenia operatorów. Zasób grafu przepływu danych opakowuje tę definicję i mapuje swoje abstrakcyjne operacje źródła i ujścia do konkretnych punktów końcowych, takich jak tematy MQTT i tematy platformy Kafka.

Wskazówka

W przypadku większości scenariuszy przetwarzania danych zacznij od wbudowanych przekształceń. Używaj przekształceń WASM, gdy potrzebujesz niestandardowej logiki biznesowej, wyspecjalizowanych algorytmów lub przetwarzania, których wbudowane opcje nie obejmują.

Ważne

Wykresy przepływu danych obsługują obecnie tylko punkty końcowe MQTT, Kafka i OpenTelemetry. Inne typy punktów końcowych, takie jak Data Lake, Microsoft Fabric OneLake, Azure Data Explorer i Magazyn lokalny, nie są obsługiwane.

Jak działają wykresy przepływu danych WASM

Implementacja przepływu danych WASM jest zgodna z tym przepływem pracy:

  1. Opracowywanie modułów WASM: pisanie niestandardowej logiki przetwarzania w obsługiwanym języku i kompilowanie jej do formatu modelu składników zestawu WebAssembly. Aby dowiedzieć się więcej, zobacz Tworzenie modułów WASM dla przepływów danych.
  2. Tworzenie definicji grafu: definiowanie sposobu, w jaki dane przechodzą przez moduły przy użyciu plików konfiguracji YAML. Aby dowiedzieć się więcej, zobacz Konfigurowanie definicji grafu zestawu WebAssembly.
  3. Przechowywanie artefaktów w rejestrze: wypychanie skompilowanych modułów WASM i definicji grafów do rejestru kontenerów przy użyciu narzędzi zgodnych z protokołem OCI, takich jak ORAS. Aby dowiedzieć się więcej, zobacz Deploy WebAssembly (WASM) modules and graph definitions (Wdrażanie modułów WebAssembly (WASM) i definicji grafu.
  4. Konfiguruj punkty końcowe rejestru: Skonfiguruj szczegóły uwierzytelniania i połączenia, aby Operacje Azure IoT mógł uzyskać dostęp do rejestru kontenerów. Aby dowiedzieć się więcej, zobacz Konfigurowanie punktów końcowych rejestru.
  5. Utwórz graf przepływu danych: Użyj sieciowego interfejsu użytkownika do zarządzania operacjami lub plików Bicep, aby zdefiniować przepływ danych korzystający z definicji grafu.
  6. Deploy and execute: Operacje Azure IoT ściąga definicje grafu i moduły WASM z rejestru kontenerów i uruchamia je.

W poniższych przykładach pokazano, jak skonfigurować wykresy przepływu danych WASM dla typowych scenariuszy. W przykładach są używane wartości zakodowane na stałe i uproszczone konfiguracje, dzięki czemu można szybko rozpocząć pracę.

Przykład 1. Podstawowe wdrożenie z jednym modułem WASM

W tym przykładzie dane temperatury są konwertowane z fahrenheita na stopnie Celsjusza przy użyciu modułu WASM. Kod źródłowy modułu temperature jest dostępny w GitHub. Jeśli wykonano przykładowe kroki opisane w artykule Wdrażanie modułów WebAssembly (WASM) i definicji grafu, graph-simple:1.0.0 definicja grafu i wstępnie skompilowany temperature:1.0.0 moduł znajdują się już w rejestrze kontenerów. Przykładowa ścieżka artefaktu grafu to azure-samples/explore-iot-operations/graph-simple:1.0.0. Użyj tej ścieżki w przypadku publicznych przykładów GHCR lub podczas kopiowania przykładowych artefaktów do własnego rejestru z użyciem tej samej ścieżki repozytorium.

Jak działa prosty graf

Definicja grafu tworzy prosty, trzyetapowy potok.

  1. Źródło: Odbiera dane temperatury z MQTT
  2. Mapa: przetwarza dane za pomocą modułu Wasm do obsługi temperatury
  3. Sink: wysyła przekonwertowane dane z powrotem do MQTT

Aby uzyskać więcej informacji na temat sposobu działania prostej definicji grafu i jej struktury, zobacz Przykład 1: Prosta definicja grafu.

Format danych wejściowych:

{"temperature": {"value": 100.0, "unit": "F"}}

Format danych wyjściowych:

{"temperature": {"value": 37.8, "unit": "C"}}

Poniższa konfiguracja tworzy graf przepływu danych, który używa tego potoku konwersji temperatury. Wykres przepływu danych odwołuje się do graph-simple:1.0.0 definicji grafu YAML i ściąga moduł temperatury z rejestru kontenerów. Przykładowy graf platformy Azure znajduje się pod ścieżką repozytorium azure-samples/explore-iot-operations, dlatego uwzględnij tę ścieżkę w wartości artifact.

Konfigurowanie grafu przepływu danych

Ta konfiguracja definiuje trzy węzły, które implementują przepływ pracy konwersji temperatury: węzeł źródłowy, który subskrybuje przychodzące dane temperatury, węzeł przetwarzania grafu z uruchomionym modułem WASM i węzeł docelowy, który publikuje przekonwertowane wyniki.

Zasób grafu przepływu danych opakowuje artefakt definicji grafu i łączy jego abstrakcyjne operacje źródła i ujścia z konkretnymi punktami końcowymi:

  • Operacja definicji source grafu łączy się z węzłem źródłowym przepływu danych (temat MQTT)
  • Operacja definicji sink grafu łączy się z węzłem docelowym przepływu danych (temat MQTT)
  • Operacje przetwarzania definicji grafu są uruchamiane w węźle przetwarzania grafu

Ta separacja umożliwia wdrożenie tej samej definicji grafu z różnymi punktami końcowymi w różnych środowiskach przy zachowaniu logiki przetwarzania bez zmian.

  1. Aby utworzyć wykres przepływu danych w środowisku operacji, przejdź do karty Przepływ danych .

  2. Wybierz menu rozwijane obok + Create i wybierz Utwórz graf przepływu danych.

    Zrzut ekranu przedstawiający interfejs środowiska operacji przedstawiający sposób tworzenia grafu przepływu danych.

  3. Wybierz nazwę zastępczą new-data-flow , aby ustawić właściwości przepływu danych. Wprowadź nazwę grafu przepływu danych i wybierz profil przepływu danych do użycia.

  4. Na diagramie przepływu danych wybierz pozycję Źródło , aby skonfigurować węzeł źródłowy. W obszarze Szczegóły źródła wybierz pozycję Punktkońcowy zasobu lub przepływu danych.

    Zrzut ekranu przedstawiający interfejs środowiska operacji przedstawiający sposób wybierania źródła wykresu przepływu danych.

    1. W przypadku wybrania pozycji Zasób wybierz zasób do ściągnięcia danych, a następnie wybierz pozycję Zastosuj.

    2. Jeśli wybierzesz pozycję Punkt końcowy przepływu danych, wprowadź następujące szczegóły i wybierz pozycję Zastosuj.

      Setting Description
      Punkt końcowy przepływu danych Wybierz wartość domyślną , aby użyć domyślnego punktu końcowego brokera komunikatów MQTT.
      Temat Filtr tematu do subskrybowania wiadomości przychodzących. Użyj Temat(y)>Dodaj wiersz, aby dodać wiele tematów.
      Schemat komunikatu Schemat używany do deserializacji przychodzących komunikatów.
  5. Na diagramie przepływu danych wybierz pozycję Dodaj przekształcenie grafu (opcjonalnie), aby dodać węzeł przetwarzania grafu. W okienku Wybór grafu wybierz pozycję graph-simple:1 i wybierz pozycję Zastosuj.

    Zrzut ekranu przedstawiający interfejs środowiska operacji przedstawiający sposób tworzenia prostego grafu przepływu danych.

  6. Niektóre ustawienia operatora grafu można skonfigurować, wybierając węzeł grafu na diagramie. Możesz na przykład wybrać operator module-temperature/map i wprowadzić wartość key2example-value-2. Wybierz Zastosuj, aby zapisać zmiany.

    Zrzut ekranu przedstawiający interfejs środowiska operacji przedstawiający sposób konfigurowania prostego grafu przepływu danych.

  7. Na diagramie przepływu danych wybierz pozycję Miejsce docelowe , aby skonfigurować węzeł docelowy.

  8. Wybierz pozycję Zapisz pod nazwą grafu przepływu danych, aby zapisać wykres przepływu danych.

Uwaga / Notatka

Odwołanie do artefaktu jest powiązane z hostem punktu końcowego rejestru. W przypadku publicznych przykładowych obrazów GHCR hostem punktu końcowego rejestru jest ghcr.io, więc użyj azure-samples/explore-iot-operations/graph-simple:1.0.0. Użyj tej samej ścieżki artefaktu, jeśli skopiowano przykładowe artefakty do własnego rejestru w ramach tej samej ścieżki repozytorium. W przypadku własnego płaskiego układu rejestru prywatnego użyj płaskiego odwołania do artefaktu, takiego jak graph-simple:1.0.0.

Testowanie przepływu danych

Aby przetestować przepływ danych, wyślij komunikaty MQTT z klastra. Moduł temperatury oczekuje komunikatów w określonym formacie JSON z zagnieżdżonym obiektem temperature, który zawiera pola: value (liczbowe) i unit (ciąg). Na przykład: {"temperature":{"value":72,"unit":"F"}}.

Najpierw wdróż zasobnik klienta MQTT, postępując zgodnie z instrukcjami w temacie Testowanie łączności z brokerem MQTT przy użyciu klientów MQTT. Klient MQTT udostępnia tokeny uwierzytelniania i certyfikaty do nawiązywania połączenia z brokerem. Aby wdrożyć klienta MQTT, uruchom następujące polecenie:

kubectl apply -f https://raw.githubusercontent.com/Azure-Samples/explore-iot-operations/main/samples/quickstarts/mqtt-client.yaml

Wysyłanie komunikatów o temperaturze

W pierwszej sesji terminalu utwórz i uruchom skrypt w celu wysyłania danych temperatury w fahrenheit:

# Connect to the MQTT client pod
kubectl exec --stdin --tty mqtt-client -n azure-iot-operations -- sh -c '
# Create and run temperature.sh from within the MQTT client pod
while true; do
  # Generate a random temperature value between 0 and 6000 Fahrenheit
  random_value=$(shuf -i 0-6000 -n 1)
  payload="{\"temperature\":{\"value\":$random_value,\"unit\":\"F\"}}"

  echo "Publishing temperature: $payload"

  # Publish to the input topic
  mosquitto_pub -h aio-broker -p 18883 \
    -m "$payload" \
    -t "sensor/temperature/raw" \
    -d \
    --cafile /var/run/certs/ca.crt \
    -D PUBLISH user-property __ts $(date +%s)000:0:df \
    -D CONNECT authentication-method 'K8S-SAT' \
    -D CONNECT authentication-data $(cat /var/run/secrets/tokens/broker-sat)

  sleep 1
done'

Uwaga / Notatka

Własność __ts użytkownika MQTT dodaje znacznik czasu do wiadomości, aby zapewnić terminowe przetwarzanie, korzystając z Hybrid Logical Clock (HLC). Znacznik czasu pomaga przepływowi danych zdecydować, czy zaakceptować lub usunąć komunikat. Struktura właściwości to <timestamp>:<counter>:<nodeid>. Dzięki temu przetwarzanie przepływu danych jest dokładniejsze, ale nie jest obowiązkowe.

Skrypt publikuje przypadkowe dane o temperaturze na temat sensor/temperature/raw co sekundę. Powinien wyglądać następująco:

Publishing temperature: {"temperature":{"value":1234,"unit":"F"}}
Publishing temperature: {"temperature":{"value":5678,"unit":"F"}}

Pozostaw skrypt uruchomiony, aby kontynuować publikowanie danych dotyczących temperatury.

Subskrybowanie przetworzonych komunikatów

W drugiej sesji terminalu (również podłączonej do zasobnika klienta MQTT) zasubskrybuj temat wyjściowy, aby wyświetlić przekonwertowane wartości temperatury:

# Connect to the MQTT client pod
kubectl exec --stdin --tty mqtt-client -n azure-iot-operations -- sh -c '
mosquitto_sub -h aio-broker -p 18883 -t "sensor/temperature/processed" --cafile /var/run/certs/ca.crt \
-D CONNECT authentication-method "K8S-SAT" \
-D CONNECT authentication-data "$(cat /var/run/secrets/tokens/broker-sat)"'

Widzisz, jak moduł WASM konwertuje dane temperatury z Fahrenheita na Celsjusza.

{"temperature":{"value":1292.2222222222222,"count":0,"max":0.0,"min":0.0,"average":0.0,"last":0.0,"unit":"C","overtemp":false}}
{"temperature":{"value":203.33333333333334,"count":0,"max":0.0,"min":0.0,"average":0.0,"last":0.0,"unit":"C","overtemp":false}}

Przykład 2. Wdrażanie złożonego grafu

W tym przykładzie przedstawiono zaawansowany przepływ pracy przetwarzania danych, który obsługuje wiele typów danych, takich jak temperatura, wilgotność i dane obrazu. Definicja złożonego grafu organizuje wiele modułów WASM w celu przeprowadzania zaawansowanej analizy i wykrywania obiektów.

Jak działa złożony graf

Złożony graf przetwarza trzy strumienie danych i łączy je w wzbogaconą analizę czujników:

  • Przetwarzanie temperatury: konwertuje fahrenheita na stopnie Celsjusza, filtruje nieprawidłowe odczyty i oblicza statystyki
  • Przetwarzanie wilgotności: Gromadzi pomiary wilgotności w odstępach czasu
  • Przetwarzanie obrazów: wykonuje wykrywanie obiektów na migawkach aparatu i formatuje wyniki

Aby uzyskać więcej informacji na temat sposobu działania złożonej definicji grafu, jego struktury i przepływu danych przez wiele etapów przetwarzania, zobacz Przykład 2: złożona definicja grafu.

Wykres używa wyspecjalizowanych modułów z kolekcji operatorów Rust.

Konfigurowanie złożonego grafu przepływu danych

Ta konfiguracja implementuje przepływ pracy przetwarzania wielu czujników przy użyciu graph-complex:1.0.0 definicji grafu YAML. Przykładowa ścieżka artefaktu grafu to azure-samples/explore-iot-operations/graph-complex:1.0.0. Zwróć uwagę, że wdrożenie grafu przepływu danych jest podobne do przykładu Przykład 1 — oba używają tego samego wzorca z trzema węzłami (źródła, procesora grafu, miejsca docelowego), mimo że logika przetwarzania jest inna.

To podobieństwo występuje, ponieważ zasób grafu przepływu danych działa jako osadzone środowisko, które ładuje i wykonuje definicje grafów. Rzeczywista logika przetwarzania znajduje się w definicji grafu (graph-simple:1.0.0 lub graph-complex:1.0.0), która zawiera specyfikację YAML operacji i połączeń między modułami WASM. Zasób grafu przepływu danych udostępnia infrastrukturę środowiska uruchomieniowego do ściągania definicji grafu, tworzenia wystąpień modułów i kierowania danych za pośrednictwem zdefiniowanego przepływu pracy.

  1. Aby utworzyć wykres przepływu danych w środowisku operacji, przejdź do karty Przepływ danych .

  2. Wybierz menu rozwijane obok + Create i wybierz Utwórz wykres przepływu danych.

    Zrzut ekranu przedstawiający interfejs środowiska operacji przedstawiający sposób tworzenia złożonego grafu przepływu danych.

  3. Wybierz nazwę zastępczą new-data-flow , aby ustawić właściwości przepływu danych. Wprowadź nazwę grafu przepływu danych i wybierz profil przepływu danych do użycia.

  4. Na diagramie przepływu danych wybierz pozycję Źródło , aby skonfigurować węzeł źródłowy. W obszarze Szczegóły źródła wybierz pozycję Punktkońcowy zasobu lub przepływu danych.

    Zrzut ekranu przedstawiający interfejs środowiska operacji przedstawiający sposób wybierania źródła wykresu przepływu danych.

    1. W przypadku wybrania pozycji Zasób wybierz zasób do ściągnięcia danych, a następnie wybierz pozycję Zastosuj.

    2. Jeśli wybierzesz pozycję Punkt końcowy przepływu danych, wprowadź następujące szczegóły i wybierz pozycję Zastosuj.

      Setting Description
      Punkt końcowy przepływu danych Wybierz wartość domyślną , aby użyć domyślnego punktu końcowego brokera komunikatów MQTT.
      Temat Filtr tematu do subskrybowania wiadomości przychodzących. Użyj Temat(y)>Dodaj wiersz, aby dodać wiele tematów.
      Schemat komunikatu Schemat używany do deserializacji przychodzących komunikatów.
  5. Na diagramie przepływu danych wybierz pozycję Dodaj przekształcenie grafu (opcjonalnie), aby dodać węzeł przetwarzania grafu. W okienku Wybór grafu wybierz pozycję graf-złożony:1 i wybierz pozycję Zastosuj.

    Zrzut ekranu przedstawiający interfejs środowiska operacji przedstawiający sposób tworzenia złożonego grafu przepływu danych.

  6. Wybierz węzeł grafu na diagramie, aby skonfigurować ustawienia operatora grafu.

    Zrzut ekranu przedstawiający interfejs środowiska operacji przedstawiający sposób konfigurowania złożonego grafu przepływu danych.

    Operator Description
    migawka modułu/gałąź Konfiguruje moduł snapshot do wykrywania obiektów na obrazach. Klucz konfiguracji snapshot_topic można ustawić, aby określić temat wejściowy dla danych obrazu.
    temperatura modułu/mapa key2 Przekształca wartości temperatury w inną skalę.
  7. Wybierz Zastosuj, aby zapisać zmiany.

  8. Na diagramie przepływu danych wybierz pozycję Miejsce docelowe , aby skonfigurować węzeł docelowy.

  9. Wybierz pozycję Zapisz pod nazwą grafu przepływu danych, aby zapisać wykres przepływu danych.

Testowanie złożonego przepływu danych

Zanim zobaczysz jakiekolwiek dane wyjściowe, skonfiguruj dane źródłowe.

Przekazywanie plików obrazów RAW do zasobnika mqtt-client

Pliki obrazów są przeznaczone dla modułu snapshot do wykrywania obiektów na obrazach. Pliki są w folderze obrazów na GitHub.

Najpierw sklonuj repozytorium, aby uzyskać dostęp do plików obrazów:

git clone https://github.com/Azure-Samples/explore-iot-operations.git
cd explore-iot-operations

Aby przekazać pliki obrazów RAW z ./samples/wasm/images folderu do mqtt-client zasobnika, użyj następującego polecenia:

kubectl cp ./samples/wasm/images azure-iot-operations/mqtt-client:/tmp

Sprawdź, czy pliki zostały przesłane.

kubectl exec -it mqtt-client -n azure-iot-operations -- ls /tmp/images

Powinna zostać wyświetlona lista plików w folderze /tmp/images .

beaker.raw          laptop.raw          sunny2.raw
binoculars.raw      lawnmower.raw       sunny4.raw
broom.raw           milkcan.raw         thimble.raw
camera.raw          photocopier.raw     tripod.raw
computer_mouse.raw  radiator.raw        typewriter.raw
daisy3.raw          screwdriver.raw     vacuum_cleaner.raw
digital_clock.raw   sewing_machine.raw
hammer.raw          sliding_door.raw

Publikowanie symulowanych danych dotyczących temperatury i wilgotności oraz wysyłanie obrazów

Możesz połączyć polecenia do publikowania danych dotyczących temperatury i wilgotności oraz wysyłania obrazów do jednego skryptu. Użyj następującego polecenia:

# Connect to the MQTT client pod and run the script
kubectl exec --stdin --tty mqtt-client -n azure-iot-operations -- sh -c '
while true; do 
  # Generate a random temperature value between 0 and 6000
  temp_value=$(shuf -i 0-6000 -n 1)
  temp_payload="{\"temperature\":{\"value\":$temp_value,\"unit\":\"F\"}}"
  echo "Publishing temperature: $temp_payload"
  mosquitto_pub -h aio-broker -p 18883 \
    -m "$temp_payload" \
    -t "sensor/temperature/raw" \
    --cafile /var/run/certs/ca.crt \
    -D CONNECT authentication-method "K8S-SAT" \
    -D CONNECT authentication-data "$(cat /var/run/secrets/tokens/broker-sat)" \
    -D PUBLISH user-property __ts $(date +%s)000:0:df

  # Generate a random humidity value between 30 and 90
  humidity_value=$(shuf -i 30-90 -n 1)
  humidity_payload="{\"humidity\":{\"value\":$humidity_value}}"
  echo "Publishing humidity: $humidity_payload"
  mosquitto_pub -h aio-broker -p 18883 \
    -m "$humidity_payload" \
    -t "sensor/humidity/raw" \
    --cafile /var/run/certs/ca.crt \
    -D CONNECT authentication-method "K8S-SAT" \
    -D CONNECT authentication-data "$(cat /var/run/secrets/tokens/broker-sat)" \
    -D PUBLISH user-property __ts $(date +%s)000:0:df

  # Send an image every 2 seconds
  if [ $(( $(date +%s) % 2 )) -eq 0 ]; then
    file=$(ls /tmp/images/*.raw | shuf -n 1)
    echo "Sending file: $file"
    mosquitto_pub -h aio-broker -p 18883 \
      -f $file \
      -t "sensor/images/raw" \
      --cafile /var/run/certs/ca.crt \
      -D CONNECT authentication-method "K8S-SAT" \
      -D CONNECT authentication-data "$(cat /var/run/secrets/tokens/broker-sat)" \
      -D PUBLISH user-property __ts $(date +%s)000:0:df
  fi

  # Wait for 1 second before the next iteration
  sleep 1
done'

Sprawdzanie danych wyjściowych

W nowym terminalu zasubskrybuj temat wyjściowy:

kubectl exec --stdin --tty mqtt-client -n azure-iot-operations -- sh -c '
mosquitto_sub -h aio-broker -p 18883 -t "analytics/sensor/processed" --cafile /var/run/certs/ca.crt \
-D CONNECT authentication-method "K8S-SAT" \
-D CONNECT authentication-data "$(cat /var/run/secrets/tokens/broker-sat)"'

Dane wyjściowe wyglądają jak w poniższym przykładzie:

{"temperature":[{"count":9,"max":2984.4444444444443,"min":248.33333333333337,"average":1849.6296296296296,"last":2612.222222222222,"unit":"C","overtemp":true}],"humidity":[{"count":10,"max":76.0,"min":30.0,"average":49.7,"last":38.0}],"object":[{"result":"milk can; broom; screwdriver; binoculars, field glasses, opera glasses; toy terrier"}]}
{"temperature":[{"count":10,"max":2490.5555555555557,"min":430.55555555555554,"average":1442.6666666666667,"last":1270.5555555555557,"unit":"C","overtemp":true}],"humidity":[{"count":9,"max":87.0,"min":34.0,"average":57.666666666666664,"last":42.0}],"object":[{"result":"broom; Saint Bernard, St Bernard; radiator"}]}

W tym miejscu dane wyjściowe zawierają dane dotyczące temperatury i wilgotności, a także wykryte obiekty na obrazach.

Konfiguracja niestandardowych wykresów przepływu danych

Ta sekcja zawiera szczegółowe informacje na temat konfigurowania wykresów przepływu danych przy użyciu modułów WASM. Obejmuje wszystkie opcje konfiguracji, punkty końcowe przepływu danych i ustawienia zaawansowane.

Omówienie wykresu przepływu danych

Wykres przepływu danych definiuje sposób przepływu danych za pośrednictwem modułów WebAssembly na potrzeby przetwarzania. Każdy graf składa się z:

  • Tryb, który kontroluje, czy wykres jest włączony, czy wyłączony
  • Odwołanie do profila powiązanego z profilem przepływu danych, który definiuje skalowanie i ustawienia zasobów.
  • Trwałość dysku, która opcjonalnie włącza trwały magazyn dla stanu grafu
  • Węzły definiujące składniki źródłowe, obliczeniowe i docelowe
  • Połączenia węzłów określające sposób przepływu danych między węzłami

Konfiguracja trybu

Właściwość mode określa, czy wykres przepływu danych aktywnie przetwarza dane. Ustaw tryb na Enabled lub Disabled (bez uwzględniania wielkości liter). Po wyłączeniu wykres zatrzymuje przetwarzanie danych, ale zachowuje konfigurację.

Podczas tworzenia lub edytowania grafu przepływu danych w okienku Właściwości przepływu danych w obszarze Włącz przepływ danych zaznacz opcję Tak , aby ustawić tryb na Włączone. Jeśli zostawisz to niezaznaczone, tryb jest wyłączony.

Zrzut ekranu przedstawiający interfejs środowiska operacji przedstawiający sposób włączania lub wyłączania konfiguracji trybu.

Odniesienie do profilu

Dokumentacja profilu łączy wykres przepływu danych z profilem przepływu danych, który definiuje ustawienia skalowania, liczby wystąpień i limity zasobów. Jeśli nie określisz odwołania do profilu, musisz zamiast tego użyć odwołania właściciela platformy Kubernetes. Większość scenariuszy korzysta z domyślnego profilu, który oferuje Operacje Azure IoT.

Podczas tworzenia lub edytowania grafu przepływu danych w okienku Właściwości przepływu danych wybierz profil przepływu danych. Doświadczenie operacyjne automatycznie wybiera domyślny profil przepływu danych. Aby uzyskać więcej informacji na temat profilów przepływu danych, zobacz Konfigurowanie profilu przepływu danych.

Ważne

Profil przepływu danych możesz wybrać tylko wtedy, gdy tworzysz wykres przepływu danych. Nie możesz zmienić profilu przepływu danych po utworzeniu wykresu przepływu danych. Jeśli chcesz zmienić profil przepływu danych istniejącego grafu przepływu danych, usuń oryginalny wykres przepływu danych i utwórz nowy z nowym profilem przepływu danych.

Żądanie trwałości dysku

Żądanie persistencji dysku pomaga grafom przepływu danych utrzymywać stan przy restartach. Po włączeniu tej funkcji graf może przywrócić stan przetwarzania, jeśli podłączony broker się zrestartuje. Ta funkcja jest przydatna w scenariuszach przetwarzania stanowego, w których utrata danych pośrednich byłaby problematyczna. Po włączeniu trwałości dysku żądania broker utrwala dane MQTT, takie jak komunikaty w kolejce subskrybentów, na dysku. Takie podejście zapewnia, że źródło danych przepływu danych nie doświadczy utraty danych podczas awarii zasilania ani ponownego uruchomienia brokera. Broker utrzymuje optymalną wydajność, ponieważ konfigurujesz trwałość dla każdego przepływu danych, więc tylko przepływy danych wymagające trwałości korzystają z tej funkcji.

Wykres przepływu danych wykonuje to żądanie trwałości podczas subskrypcji, korzystając z właściwości użytkownika MQTTv5. Ta funkcja działa tylko wtedy, gdy:

  • Przepływ danych korzysta z brokera MQTT jako źródła (węzeł źródłowy z punktem końcowym MQTT)
  • Broker MQTT ma włączoną trwałość w trybie dynamicznym ustawionym na Enabled dla typu danych, takich jak kolejki subskrybentów.

W tej konfiguracji klienci MQTT, podobnie jak grafy przepływu danych, mogą żądać trwałości dysku dla swoich subskrypcji, korzystając z właściwości użytkownika MQTTv5. Aby uzyskać szczegółową konfigurację trwałości brokera MQTT, zobacz Konfigurowanie trwałości brokera MQTT.

Ustawienie akceptuje Enabled lub Disabled, z Disabled jako domyślnym.

Podczas tworzenia lub edytowania grafu przepływu danych w okienku Właściwości przepływu danych, w obszarze Żądaj trwałości danych zaznacz Tak w celu ustawienia trwałości dysku żądania jako Włączone. Jeśli pozostawisz to pole niezaznaczone, ustawienie to Wyłączone.

Reguły i limity nazewnictwa

Zasoby grafu przepływu danych i ich składniki mają ograniczenia nazewnictwa wymuszane w różnych warstwach:

Składnik Dozwolone znaki Długość Notatki
Nazwa zasobu grafu przepływu danych Małe litery alfanumeryczne i łączniki (a-z, 0-9, ). - Musi zaczynać się i kończyć znakiem alfanumerycznym. 3–63 znaki Wymuszane przez interfejs API usługi Azure Resource Manager.
Nazwa węzła Znaki alfanumeryczne, podkreślenia i łączniki (a-zA-Z0-9, _, -). Brak udokumentowanych limitów Musi być unikatowa w obrębie grafu.
Klucz konfiguracji Znaki alfanumeryczne, podkreślenia i łączniki (a-zA-Z0-9, _, -). Brak udokumentowanych limitów Pary kluczy i wartości przekazywane do modułów WASM.
Nazwa profilu przepływu danych Małe litery alfanumeryczne i łączniki. 3–39 znaków Limit 39 znaków wynika z ograniczeń nazwy podu w Kubernetes (limit 63 znaków minus prefiks aio-dataflow- i sufiks rewizji).
Odniesienie do schematu Musi być zgodny z formatem aio-sr://<namespace>/<name>:<version> lub aio-sr://<name>:<version>. N/A Używane w schematach połączeń węzłów.

Wykres przepływu danych wymusza również następujące reguły strukturalne:

  • Brak zduplikowanych nazw węzłów: każdy węzeł na grafie musi mieć unikatową nazwę.
  • Dozwolone typy połączeń: W grafie dozwolone są tylko następujące typy połączeń między węzłami: ze źródła do grafu, ze źródła do miejsca docelowego, z grafu do grafu oraz z grafu do miejsca docelowego.
  • Brak cykli: wykres nie może zawierać połączeń okrągłych, które tworzą nieskończone pętle przetwarzania.
  • Brak pętli własnych: węzeł nie może nawiązać połączenia z samym sobą.
  • Brak nakładania się tematów: jeśli źródło i miejsce docelowe używają tego samego punktu końcowego, tematy MQTT nie mogą się nakładać, co spowoduje utworzenie nieskończonej pętli komunikatów.

Konfiguracja węzła

Węzły to bloki konstrukcyjne grafu przepływu danych. Każdy węzeł ma unikatową nazwę w grafie i wykonuje określoną funkcję. Wykres zawiera trzy typy węzłów:

Węzły źródłowe

Węzły źródłowe definiują miejsce wprowadzania danych do grafu. Łączą się z punktami końcowymi przepływu danych, które odbierają dane z brokerów MQTT lub tematów platformy Kafka. Każdy węzeł źródłowy musi określać:

  • Odwołanie do punktu końcowego, które wskazuje na skonfigurowany punkt przepływu danych.
  • Źródła danych jako lista tematów MQTT lub tematów platformy Kafka do subskrypcji
  • Odwołanie do zasobu (opcjonalnie) powiązane z zasobem Rejestru Urządzeń Azure na potrzeby wnioskowania schematu.

Tablica źródeł danych obsługuje subskrypcję wielu tematów bez zmiany konfiguracji końcowych. Ta elastyczność pozwala ponownie wykorzystać endpointy w różnych przepływach danych.

Uwaga / Notatka

Obecnie wykresy przepływu danych obsługują jedynie punkty końcowe MQTT i Kafka jako źródła danych. Aby uzyskać więcej informacji, zobacz Konfigurowanie punktów końcowych przepływu danych.

Na diagramie przepływu danych wybierz pozycję Źródło , aby skonfigurować węzeł źródłowy. W obszarze Szczegóły źródła wybierz pozycję Punkt końcowy przepływu danych, a następnie użyj pola Temat(s), aby określić filtry tematu MQTT, aby subskrybować komunikaty przychodzące. Dodaj wiele tematów MQTT, wybierając pozycję Dodaj wiersz i wprowadzając nowy temat.

Węzły przetwarzania grafu

Węzły przetwarzania grafu zawierają moduły WebAssembly, które przekształcają dane. Te węzły ściągają artefakty WASM z rejestrów kontenerów i wykonują je przy użyciu określonych parametrów konfiguracji. Każdy węzeł grafu wymaga:

  • Odwołanie do punktu końcowego rejestru wskazującego na ten punkt, umożliwiającego pobieranie artefaktów.
  • Specyfikacja artefaktu definiująca nazwę i wersję modułu do ściągnięcia
  • Parametry konfiguracji jako pary klucz-wartość przekazywane do modułu WASM

Tablica konfiguracyjna umożliwia dostosowanie działania modułu bez przebudowy artefaktu WASM. Typowe opcje konfiguracji obejmują parametry przetwarzania, progi, ustawienia konwersji i flagi funkcji.

Na diagramie przepływu danych wybierz pozycję Dodaj przekształcenie grafu (opcjonalnie), aby dodać węzeł przetwarzania grafu. W okienku Wybór grafu wybierz żądany artefakt grafu, prosty lub złożony graf, a następnie wybierz pozycję Zastosuj. Niektóre ustawienia operatora grafu można skonfigurować, wybierając węzeł grafu na diagramie.

Pary klucz-wartość konfiguracji są przekazywane do modułu WASM w czasie wykonywania. Moduł może uzyskać dostęp do tych wartości, aby dostosować jego zachowanie. Stosując to podejście, możesz:

  • Wdróż ten sam moduł WASM z różnymi konfiguracjami.
  • Dostosuj parametry przetwarzania bez ponownego kompilowania modułów.
  • Włączanie lub wyłączanie funkcji na podstawie wymagań dotyczących wdrożenia.
  • Ustaw wartości specyficzne dla środowiska, takie jak progi lub punkty końcowe.

Ważne

Sprawdź dokumentację modułu WASM lub kod źródłowy, aby uzyskać wymagane parametry konfiguracji. Jeśli moduł oczekuje konkretnych parametrów (takich jak ograniczenia filtrów czy progi), a ty ich nie podasz, moduł może zawiść w czasie działania. Aby uzyskać szczegółowe informacje na temat definiowania parametrów w definicjach grafu, zobacz Parametry konfiguracji modułu.

Węzły docelowe

Węzły docelowe definiują, gdzie graf przesyła przetworzone dane. Łączą się z punktami końcowymi przepływu danych, które wysyłają dane do brokerów MQTT lub innych systemów. Każdy węzeł docelowy określa:

  • Odwołanie do punktu końcowego, które wskazuje na skonfigurowany punkt przepływu danych.
  • Miejsce docelowe danych jako konkretny temat, ścieżka lub lokalizacja dla danych wyjściowych.
  • Ustawienia schematu wyjściowego (opcjonalnie) definiujące format serializacji i walidację schematu.

Uwaga / Notatka

Obecnie wykresy przepływu danych obsługują jako miejsca docelowe tylko MQTT, Kafka i OpenTelemetry. Aby uzyskać więcej informacji, zobacz Konfigurowanie punktów końcowych przepływu danych.

  1. Na diagramie przepływu danych wybierz węzeł Docelowy .
  2. Wybierz punkt końcowy żądanego przepływu danych z listy rozwijanej Szczegóły punktu końcowego przepływu danych .
  3. Wybierz pozycję Kontynuuj , aby skonfigurować miejsce docelowe.
  4. Wprowadź wymagane ustawienia dla miejsca docelowego, w tym temat lub tabelę do wysłania danych. Portal automatycznie interpretuje pole miejsca docelowego danych na podstawie typu punktu końcowego. Jeśli na przykład punkt końcowy przepływu danych jest punktem końcowym MQTT, na stronie szczegółów docelowej zostanie wyświetlony monit o wprowadzenie tematu.

Połączenia węzłów

Połączenia węzłów definiują ścieżkę przepływu danych między węzłami. Każde połączenie określa węzeł źródłowy i węzeł docelowy, tworząc potok przetwarzania. Opcjonalnie możesz dołączyć schemat do połączenia. Moduł otrzymuje schemat przy inicjalizacji, co wspiera walidację schematu, jak w tym przykładzie.

Środowisko operacji automatycznie tworzy połączenia węzłów po wybraniu węzła przetwarzania grafu. Nie możesz modyfikować połączeń po utworzeniu wykresu.

Punkty końcowe przepływu danych

Wykresy przepływu danych łączą się z systemami zewnętrznymi za pośrednictwem punktów końcowych przepływu danych. Typ punktu końcowego określa, czy można go użyć jako źródła, miejsca docelowego, czy obu tych elementów.

Punkty końcowe MQTT

Punkty końcowe MQTT mogą służyć zarówno jako źródła, jak i miejsca docelowe. Łączą się z brokerami MQTT, w tym:

  • Operacje Azure IoT lokalny broker MQTT (wymagany w każdym przepływie danych)
  • Azure Event Grid MQTT
  • Niestandardowe brokery MQTT

Aby uzyskać szczegółowe informacje o konfiguracji, zobacz Konfigurowanie punktów końcowych przepływu danych MQTT.

Punkty końcowe platformy Kafka

Punkty końcowe platformy Kafka mogą służyć zarówno jako źródła, jak i miejsca docelowe. Łączą się one z systemami zgodnymi z platformą Kafka, w tym:

  • Azure Event Hubs (zgodne z platformą Kafka)
  • Klastry Apache Kafka
  • Confluent Cloud

Aby uzyskać szczegółowe informacje o konfiguracji, zobacz Konfigurowanie punktów końcowych przepływu danych platformy Azure Event Hubs i Kafka.

Punkty końcowe rejestru

Punkty końcowe rejestru zapewniają dostęp do rejestrów kontenerów na potrzeby ściągania modułów WASM i definicji grafu. Nie są one używane bezpośrednio w przepływie danych, ale węzły przetwarzania grafu odwołują się do nich.

Aby uzyskać szczegółowe informacje o konfiguracji, zobacz Konfigurowanie punktów końcowych rejestru.

Rozwiązywanie problemów z wykresami przepływu danych

Ta sekcja zawiera porady dotyczące rozwiązywania typowych problemów podczas pracy z grafami przepływu danych.

Nie można odnaleźć punktu końcowego rejestru

Jeśli nie można uruchomić grafu przepływu danych i zgłasza, że nie może znaleźć punktu końcowego rejestru, sprawdź następujące kwestie:

  1. Nazwa punktu końcowego rejestru musi być zgodna: wartość w registryEndpointRef twojego wykresu przepływu danych musi dokładnie pasować do name twojego zasobu RegistryEndpoint. Sprawdź, czy nie ma literówek i czy zachowano wielkość liter.

    # List all registry endpoints in the namespace
    kubectl get registryendpoints -n azure-iot-operations
    
  2. Punkt końcowy rejestru znajduje się w prawidłowej przestrzeni nazw: punkt końcowy rejestru musi znajdować się w azure-iot-operations przestrzeni nazw (lub w tej samej przestrzeni nazw co wykres przepływu danych).

  3. Punkt końcowy rejestru jest gotowy: Sprawdź stan punktu końcowego rejestru:

    kubectl describe registryendpoint $REGISTRY_ENDPOINT -n azure-iot-operations
    
  4. Uwierzytelnianie jest poprawnie konfigurowane: Jeśli używasz tożsamości zarządzanej, upewnij się, że rozszerzenie Operacje Azure IoT Arc ma uprawnienia AcrPull w rejestrze. Jeśli używasz uwierzytelniania anonimowego z publicznym rejestrem, sprawdź, czy adres URL hosta jest poprawny.

  5. Artefakty istnieją w rejestrze: Sprawdź, czy definicje grafu i moduły WASM, do których odwołuje się wykres, są dostępne w oczekiwanych tagach w rejestrze:

    # Check if artifacts exist (example with ORAS)
    oras manifest fetch $REGISTRY_HOST/graph-simple:1.0.0
    

Wykres przepływu danych jest uruchomiony, ale nie przetwarza danych

Jeśli wdrażasz graf przepływu danych, ale nie przetwarza komunikatów:

  1. Sprawdź stan grafu przepływu danych: poszukaj błędów w stanie zasobu grafu przepływu danych.

    kubectl get dataflowgraph $GRAPH_NAME -n azure-iot-operations -o yaml
    
  2. Sprawdź tematy MQTT: Upewnij się, że tematy źródłowe na wykresie przepływu danych są zgodne z tematami, w których są publikowane dane.

  3. Sprawdź znaczniki czasu: wykresy przepływu danych używają sygnatur czasowych hybrydowego zegara logicznego (HLC) do przetwarzania komunikatów. Uwzględnij właściwość __ts użytkownika podczas publikowania wiadomości MQTT, aby zapewnić terminowe przetwarzanie.