Alguns cenários exigem que as mensagens cheguem em tópicos MQTT diferentes, dependendo do conteúdo. Por exemplo, as leituras de sensor acima de um limite crítico podem precisar ir para um alerts tópico, enquanto as leituras normais vão para um historian tópico. Com grafos de fluxo de dados, você pode definir o tópico de saída dinamicamente, mesmo que o fluxo de dados tenha um único destino.
O roteamento dinâmico de tópicos é uma técnica construída sobre a transformação do mapa: uma regra de mapa escreve os metadados do tópico alvo para a mensagem, e o destino publica para esse tópico. Para direcionar mensagens por diferentes caminhos de processamento dentro do grafo, veja Filtrar, desviar e mesclar dados.
Para obter uma visão geral dos grafos de fluxo de dados e como as transformações compõem em um pipeline, consulte a visão geral dos grafos de fluxo de dados.
Prerequisites
- Uma instância de Operações do Azure IoT implantada em um cluster do Kubernetes. Para obter mais informações, consulte Deploy Operações do Azure IoT.
- Um endpoint padrão chamado
default que aponta para mcr.microsoft.com é criado automaticamente durante a implantação. As transformações internas usam esse ponto de extremidade.
Os CLI do Azure exemplos deste artigo usam variáveis de ambiente para que você possa definir cada valor uma vez e então copiar e colar os comandos as-is. Se você está usando o ambiente Operações do Azure IoT Codespaces do quickstart, essas variáveis já estão definidas para você e você pode pular essa etapa. 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 comumente usadas:
| Variável de ambiente |
Description |
SUBSCRIPTION_ID |
O ID da assinatura que contém sua instância do Operações do Azure IoT. |
RESOURCE_GROUP |
O nome do grupo de recursos que contém sua instância do Operações do Azure IoT. |
AIO_INSTANCE_NAME |
O nome da sua instância do Operações do Azure IoT. Para listar suas instâncias, execute az iot ops list -o table. |
CLUSTER_NAME |
O nome do cluster Kubernetes habilitado para Azure Arc que hospeda sua instância. |
LOCATION |
A região Azure para usar para novos recursos, por exemploeastus, . |
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>"
Você só precisa definir as variáveis que este artigo utiliza. Este artigo pode usar variáveis adicionais de ambiente para nomes de recursos que você escolher. O artigo explica como posicioná-los onde são apresentados.
Como funciona o roteamento dinâmico de tópicos
Uma transformação de mapa pode gravar em metadados de mensagem, incluindo o tópico MQTT, por meio do caminho de saída $metadata.topic. Em seguida, o destino usa a variável ${outputTopic} para publicar no tópico definido pela transformação.
Duas partes funcionam juntas:
-
Dentro da transformação: uma regra de mapa grava um valor de cadeia de caracteres em
$metadata.topic.
-
No destino: o campo
dataDestination se refere a ${outputTopic}, que resulta no valor que a transformação registrou.
Transformadas usam uma linguagem de expressão 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 por diante. Funções integradas como cToF convertem e manipulam esses valores.
Para a lista completa de operadores, funções, tipos de dados e campos de metadados, veja a referência Expressions.
Este artigo fala sobre metadados de mensagens. Para os caminhos de metadados que você pode ler e escrever, veja Campos de Metadados.
A abordagem mais simples usa uma transformação de mapa com uma expressão if que escolhe o tópico.
Na experiência de Operações, crie um grafo de fluxo de dados:
- Adicionar uma origem que lê de
sensors/temperature.
- Adicione uma transformação de mapa com duas regras:
- Uma regra de passagem curinga (entrada
*, saída *).
- Uma regra de computação com entrada
temperature, saída $metadata.topice expressão if($1 > 1000, "alerts", "historian").
- Adicionar um destino com o tópico
factory/${outputTopic}.
Quando a transformação do mapa escreve "alerts" para $metadata.topic, o destino resolve factory/${outputTopic} para factory/alerts.
A CLI do Azure usa um grafo de fluxo de dados a partir de um único arquivo de configuração JSON. Crie um graph.json arquivo com as propriedades do grafo. No graph.json arquivo, cada transformação armazena suas regras no value campo como uma string JSON escapada. Para a forma legível das regras de cada transformação, veja o artigo de instruções 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"
}
}
]
}
Dica
Para gerar a string escapada, salve as regras em um arquivo como rules.json, execute jq -c . rules.json e cole a saída em uma única linha no campo value.
Aplique o arquivo 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
O uso de manifestos de implantação do Kubernetes não tem suporte em ambientes de produção e só deve ser usado para depuração e teste.
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 caminho e uma fusão
Se você precisar de transformações diferentes em cada caminho (não apenas um tópico diferente), use uma transformação de ramificação para dividir o fluxo, uma transformação de mapeamento em cada braço para definir o tópico e aplicar regras específicas de caminho, e uma transformação de concatenação para mesclar os caminhos.
Na experiência em Operações:
- Adicionar uma origem que lê de
sensors/temperature.
- Adicione uma transformação de ramificação com condição
$1 > 1000 no temperature campo.
- No caminho verdadeiro, adicione uma transformação mapa com uma passagem curinga e uma regra que define
$metadata.topic como "alerts".
- No caminho verdadeiro, adicione a transformação mapa com uma passagem curinga e uma regra que define
$metadata.topic como "historian".
- Adicione uma transformação concatena para mesclar os dois caminhos.
- Adicionar um destino com o tópico
factory/${outputTopic}.
A CLI do Azure usa um grafo de fluxo de dados a partir de um único arquivo de configuração JSON. Crie um graph.json arquivo com as propriedades do grafo. No graph.json arquivo, cada transformação armazena suas regras no value campo como uma string JSON escapada. Para a forma legível das regras de cada transformação, veja o artigo de instruções 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"
}
}
]
}
Aplique o arquivo 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
O uso de manifestos de implantação do Kubernetes não tem suporte em ambientes de produção e só deve ser usado para depuração e teste.
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 (ramo + mapas) |
| Simplicidade |
Menos nós, mais fácil de ler |
Mais nós, mais explícitos |
| Roteamento exclusivo por tópico |
Ideal |
Funciona, mas mais configuração do que o necessário |
| Transformações distintas para cada caminho |
Possível com uma estrutura aninhada if(), fica complexo |
Natural: cada ramificação tem suas próprias regras de mapeamento |
| Adicionando mais caminhos |
Chamadas em if() cadeia |
Requer ramificações aninhadas |
Para o roteamento de tópico simples com base em uma única condição, a opção 1 é mais simples. Use a opção 2 quando cada caminho precisar de processamento diferente além do nome do tópico.
Como a variável outputTopic resolve o tópico de destino
A variável ${outputTopic} em dataDestination se resolve para o valor total de $metadata.topic conforme definido pela última transformação no pipeline. Você também pode usar segmentos com ${outputTopic.N} (indexado por 1). Por exemplo, se a transformação definir $metadata.topic como "region/west":
dataDestination |
Tópico 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 ativado), ele deixa a mensagem cair e registra um erro.
Conteúdo relacionado