Alcuni scenari richiedono l'arrivo di messaggi su argomenti MQTT diversi a seconda del contenuto. Ad esempio, le letture dei sensori al di sopra di una soglia critica devono passare a un alerts topic, mentre le letture normali passano a un historian topic. Con i grafici del flusso di dati, è possibile impostare l'argomento di output in modo dinamico, anche se il flusso di dati ha una singola destinazione.
Il routing dinamico dei topic è una tecnica basata sulla map transform: una regola mappa scrive l'argomento di destinazione in metadati del messaggio, e la destinazione pubblica su quell'argomento. Per instradare i messaggi lungo diversi percorsi di elaborazione all'interno del grafo, vedi Filtra, ramifica e unisci dati.
Per una panoramica dei grafici del flusso di dati e della composizione delle trasformazioni in una pipeline, vedere Panoramica dei grafici del flusso di dati.
Prerequisiti
- Un endpoint del Registro di sistema predefinito denominato
default che punta a mcr.microsoft.com viene creato automaticamente durante la distribuzione. Le trasformazioni predefinite usano questo endpoint.
Gli esempi interfaccia della riga di comando di Azure in questo articolo usano variabili di ambiente così puoi impostare ogni valore una volta e poi copiare e incollare i comandi as-is. Se stai usando l'ambiente Operazioni di Azure IoT Codespaces dal quickstart, queste variabili sono già impostate per te e puoi saltare questo passaggio. Altrimenti, imposta le seguenti variabili di ambiente nella tua shell prima di eseguire i comandi.
I seguenti script impostano le variabili di ambiente più comunemente utilizzate:
| Variabile di ambiente |
Description |
SUBSCRIPTION_ID |
L'ID dell'abbonamento che contiene la tua istanza Operazioni di Azure IoT. |
RESOURCE_GROUP |
Il nome del gruppo di risorse che contiene la tua istanza Operazioni di Azure IoT. |
AIO_INSTANCE_NAME |
Il nome della tua istanza Operazioni di Azure IoT. Per elencare le tue istanze, esegui az iot ops list -o table. |
CLUSTER_NAME |
Il nome del cluster Kubernetes abilitato Azure Arc che ospita la tua istanza. |
LOCATION |
La regione Azure da utilizzare per nuove risorse, ad esempio eastus. |
SUBSCRIPTION_ID=<subscription-id>
RESOURCE_GROUP=<resource-group-name>
AIO_INSTANCE_NAME=<instance-name>
CLUSTER_NAME=<cluster-name>
LOCATION=<region>
$SUBSCRIPTION_ID = "<subscription-id>"
$RESOURCE_GROUP = "<resource-group-name>"
$AIO_INSTANCE_NAME = "<instance-name>"
$CLUSTER_NAME = "<cluster-name>"
$LOCATION = "<region>"
Devi solo impostare le variabili utilizzate in questo articolo. Questo articolo potrebbe utilizzare variabili ambientali aggiuntive per i nomi delle risorse che scegli. L'articolo spiega come posizionarli dove vengono introdotti.
Come funziona l'instradamento dinamico degli argomenti
Una trasformazione mappa può scrivere nei metadati dei messaggi, incluso l'argomento MQTT, usando il percorso di output $metadata.topic. Quindi, la destinazione usa la variabile ${outputTopic} per pubblicare in qualsiasi argomento impostato dalla trasformazione.
Due pezzi interagiscono:
-
All'interno della trasformazione: una regola della mappa scrive un valore stringa in
$metadata.topic.
-
Nella destinazione: il campo
dataDestination si riferisce a ${outputTopic}, che viene risolto nel valore scritto dalla trasformazione.
Le trasformazioni utilizzano un linguaggio di espressione per calcolare valori, condizioni di test e campi di riferimento. Le espressioni si riferiscono agli input per posizione, non per nome: il primo input nella inputs lista è $1, il secondo è $2, e così via. Funzioni integrate come cToF convertono e manipolano tali valori.
Per l'elenco completo di operatori, funzioni, tipi di dati e campi di metadati, consulta il riferimento Expressions.
Questo articolo scrive nei metadati dei messaggi. Per i percorsi di metadati che puoi leggere e scrivere, vedi i campi Metadata.
L'approccio più semplice utilizza una trasformazione della mappa con un'espressione if che seleziona l'argomento.
Nell'esperienza Operazioni creare un grafico del flusso di dati:
- Aggiungere una sorgente che legga da
sensors/temperature.
- Aggiungere una trasformazione mappa con due regole:
- Regola di trasmissione diretta con caratteri jolly (input
*, output *).
- Regola di calcolo con input
temperature, output $metadata.topiced espressione if($1 > 1000, "alerts", "historian").
- Aggiungere una destinazione con l'argomento
factory/${outputTopic}.
Quando la trasformazione della mappa scrive "alerts" in $metadata.topic, la destinazione risolve factory/${outputTopic} in factory/alerts.
L'interfaccia della riga di comando di Azure utilizza un grafo di flusso dati da un singolo file di configurazione JSON. Creare un graph.json file con le proprietà del grafo. Nel graph.json file, ogni trasformazione memorizza le proprie regole nel value campo come stringa JSON sfuggita. Per la forma leggibile delle regole di ciascuna trasformazione, vedi l'articolo pratico per quel tipo di trasformazione.
{
"mode": "Enabled",
"nodes": [
{
"nodeType": "Source",
"name": "sensors",
"sourceSettings": {
"endpointRef": "default",
"dataSources": [
"sensors/temperature"
]
}
},
{
"nodeType": "Graph",
"name": "route-by-temperature",
"graphSettings": {
"registryEndpointRef": "default",
"artifact": "azureiotoperations/graph-dataflow-map:1.0.0",
"configuration": [
{
"key": "rules",
"value": "{\"map\":[{\"inputs\":[\"*\"],\"output\":\"*\"},{\"description\":\"Set topic based on temperature threshold\",\"inputs\":[\"temperature\"],\"output\":\"$metadata.topic\",\"expression\":\"if($1 > 1000, \\\"alerts\\\", \\\"historian\\\")\"}]}"
}
]
}
},
{
"nodeType": "Destination",
"name": "output",
"destinationSettings": {
"endpointRef": "default",
"dataDestination": "factory/${outputTopic}"
}
}
],
"nodeConnections": [
{
"from": {
"name": "sensors"
},
"to": {
"name": "route-by-temperature"
}
},
{
"from": {
"name": "route-by-temperature"
},
"to": {
"name": "output"
}
}
]
}
Suggerimento
Per generare la stringa sfuggita, salva le regole in un file come rules.json, esegui jq -c . rules.json, e incolla l'output a singola riga nel value campo.
Applicare il file di configurazione.
az iot ops dataflowgraph apply \
--name dynamic-topic-routing \
--instance $AIO_INSTANCE_NAME \
--resource-group $RESOURCE_GROUP \
--config-file graph.json
resource dataflowGraph 'Microsoft.IoTOperations/instances/dataflowProfiles/dataflowGraphs@2026-03-01' = {
name: 'dynamic-topic-routing'
parent: dataflowProfile
properties: {
profileRef: dataflowProfileName
mode: 'Enabled'
nodes: [
{
nodeType: 'Source'
name: 'sensors'
sourceSettings: {
endpointRef: 'default'
dataSources: [ 'sensors/temperature' ]
}
}
{
nodeType: 'Graph'
name: 'route-by-temperature'
graphSettings: {
registryEndpointRef: 'default'
artifact: 'azureiotoperations/graph-dataflow-map:1.0.0'
configuration: [
{
key: 'rules'
value: '{"map":[{"inputs":["*"],"output":"*"},{"description":"Set topic based on temperature threshold","inputs":["temperature"],"output":"$metadata.topic","expression":"if($1 > 1000, \\"alerts\\", \\"historian\\")"}]}'
}
]
}
}
{
nodeType: 'Destination'
name: 'output'
destinationSettings: {
endpointRef: 'default'
dataDestination: 'factory/${outputTopic}'
}
}
]
nodeConnections: [
{ from: { name: 'sensors' }, to: { name: 'route-by-temperature' } }
{ from: { name: 'route-by-temperature' }, to: { name: 'output' } }
]
}
}
Importante
L'uso dei manifesti di distribuzione Kubernetes non è supportato negli ambienti di produzione e deve essere usato solo per il debug e il test.
apiVersion: connectivity.iotoperations.azure.com/v1
kind: DataflowGraph
metadata:
name: dynamic-topic-routing
namespace: azure-iot-operations
spec:
profileRef: default
nodes:
- nodeType: Source
name: sensors
sourceSettings:
endpointRef: default
dataSources:
- sensors/temperature
- nodeType: Graph
name: route-by-temperature
graphSettings:
registryEndpointRef: default
artifact: azureiotoperations/graph-dataflow-map:1.0.0
configuration:
- key: rules
value: |
{
"map": [
{
"inputs": ["*"],
"output": "*"
},
{
"description": "Set topic based on temperature threshold",
"inputs": ["temperature"],
"output": "$metadata.topic",
"expression": "if($1 > 1000, \"alerts\", \"historian\")"
}
]
}
- nodeType: Destination
name: output
destinationSettings:
endpointRef: default
dataDestination: "factory/${outputTopic}"
nodeConnections:
- from: { name: sensors }
to: { name: route-by-temperature }
- from: { name: route-by-temperature }
to: { name: output }
Opzione 2: Percorso con diramazione, mappe per percorso e un'immergibile
Se sono necessarie trasformazioni diverse in ogni percorso (non solo in un argomento diverso), usare una trasformazione di diramazione per suddividere il flusso, una trasformazione di mappatura in ogni ramo per impostare l'argomento e applicare regole specifiche del percorso e una trasformazione di concatenazione per unire i percorsi.
Nell'esperienza operativa:
- Aggiungere una sorgente che legga da
sensors/temperature.
- Aggiungere una trasformazione di ramo con condizione
$1 > 1000 nel temperature campo.
- Nel percorso true aggiungere una trasformazione mappa con un pass-through con caratteri jolly e una regola che imposta
$metadata.topic su "alerts".
- Nel percorso false aggiungere una trasformazione mappa con un pass-through con caratteri jolly e una regola che imposta
$metadata.topic su "historian".
- Aggiungere una trasformazione concatenata per unire entrambi i percorsi.
- Aggiungere una destinazione con l'argomento
factory/${outputTopic}.
L'interfaccia della riga di comando di Azure utilizza un grafo di flusso dati da un singolo file di configurazione JSON. Creare un graph.json file con le proprietà del grafo. Nel graph.json file, ogni trasformazione memorizza le proprie regole nel value campo come stringa JSON sfuggita. Per la forma leggibile delle regole di ciascuna trasformazione, vedi l'articolo pratico per quel tipo di trasformazione.
{
"mode": "Enabled",
"nodes": [
{
"nodeType": "Source",
"name": "sensors",
"sourceSettings": {
"endpointRef": "default",
"dataSources": [
"sensors/temperature"
]
}
},
{
"nodeType": "Graph",
"name": "check-temperature",
"graphSettings": {
"registryEndpointRef": "default",
"artifact": "azureiotoperations/graph-dataflow-branch:1.0.0",
"configuration": [
{
"key": "rules",
"value": "{\"branch\":{\"inputs\":[\"temperature\"],\"expression\":\"$1 > 1000\",\"description\":\"Route critical temperatures to alerts\"}}"
}
]
}
},
{
"nodeType": "Graph",
"name": "set-alerts-topic",
"graphSettings": {
"registryEndpointRef": "default",
"artifact": "azureiotoperations/graph-dataflow-map:1.0.0",
"configuration": [
{
"key": "rules",
"value": "{\"map\":[{\"inputs\":[\"*\"],\"output\":\"*\"},{\"inputs\":[],\"output\":\"$metadata.topic\",\"expression\":\"\\\"alerts\\\"\"}]}"
}
]
}
},
{
"nodeType": "Graph",
"name": "set-historian-topic",
"graphSettings": {
"registryEndpointRef": "default",
"artifact": "azureiotoperations/graph-dataflow-map:1.0.0",
"configuration": [
{
"key": "rules",
"value": "{\"map\":[{\"inputs\":[\"*\"],\"output\":\"*\"},{\"inputs\":[],\"output\":\"$metadata.topic\",\"expression\":\"\\\"historian\\\"\"}]}"
}
]
}
},
{
"nodeType": "Graph",
"name": "merge",
"graphSettings": {
"registryEndpointRef": "default",
"artifact": "azureiotoperations/graph-dataflow-concatenate:1.0.0"
}
},
{
"nodeType": "Destination",
"name": "output",
"destinationSettings": {
"endpointRef": "default",
"dataDestination": "factory/${outputTopic}"
}
}
],
"nodeConnections": [
{
"from": {
"name": "sensors"
},
"to": {
"name": "check-temperature"
}
},
{
"from": {
"name": "check-temperature.output.true"
},
"to": {
"name": "set-alerts-topic"
}
},
{
"from": {
"name": "check-temperature.output.false"
},
"to": {
"name": "set-historian-topic"
}
},
{
"from": {
"name": "set-alerts-topic"
},
"to": {
"name": "merge"
}
},
{
"from": {
"name": "set-historian-topic"
},
"to": {
"name": "merge"
}
},
{
"from": {
"name": "merge"
},
"to": {
"name": "output"
}
}
]
}
Applicare il file di configurazione.
az iot ops dataflowgraph apply \
--name dynamic-topic-routing-branched \
--instance $AIO_INSTANCE_NAME \
--resource-group $RESOURCE_GROUP \
--config-file graph.json
resource dataflowGraph 'Microsoft.IoTOperations/instances/dataflowProfiles/dataflowGraphs@2026-03-01' = {
name: 'dynamic-topic-routing-branched'
parent: dataflowProfile
properties: {
profileRef: dataflowProfileName
mode: 'Enabled'
nodes: [
{
nodeType: 'Source'
name: 'sensors'
sourceSettings: {
endpointRef: 'default'
dataSources: [ 'sensors/temperature' ]
}
}
{
nodeType: 'Graph'
name: 'check-temperature'
graphSettings: {
registryEndpointRef: 'default'
artifact: 'azureiotoperations/graph-dataflow-branch:1.0.0'
configuration: [
{
key: 'rules'
value: '{"branch":{"inputs":["temperature"],"expression":"$1 > 1000","description":"Route critical temperatures to alerts"}}'
}
]
}
}
{
nodeType: 'Graph'
name: 'set-alerts-topic'
graphSettings: {
registryEndpointRef: 'default'
artifact: 'azureiotoperations/graph-dataflow-map:1.0.0'
configuration: [
{
key: 'rules'
value: '{"map":[{"inputs":["*"],"output":"*"},{"inputs":[],"output":"$metadata.topic","expression":"\\"alerts\\""}]}'
}
]
}
}
{
nodeType: 'Graph'
name: 'set-historian-topic'
graphSettings: {
registryEndpointRef: 'default'
artifact: 'azureiotoperations/graph-dataflow-map:1.0.0'
configuration: [
{
key: 'rules'
value: '{"map":[{"inputs":["*"],"output":"*"},{"inputs":[],"output":"$metadata.topic","expression":"\\"historian\\""}]}'
}
]
}
}
{
nodeType: 'Graph'
name: 'merge'
graphSettings: {
registryEndpointRef: 'default'
artifact: 'azureiotoperations/graph-dataflow-concatenate:1.0.0'
}
}
{
nodeType: 'Destination'
name: 'output'
destinationSettings: {
endpointRef: 'default'
dataDestination: 'factory/${outputTopic}'
}
}
]
nodeConnections: [
{ from: { name: 'sensors' }, to: { name: 'check-temperature' } }
{ from: { name: 'check-temperature.output.true' }, to: { name: 'set-alerts-topic' } }
{ from: { name: 'check-temperature.output.false' }, to: { name: 'set-historian-topic' } }
{ from: { name: 'set-alerts-topic' }, to: { name: 'merge' } }
{ from: { name: 'set-historian-topic' }, to: { name: 'merge' } }
{ from: { name: 'merge' }, to: { name: 'output' } }
]
}
}
Importante
L'uso dei manifesti di distribuzione Kubernetes non è supportato negli ambienti di produzione e deve essere usato solo per il debug e il test.
apiVersion: connectivity.iotoperations.azure.com/v1
kind: DataflowGraph
metadata:
name: dynamic-topic-routing-branched
namespace: azure-iot-operations
spec:
profileRef: default
nodes:
- nodeType: Source
name: sensors
sourceSettings:
endpointRef: default
dataSources:
- sensors/temperature
- nodeType: Graph
name: check-temperature
graphSettings:
registryEndpointRef: default
artifact: azureiotoperations/graph-dataflow-branch:1.0.0
configuration:
- key: rules
value: |
{
"branch": {
"inputs": ["temperature"],
"expression": "$1 > 1000",
"description": "Route critical temperatures to alerts"
}
}
- nodeType: Graph
name: set-alerts-topic
graphSettings:
registryEndpointRef: default
artifact: azureiotoperations/graph-dataflow-map:1.0.0
configuration:
- key: rules
value: |
{
"map": [
{ "inputs": ["*"], "output": "*" },
{ "inputs": [], "output": "$metadata.topic", "expression": "\"alerts\"" }
]
}
- nodeType: Graph
name: set-historian-topic
graphSettings:
registryEndpointRef: default
artifact: azureiotoperations/graph-dataflow-map:1.0.0
configuration:
- key: rules
value: |
{
"map": [
{ "inputs": ["*"], "output": "*" },
{ "inputs": [], "output": "$metadata.topic", "expression": "\"historian\"" }
]
}
- nodeType: Graph
name: merge
graphSettings:
registryEndpointRef: default
artifact: azureiotoperations/graph-dataflow-concatenate:1.0.0
- nodeType: Destination
name: output
destinationSettings:
endpointRef: default
dataDestination: "factory/${outputTopic}"
nodeConnections:
- from: { name: sensors }
to: { name: check-temperature }
- from: { name: check-temperature.output.true }
to: { name: set-alerts-topic }
- from: { name: check-temperature.output.false }
to: { name: set-historian-topic }
- from: { name: set-alerts-topic }
to: { name: merge }
- from: { name: set-historian-topic }
to: { name: merge }
- from: { name: merge }
to: { name: output }
| Considerazione |
Opzione 1 (singola mappa) |
Opzione 2 (ramo + mappe) |
| Semplicità |
Meno nodi, più semplice da leggere |
Altri nodi, più espliciti |
| Indirizzamento basato su argomenti |
Ideal |
Funziona, ma più configurazione del necessario |
| Trasformazioni diverse per percorso |
È possibile con if() annidato, diventa complesso. |
Naturale: ogni ramo ha regole di mappa proprie |
| Aggiunta di altri percorsi |
Chiamate a catena if() |
Richiede rami annidati |
Per semplificare il routing degli argomenti in base a una singola condizione, l'opzione 1 è più semplice. Usare l'opzione 2 quando ogni percorso richiede un'elaborazione diversa rispetto al nome dell'argomento.
Come la variabile outputTopic risolve l'argomento di destinazione
La ${outputTopic} variabile in dataDestination viene risolta nel valore completo di $metadata.topic come impostato dall'ultima trasformazione nella pipeline. È anche possibile usare segmenti con ${outputTopic.N} (1 indicizzato). Ad esempio, se la trasformazione imposta $metadata.topic su "region/west":
dataDestination |
Argomento risolto |
factory/${outputTopic} |
factory/region/west |
factory/${outputTopic.1} |
factory/region |
factory/${outputTopic.2} |
factory/west |
Se il flusso di dati non riesce a risolvere la variabile topic (ad esempio, $metadata.topic non è mai stata impostata), lascia cadere il messaggio e registra un errore.
Contenuti correlati