Encaminhar mensagens para diferentes tópicos MQTT em grafos de fluxo de dados

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>

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:

  1. Dentro da transformação: Uma regra de mapeamento escreve um valor de cadeia de caracteres em $metadata.topic.
  2. 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.

Opção 1: Roteamento com uma transformação de mapa único e uma expressão condicional

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:

  1. Adicione uma fonte que leia de sensors/temperature.
  2. 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").
  3. 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.

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:

  1. Adicione uma fonte que leia de sensors/temperature.
  2. Adiciona uma transformada de ramificação com condição $1 > 1000 no temperature campo.
  3. No caminho verdadeiro , adicione uma transformação de mapa com uma passagem de curinga e uma regra que defina $metadata.topic como "alerts".
  4. No caminho falso , adicione uma transformação de mapa com uma passagem coringa e uma regra que defina $metadata.topic para "historian".
  5. Adicione uma transformação de concatenação para unir os dois caminhos.
  6. Adicione um destino com o tema factory/${outputTopic}.

Escolha entre uma transformação de um único mapa e um ramo

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.