Alguns cenários exigem que as mensagens cheguem sobre diferentes tópicos MQTT, dependendo do seu conteúdo. Por exemplo, leituras dos sensores acima de um limiar crítico podem precisar de ser enviadas para um alerts tópico, enquanto leituras normais vão para um historian tópico. Com gráficos de fluxo de dados, pode definir o tópico de saída dinamicamente, mesmo que o fluxo de dados tenha um único destino.
O encaminhamento dinâmico de tópicos é uma técnica baseada na transformação de mapa: uma regra de mapa escreve os metadados do tópico alvo para mensagem, e o destino publica para esse tópico. Para encaminhar mensagens por diferentes caminhos de processamento dentro do grafo, veja Filtrar, desviar e mesclar dados.
Para uma visão geral dos gráficos de fluxo de dados e de como as transformações se compõem num pipeline, consulte Visão Geral dos Gráficos de Fluxo de Dados.
Pré-requisitos
- Um endpoint de registo predefinido chamado
default que aponta mcr.microsoft.com é criado automaticamente durante a implementação. As transformações incorporadas usam este endpoint.
Os CLI do Azure exemplos deste artigo usam variáveis de ambiente para que possas definir cada valor uma vez e depois copiar e colar os comandos as-is. Se estiver a usar o ambiente Operações IoT do Azure Codespaces do quickstart, estas variáveis já estão definidas para si e pode saltar este passo. Caso contrário, defina as seguintes variáveis de ambiente no seu shell antes de executar os comandos.
Os seguintes scripts definem as variáveis de ambiente mais usadas:
| Variável de ambiente |
Description |
SUBSCRIPTION_ID |
O ID da subscrição que contém a sua instância Operações IoT do Azure. |
RESOURCE_GROUP |
O nome do grupo de recursos que contém a sua instância do Operações IoT do Azure. |
AIO_INSTANCE_NAME |
O nome da sua instância do Operações IoT do Azure. Para listar as suas instâncias, execute az iot ops list -o table. |
CLUSTER_NAME |
O nome do cluster Kubernetes com Azure Arc que hospeda a sua instância. |
LOCATION |
A região do Azure deve ser usada para novos recursos, por exemplo 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>"
Só precisa de definir as variáveis que este artigo utiliza. Este artigo pode usar variáveis adicionais de ambiente para nomes de recursos que escolher. O artigo explica como posicioná-los onde são apresentados.
Como funciona o encaminhamento dinâmico de tópicos
Uma transformação de mapa pode escrever para metadados de mensagens, incluindo o tópico MQTT, usando o $metadata.topic caminho de saída. O destino usa então a ${outputTopic} variável para publicar para o tópico que a transformação definiu.
Duas peças funcionam juntas:
-
Dentro da transformação: Uma regra de mapeamento escreve um valor de cadeia de caracteres em
$metadata.topic.
-
No destino: O
dataDestination campo faz referência a ${outputTopic}, que se resolve no valor que a transformação escreveu.
As transformações utilizam uma linguagem de expressões para calcular valores, condições de teste e campos de referência. Expressões referem-se a entradas por posição, não pelo nome: a primeira entrada na inputs lista é $1, a segunda é $2, e assim sucessivamente. Funções incorporadas como cToF convertem e manipulam esses valores.
Para a lista completa de operadores, funções, tipos de dados e campos de metadados, consulte a referência Expressões.
Este artigo escreve para metadados de mensagens. Para os caminhos de metadados que pode ler e escrever, consulte campos de metadados.
A abordagem mais simples usa uma transformação de mapa com uma if expressão que seleciona o tema.
Na experiência de Operações, crie um grafo de fluxo de dados:
- Adicione uma fonte que leia de
sensors/temperature.
- Adicione uma transformação de mapa com duas regras:
- Uma regra de passagem coringa (entrada
*, saída *).
- Uma regra de cálculo com entrada
temperature, saída $metadata.topic, e expressão if($1 > 1000, "alerts", "historian").
- Adicione um destino com o tema
factory/${outputTopic}.
Quando a transformação do mapa escreve "alerts" em $metadata.topic, o destino resolve factory/${outputTopic} em factory/alerts.
A CLI do Azure utiliza um grafo de fluxo de dados a partir de um único ficheiro de configuração JSON. Cria um graph.json ficheiro com as propriedades do grafo. No graph.json ficheiro, cada transformação armazena as suas regras no value campo como uma string JSON escapada. Para a forma legível das regras de cada transformação, consulte o artigo prático para esse tipo de transformação.
{
"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"
}
}
]
}
Sugestão
Para gerar a string escapada, guarde as regras num ficheiro como rules.json, execute jq -c . rules.json, e cole a saída de linha única no value campo.
Aplica o ficheiro de configuração.
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
A utilização dos manifestos de implementação do Kubernetes não é suportada em ambientes de produção e deve ser usada apenas para depuração e testes.
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 }
Opção 2: Rota com ramificação, mapas por percurso e uma fusão
Se precisares de transformações diferentes em cada caminho (não apenas um tópico diferente), usa uma transformação de ramificação para separar o fluxo, uma transformação de mapa em cada ramo para definir o tópico e aplicar regras específicas de cada caminho, e uma transformação de concatenação para unir os caminhos.
Na experiência de operações:
- Adicione uma fonte que leia de
sensors/temperature.
- Adiciona uma transformada de ramificação com condição
$1 > 1000 no temperature campo.
- No caminho verdadeiro , adicione uma transformação de mapa com uma passagem de curinga e uma regra que defina
$metadata.topic como "alerts".
- No caminho falso , adicione uma transformação de mapa com uma passagem coringa e uma regra que defina
$metadata.topic para "historian".
- Adicione uma transformação de concatenação para unir os dois caminhos.
- Adicione um destino com o tema
factory/${outputTopic}.
A CLI do Azure utiliza um grafo de fluxo de dados a partir de um único ficheiro de configuração JSON. Cria um graph.json ficheiro com as propriedades do grafo. No graph.json ficheiro, cada transformação armazena as suas regras no value campo como uma string JSON escapada. Para a forma legível das regras de cada transformação, consulte o artigo prático para esse tipo de transformação.
{
"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"
}
}
]
}
Aplica o ficheiro de configuração.
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
A utilização dos manifestos de implementação do Kubernetes não é suportada em ambientes de produção e deve ser usada apenas para depuração e testes.
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 }
| Consideração |
Opção 1 (mapa único) |
Opção 2 (ramificação + mapas) |
| Simplicidade |
Menos nós, mais simples de ler |
Mais nós, mais explícito |
| Encaminhamento baseado apenas em tópicos |
Ideal |
Funciona, mas há mais preparação do que o necessário |
| Transformações diferentes para cada caminho |
Possível com aninhamentos if(), torna-se complexo |
Natural: cada ramo tem as suas próprias regras de mapa |
| Adicionar mais caminhos |
Chamadas em cadeia if() |
Requer ramos aninhados |
Para um encaminhamento direto de tópicos baseado numa única condição, a opção 1 é mais simples. Use a opção 2 quando cada caminho precisar de processamento diferente para além do nome do tópico.
Como a variável outputTopic resolve o tópico de destino
A variável ${outputTopic} em dataDestination resolve para o valor completo de $metadata.topic conforme definido pela última transformação do pipeline. Também pode usar segmentos com ${outputTopic.N} (1-indexado). Por exemplo, se a transformada define $metadata.topic para "region/west":
dataDestination |
Tema resolvido |
factory/${outputTopic} |
factory/region/west |
factory/${outputTopic.1} |
factory/region |
factory/${outputTopic.2} |
factory/west |
Se o fluxo de dados não conseguir resolver a variável tópico (por exemplo, $metadata.topic nunca foi definida), deixa a mensagem cair e regista um erro.
Conteúdo relacionado