Filtrar, desviar e mesclar dados em grafos de fluxo de dados

Os grafos de fluxo de dados fornecem duas maneiras de controlar quais mensagens fluem pelo pipeline: as transformações de filtro descartam mensagens indesejadas e as transformações de ramificação roteiam cada mensagem para um dos dois caminhos com base em uma condição. Após a ramificação, uma transformação de concatenação une os caminhos novamente.

Essas transformações roteiam mensagens dentro do grafo. Para direcionar mensagens para diferentes tópicos MQTT com base em seu conteúdo, veja Encaminhar mensagens para diferentes tópicos MQTT.

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.

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.

Pré-requisitos

  • 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 integradas utilizam este 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 Descrição
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.

Transformação de filtro

Uma transformação de filtro avalia cada mensagem de entrada em relação a uma ou mais regras e decide se a mensagem continua por meio do pipeline ou é descartada.

Importante

Uma expressão de filtro seleciona as mensagens a remover, não as que devem ser mantidas. Quando a expressão é verdadeira, a mensagem é descartada. Esse comportamento é o oposto de uma expressão de mapa, que calcula um valor que é mantido.

Para manter as mensagens que correspondem a uma condição, inverta a expressão. Por exemplo, para manter apenas leituras acima de 90, filtre em $1 <= 90.

Como funcionam as regras de filtro

Cada regra de filtro tem estas propriedades:

Propriedade Obrigatório Descrição
inputs Sim Lista de caminhos de campos a serem lidos na mensagem recebida.
expression Sim Fórmula aplicada aos valores de entrada. Deve retornar um booliano. Quando retorna verdadeiro, a mensagem é descartada.
description No Rótulo legível por humanos usado em mensagens de erro.

Cada entrada mapeia para uma variável posicional baseada em sua ordem: a primeira entrada é $1, a segunda é $2, e assim por diante.

Quando você define várias regras, elas usam a lógica OR: se qualquer regra for avaliada como true, a mensagem será descartada. O mecanismo interrompe a execução assim que uma regra é encontrada.

Restrições de chave:

  • A expressão é necessária. Cada regra de filtro deve incluir um expression.
  • filter aceita uma matriz. Forneça regras como um array JSON, "filter": [ { ... } ], mesmo para uma única regra. Ao passar apenas um objeto simples, a transformação não é carregada, e o erro que surge aponta para o artefato e o registro, não para o conteúdo das regras. Essa restrição difere de branch, que toma um único objeto.
  • Nenhuma entrada curinga. Cada entrada deve fazer referência a um caminho de campo específico.
  • Campos ausentes causam erros. Se um campo referenciado em inputs não existir, o filtro retornará um erro, em vez de passar a mensagem silenciosamente.
  • Os resultados não boolianos causam erros. Se uma expressão retornar um valor não booliano (como uma cadeia de caracteres ou um número), o filtro retornará um erro.

Remover mensagens por condição

Para descartar mensagens em que a temperatura exceda 100 graus:

Na configuração de transformação de filtro, adicione uma regra:

Configurações Valor
Input temperature
Expressão $1 > 100

As mensagens em que a temperatura é igual ou inferior a 100 são aprovadas. Mensagens acima de 100 são descartadas.

Manter mensagens por condição

Frequentemente, você quer o resultado oposto: manter apenas as mensagens que correspondem a uma condição. Como uma expressão de filtro seleciona o que remover, inverta a comparação.

Para manter apenas as leituras acima de 90, descarte tudo que seja igual ou inferior a 90:

Na configuração de transformação de filtro, adicione uma regra:

Configurações Valor
Input temperature
Expressão $1 <= 90
Description Drop readings at or below 90

Apenas mensagens com valores acima de 90 continuam pelo pipeline. Escrever $1 > 90 aqui faria o oposto do que você quer: descartaria todas as leituras acima de 90 e manteria as mais baixas.

Dica

Use o campo description para registrar a intenção da regra quanto ao que ela descarta. Uma descrição como Drop readings at or below 90 permanece precisa, enquanto Keep hot readings convida ao erro de expressão invertida e aparece em mensagens de erro que depois são lidas ao contrário.

Usar várias condições

Quando você define mais de uma regra, o filtro descarta a mensagem se alguma regra corresponder:

Adicione duas regras:

Entrada Expression Descrição
temperature $1 > 100 Baixar alta temperatura
humidity $1 > 95 Reduza a umidade alta
Mensagem regra de temperatura Regra de umidade Resultado
{"temperature": 150, "humidity": 60} verdadeiro falso Dropped
{"temperature": 80, "humidity": 98} falso verdadeiro Dropped
{"temperature": 80, "humidity": 60} falso falso Passes

Dica

Use várias entradas em uma regra quando precisar de lógica AND entre campos. Use várias regras quando precisar de lógica OR em condições independentes.

Use expressões complexas

Faça referência a vários campos em uma única regra e combine-os com operadores lógicos:

Adicione uma regra com entradas temperature e humidityexpressão $1 > 30 && $2 < 60.

Para obter a lista completa de operadores e funções, consulte a referência de Expressões.

Validar mensagens de filtro contra um esquema

Configure uma transformação de filtro para validar mensagens recebidas contra um esquema JSON antes da execução das regras de filtro. O processo deixa mensagens que não seguem o esquema.

Para habilitar a validação de esquema, defina validateSchema para true na configuração de filtro. Quando habilitado, o filtro recupera o esquema de schemaRef na conexão de nó de entrada (o lado from da entrada nodeConnections que alimenta o nó de filtro).

A configuração de transformação de filtro inclui uma caixa de seleção Validar esquema. No entanto, a experiência de operações atualmente não dá suporte à configuração ou à exibição de schemaRef nas conexões nos nós. Para usar a validação de esquema, configure o schemaRef da conexão do nó usando manifestos Bicep ou Kubernetes.

