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

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>

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:

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

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

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:

  1. Adicionar uma origem que lê de sensors/temperature.
  2. 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").
  3. 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.

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:

  1. Adicionar uma origem que lê de sensors/temperature.
  2. Adicione uma transformação de ramificação com condição $1 > 1000 no temperature campo.
  3. No caminho verdadeiro, adicione uma transformação mapa com uma passagem curinga e uma regra que define $metadata.topic como "alerts".
  4. No caminho verdadeiro, adicione a transformação mapa com uma passagem curinga e uma regra que define $metadata.topic como "historian".
  5. Adicione uma transformação concatena para mesclar os dois caminhos.
  6. Adicionar um destino com o tópico factory/${outputTopic}.

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

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.