Algunos escenarios requieren que los mensajes lleguen a diferentes tópicos MQTT dependiendo de su contenido. Por ejemplo, las lecturas del sensor por encima de un umbral crítico pueden necesitar ir a un tema alerts, mientras que las lecturas normales van a un tema historian. Con los gráficos de flujo de datos, puede establecer el tema de salida dinámicamente, aunque el flujo de datos tenga un único destino.
El enrutamiento dinámico de temas es una técnica basada en la transformación de mapa: una regla de mapa escribe el tema destino a metadatos de mensaje, y el destino publica en ese tema. Para enrutar los mensajes por diferentes rutas de procesamiento dentro del grafo, consulta Filtrar, bifurcar y fusionar datos.
Para obtener información general sobre los gráficos de flujo de datos y cómo las transformaciones se componen en una canalización, consulte Introducción a los gráficos de flujo de datos.
Prerequisites
- Un punto de conexión del Registro predeterminado denominado
default que apunta a mcr.microsoft.com se crea automáticamente durante la implementación. Las transformaciones integradas usan este endpoint.
Los CLI de Azure ejemplos de este artículo usan variables de entorno para que puedas establecer cada valor una vez y luego copiar y pegar los comandos as-is. Si usas el entorno Operaciones de IoT de Azure Codespaces del quickstart, estas variables ya están configuradas para ti y puedes saltarte este paso. De lo contrario, configura las siguientes variables de entorno en tu shell antes de ejecutar los comandos.
Los siguientes scripts establecen las variables de entorno más usadas:
| Variable del entorno |
Description |
SUBSCRIPTION_ID |
El ID de la suscripción que contiene tu instancia de Operaciones de IoT de Azure. |
RESOURCE_GROUP |
El nombre del grupo de recursos que contiene tu instancia de Operaciones de IoT de Azure. |
AIO_INSTANCE_NAME |
El nombre de tu instancia de Operaciones de IoT de Azure. Para listar tus instancias, ejecuta az iot ops list -o table. |
CLUSTER_NAME |
El nombre del clúster Kubernetes habilitado para Azure Arc que aloja tu instancia. |
LOCATION |
La región de Azure para usar para nuevos recursos, por ejemplo 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>"
Solo necesitas establecer las variables que utiliza este artículo. Este artículo podría utilizar variables de entorno adicionales para los nombres de recursos que elijas. El artículo explica cómo situarlos donde se introducen.
Cómo funciona el enrutamiento dinámico de temas
Una transformación de mapa puede escribir en los metadatos del mensaje, incluido el tema MQTT, utilizando la ruta de acceso de salida $metadata.topic. El destino usa la variable ${outputTopic} para publicar en el tema que el conjunto de transformaciones establezca.
Dos piezas funcionan juntas:
-
Dentro de la transformación: una regla de asignación escribe un valor de cadena en
$metadata.topic.
-
En el destino: el campo
dataDestination hace referencia a ${outputTopic}, que se resuelve en el valor escrito por la transformación.
Las transformadas utilizan un lenguaje de expresiones para calcular valores, condiciones de prueba y campos de referencia. Las expresiones se refieren a las entradas por posición, no por nombre: la primera entrada en la inputs lista es $1, la segunda es $2, y así sucesivamente. Las funciones integradas, como cToF, convierten y manipulan esos valores.
Para la lista completa de operadores, funciones, tipos de datos y campos de metadatos, consulte la referencia Expressions.
Este artículo escribe para los metadatos de los mensajes. Para las rutas de metadatos que puedes leer y escribir, consulta los campos de metadatos.
El enfoque más sencillo usa una transformación de mapa con una if expresión que elige el tema.
En la experiencia de operaciones, cree un gráfico de flujo de datos:
- Agregue un origen que lea de
sensors/temperature.
- Agregue una transformación de mapa con dos reglas:
- Una regla de tránsito con caracteres comodín (entrada
*, salida *).
- Una regla de proceso con entrada
temperature, salida $metadata.topicy expresión if($1 > 1000, "alerts", "historian").
- Agregue un destino con el tema
factory/${outputTopic}.
Cuando la transformación de mapa escribe "alerts" a $metadata.topic, el destino resuelve factory/${outputTopic} a factory/alerts.
La CLI de Azure utiliza un grafo de flujo de datos de un único archivo de configuración JSON. Cree un graph.json archivo con las propiedades del grafo. En el graph.json archivo, cada transformación almacena sus reglas en el value campo como una cadena JSON escapada. Para la forma legible de las reglas de cada transformada, consulta el artículo práctico para ese tipo de transformación.
{
"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"
}
}
]
}
Tip
Para generar la cadena escapada, guarda las reglas en un archivo como rules.json, ejecuta jq -c . rules.json, y pega la salida de una sola línea en el value campo.
Aplique el archivo de configuración.
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
El uso de manifiestos de implementación de Kubernetes no se admite en entornos de producción y solo se debe usar para la depuración y las pruebas.
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 }
Opción 2: Ruta con rama, mapas por ruta y una fusión
Si necesita diferentes transformaciones en cada ruta de acceso (no solo un tema diferente), use una transformación de rama para dividir el flujo, una transformación de mapa en cada brazo para establecer el tema y aplicar reglas específicas de la ruta de acceso y una transformación concatenada para combinar las rutas de acceso.
En la experiencia de operaciones:
- Agregue un origen que lea de
sensors/temperature.
- Agregue una transformación de rama con condición
$1 > 1000 en el campo temperature.
- En la ruta true, agregue una transformación de mapa con un tránsito con caracteres comodín y una regla que establezca
$metadata.topic en "alerts".
- En la ruta false, agregue una transformación de mapa con un tránsito con caracteres comodín y una regla que establezca
$metadata.topic en "historian".
- Agregue una transformación concatenada para combinar ambas rutas de acceso.
- Agregue un destino con el tema
factory/${outputTopic}.
La CLI de Azure utiliza un grafo de flujo de datos de un único archivo de configuración JSON. Cree un graph.json archivo con las propiedades del grafo. En el graph.json archivo, cada transformación almacena sus reglas en el value campo como una cadena JSON escapada. Para la forma legible de las reglas de cada transformada, consulta el artículo práctico para ese tipo de transformación.
{
"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 el archivo de configuración.
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
El uso de manifiestos de implementación de Kubernetes no se admite en entornos de producción y solo se debe usar para la depuración y las pruebas.
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 }
| Consideración |
Opción 1 (mapa único) |
Opción 2 (rama + mapas) |
| Simplicidad |
Menos nodos, más fáciles de leer |
Más nodos, más explícitos |
| Enrutamiento exclusivamente por tema |
Ideal |
Funciona, pero más configuración de la necesaria |
| Diferentes transformaciones por trayectoria |
Posible con estructura anidada if(), se vuelve complejo |
Natural: cada rama tiene sus propias reglas de mapa |
| Adición de más rutas de acceso |
Cadenas if() de llamadas |
Requiere ramas anidadas |
Para el enrutamiento sencillo de temas basado en una sola condición, la opción 1 es más sencilla. Use la opción 2 cuando cada ruta de acceso necesite un procesamiento diferente más allá del nombre del tema.
Cómo la variable outputTopic resuelve el tema destino
La variable ${outputTopic} en dataDestination se resuelve como el valor completo de $metadata.topic tal como lo define la última transformación en la canalización. También puede usar segmentos con ${outputTopic.N} (índice 1). Por ejemplo, si la transformación establece $metadata.topic a "region/west":
dataDestination |
Tema resuelto |
factory/${outputTopic} |
factory/region/west |
factory/${outputTopic.1} |
factory/region |
factory/${outputTopic.2} |
factory/west |
Si el flujo de datos no puede resolver la variable de tema (por ejemplo, $metadata.topic nunca se activó), se pierde el mensaje y se registra un error.
Contenido relacionado