Diretrizes:

  • Use apenas um filtro de validação por pipeline.
  • Coloque o filtro de validação primeiro para que as mensagens inválidas sejam descartadas antes de outro processamento.
  • As regras de filtro ainda se aplicam após a aprovação da validação do esquema. Se você precisar apenas de validação de esquema, deixe as regras de filtro vazias.
  • O schemaRef deve apontar para um esquema no registro de esquemas. Especifica serializationFormat o formato de esquema (por exemplo, Json).

Para saber mais sobre como configurar esquemas, consulte Noções básicas sobre esquemas de mensagens.

Enriquecer regras de filtro com dados externos

As regras de filtro dão suporte a conjuntos de dados, que permitem comparar valores com dados de um repositório de estado externo. Para obter detalhes sobre como configurar conjuntos de dados, consulte Enriquecer com dados externos.

Configuração de filtro completo

Na configuração de transformação de filtro, adicione uma ou mais regras com entradas e expressões boolianas. Opcionalmente, habilite a validação de esquema e configure conjuntos de dados para pesquisas de enriquecimento.

Chave Obrigatório Descrição
filter Sim Matriz de regras de filtro.
datasets No Matriz de definições de conjunto de dados para pesquisas de enriquecimento.
validateSchema No Quando true, valida mensagens em um esquema JSON antes da execução das regras de filtro. Usa false como padrão.

Transformação de ramo

Uma transformação de ramificação avalia uma condição em cada mensagem de entrada e a roteia para um dos dois caminhos de saída: true ou false. Ao contrário de um filtro (que descarta mensagens), um branch preserva cada mensagem e a direciona para o caminho apropriado.

Como funciona a ramificação

Cada mensagem vai para exatamente um dos dois caminhos. Nada é descartado.

Restrições de chave:

  • A expressão branch deve retornar um booliano. Resultados não booleanos causam um erro.
  • Nenhuma entrada curinga.
  • Exatamente uma regra de ramificação. A branch chave usa um único objeto, não uma matriz.

Importante

O ramificamento divide as mensagens em caminhos de processamento separados, mas todos os caminhos devem se fundir novamente usando uma transformada de concatenar antes de chegar ao destino. Pense na ramificação como uma maneira de aplicar transformações diferentes a mensagens diferentes, não como uma maneira de rotear para vários pontos de extremidade.

Definir uma regra de ramificação

Para ramificar mensagens com base em um limite de severidade:

Na configuração de transformação do branch, defina:

Configurações Valor
Input severity
Expressão $1 > 5

Mensagens onde severity é maior que 5 seguem para a rota true. Todos os outros vão para o caminho false.

Validar mensagens de desvio contra um esquema

A partir da versão 1.1.0, você pode configurar uma transformação de ramificação para validar mensagens recebidas em relação a um esquema JSON antes de avaliar a expressão de ramificação.

Para habilitar a validação do esquema, defina validateSchema como true na configuração da ramificação. O validateSchema campo é opcional e tem falsecomo padrão . Quando habilitado, a ramificação recupera o esquema de schemaRef na conexão de nó de entrada (o lado from da entrada nodeConnections que alimenta o nó de branch).

  • Mensagens que passam pela validação do esquema prosseguem para a avaliação de branch.
  • Mensagens que falham na validação do esquema vão para o false caminho.

A configuração da transformação de branch inclui uma caixa de seleção Validar esquema. No entanto, a experiência de operações atualmente não dá suporte à configuração ou à exibição de schemaRef nas conexões nos nós. Para usar a validação de esquema, configure o schemaRef da conexão do nó usando manifestos Bicep ou Kubernetes.

Conectar saídas de ramificações

Na configuração do pipeline, use o nome do nó seguido de .output.true ou .output.false para conectar cada caminho a uma transformação subsequente.

No editor de gráficos de fluxo de dados, arraste as conexões das saídas "true" e "false" da transformação de ramificação para as transformações a jusante apropriadas.

Mesclar os caminhos com a concatenação

Todos os caminhos de ramificação devem convergir antes de chegar a um destino. Uma transformação de concatenação os mescla. Ele não tem nenhuma configuração e nenhuma regra. As mensagens de todas as entradas conectadas passam sem modificações.

Adicione uma transformação de concatenação à tela e conecte os dois caminhos de ramificação a ela; em seguida, conecte a concatenação ao destino.

Exemplo: filtrar, ramificar e mesclar

Este exemplo de ponta a ponta filtra leituras incorretas, ramifica por severidade, aplica diferentes transformações de mapa a cada caminho e mescla os resultados.

Captura de tela do quadro de experiência operacional mostrando um filtro, uma ramificação, um mapa, uma concatenação e um pipeline de destino.

Para construir esse pipeline na experiência operacional:

  1. Crie um grafo de fluxo de dados e adicione uma fonte que leia de telemetry/sensors.
  2. Adicionar uma transformação de filtro. Configurar uma regra que descarta mensagens em que temperature > 1000.
  3. Adicione uma transformação branch. Configure a condição severity > 5 para rotear mensagens de alta gravidade para o caminho verdadeiro.
  4. Adicione uma transformação de mapa no caminho verdadeiro. Configure regras para renomear deviceId para id, temperature para temp, e adicionar um campo alert definido como true.
  5. Adicione uma transformação de mapa no caminho falso. Configurar regras para renomear deviceId para id .temperaturetemp
  6. Adicione uma transformação concatena para mesclar os dois caminhos.
  7. Adicionar um destino que envia para telemetry/processed.
  8. Conecte os elementos: fonte → filtro → ramificação → (caminho verdadeiro: mapa de alerta, caminho falso: mapa normal) → concatenação → destino.