Agregue dados com transformações de janela em grafos de fluxo de dados

Uma transformação de janela agrupa mensagens recebidas e produz uma única mensagem de saída com valores agregados quando a janela se fecha. Em vez de encaminhar cada leitura individualmente, você pode calcular estatísticas como médias, mínimos ou contagens e enviar um resultado consolidado adiante.

Atualmente, uma janela pode fechar com base na duração, contagem, memória ou condições de gatilho.

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.

Observação

Janelas que não são baseadas em duração exigem azureiotoperations/graph-dataflow-window:1.1.0 ou posteriores.

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.

Transformadas de janela adicionam funções de agregação como average, min, e max, que estão disponíveis apenas em regras de acumulação. Para a lista completa, veja Funções de agregação.

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.

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.

Limitação de escalabilidade para grafos com estado

Importante

As transformadas de janela e de acelerador são com estado. Cada instância mantém seu próprio estado e as instâncias não compartilham esse estado entre si. Quando a contagem de instâncias do perfil de fluxo de dados é maior que uma, assinaturas compartilhadas distribuem mensagens entre instâncias, de modo que cada instância vê apenas um subconjunto das mensagens. Uma transformada de janela então calcula agregações como médias, somas e contagens sobre um conjunto de dados parcial, e uma transformada de aceleração impõe o limite de taxa configurado independentemente em cada instância, em vez de em todo o pipeline.

Defina a contagem de instâncias do perfil de fluxo de dados para 1 para qualquer grafo de fluxo de dados que use uma transformação de janela ou de aceleração. Grafos de fluxo de dados sem estado que usam apenas transformadas de mapeamento, filtro, desvio e concatenação podem usar com segurança contagens maiores de instâncias para aumentar a taxa de transferência.

Quando usar uma transformação de janela

Utilize uma transformação de janela quando receber dados de sensores de alta frequência e desejar reduzir o volume antes de enviá-los para a etapa seguinte. Cenários comuns incluem:

  • Calcular médias: um sensor de temperatura publica dados todo segundo, mas seu aplicativo na nuvem só precisa de uma média de 30 segundos.
  • Acompanhe os extremos: você deseja as leituras de pressão mínima e máxima em cada intervalo de um minuto.
  • Contar eventos: é preciso saber quantos eventos de abertura de porta ocorreram nos últimos cinco minutos.
  • Crie lotes de produção: Você quer calcular estatísticas para cada lote de tamanho fixo, como a cada 100 pacotes saindo de uma linha de enchimento.
  • Responda a mudanças de estado: Você quer saber quando um sinal operacional muda, como quando um mixer muda de running para draining.

Como funciona a transformação da janela

A transformação de janela possui duas etapas internas conectadas em sequência.

  1. Janela: armazena mensagens até que uma das condições de fechamento configuradas seja disparada.
  2. Acumular: aplica suas regras de agregação quando a janela é fechada. A transformação reduz todas as mensagens na janela para uma única mensagem de saída.

Observação

Uma transformada de janela deve configurar pelo menos uma condição de fechamento: delay, count, memory, ou triggers.

Configurar condições de fechamento de janelas

A partir da versão 1.1.0, a transformação de janela adiciona três chaves de configuração junto com a chave existente delay :

Chave de configuração Tipo de janela Purpose
delay Janela baseada em duração Feche a janela após um período fixo.
count Janela baseada em contagem Feche a janela após um número fixo de mensagens.
memory Janela baseada em memória Feche a janela quando o tamanho da carga útil armazenada em buffer atingir um limite.
triggers Janela baseada em gatilho Feche a janela quando uma expressão personalizada for avaliada como true.

Janela baseada em duração

Use a delay configuração para fechar a janela após um período fixo. Essa configuração controla quanto tempo cada janela de rotação dura.

Observação

A etapa de atraso alinha os carimbos de data/hora das mensagens aos limites da janela. Se uma mensagem chegar 7 segundos após o início de uma janela de 10 segundos, ela pertence ao limite de 10 segundos.

Observação

Se você não fornecer delay, a janela usa um tempo padrão de 60 segundos como válvula de segurança.

Na configuração de transformação da janela, defina a duração da janela em segundos. Por exemplo, defina-o 30 para uma janela em cascata de 30 segundos.

Propriedade Tipo Descrição
type cadeia Deve ser "duration".
delaySeconds uint64 Número de segundos antes da janela fechar. Deve ser maior que 0.

Janela baseada em contagem

Use a count configuração para fechar a janela após um número fixo de mensagens.

Na configuração de transformação de janela, defina a contagem de mensagens para 5 e defina o comportamento da mensagem de fronteira para messageInCurrent.

Propriedade Tipo Descrição
type cadeia Deve ser "messageCount".
maxMessageCount uint64 Número de mensagens para armazenar antes do fechamento da janela. Deve ser maior que 0.
boundaryMessage cadeia Se a mensagem que fecha a janela permanece na janela atual (messageInCurrent) ou inicia a próxima janela (messageInNext).

Janela baseada em memória

Use a configuração memory para fechar a janela quando o tamanho da carga útil armazenada em buffer atingir um limite.

Na configuração de transformação de janela, defina o Tamanho do buffer como 1048576 bytes e defina o comportamento da mensagem de limite como messageInNext.

Propriedade Tipo Descrição
type cadeia Deve ser "bufferSize".
maxBufferBytes uint64 Bytes máximos acumulados de payload antes do fechamento da janela. Deve ser maior que 0.
boundaryMessage cadeia Se a mensagem que fecha a janela permanece na janela atual (messageInCurrent) ou inicia a próxima janela (messageInNext).

Janela baseada em gatilho

Use a triggers configuração quando a janela deve fechar com base no conteúdo da mensagem ou no estado em execução dentro da janela atual.

Na configuração da transformação de janela, adicione uma regra de gatilho com o campo de entrada temperature, a expressão running_sum($1) + $1 > 100 e o comportamento da mensagem de fronteira messageInCurrent.

Propriedade Obrigatório Descrição
type Sim Deve ser "expression".
rules Sim Conjunto de regras de disparo. As regras são avaliadas sequencialmente por mensagem; A primeira regra de correspondência fecha a janela.
datasets No Conjuntos de dados opcionais do armazenamento de estado que fazem referência ao armazenamento de estado.

Cada regra de gatilho suporta estes campos:

Propriedade Obrigatório Descrição
inputs Sim Array de referências a campos de entrada. A expressão liga a $1, $2, e assim por diante.
trigger Sim Expressão booleana que fecha a janela quando resulta em true.
boundaryMessage Sim Se a mensagem que fecha a janela permanece na janela atual (messageInCurrent) ou inicia a próxima janela (messageInNext).

O inputs campo suporta a mesma sintaxe de entrada usada em outros gráficos de fluxo de dados, incluindo campos simples, ?? padrão, ? $last, $context(key).field, e $metadata.*. Para mais detalhes sobre como usar $context(key), veja Enriquecer com dados externos.

Expressões de acionamento podem usar as funções regulares de expressão de grafo e as seguintes funções de estado de execução, que são redefinidas quando a janela é fechada:

Função Descrição
running_sum($1) Soma cumulativa de $1 entre mensagens anteriores na janela atual.
running_avg($1) Média cumulativa de $1 ao longo das mensagens anteriores.
running_min($1) Valor mínimo de $1 visto em mensagens anteriores. Retorna $1 na mensagem inicial (o mínimo de um elemento é ele mesmo).
running_max($1) Valor máximo de $1 visto em mensagens anteriores. Retorna $1 na primeira mensagem (máximo de um elemento é ele mesmo).
running_count($1) Contagem de mensagens em que $1 estava presente.
running_count() Contagem total de mensagens (sem filtro de campo).
first($1) Primeiro valor não vazio de $1 na janela atual. Retorna $1 na primeira mensagem.
changed($1) true se $1 difere de seu valor na mensagem anterior. false na primeira mensagem de uma janela (sem valor anterior para comparar).
prev($1) Valor não vazio mais recente de $1 de uma mensagem anterior na janela atual. Mensagens onde $1 estava Vazio são puladas (o valor armazenado não é sobrescrevido). Retorna $1 na primeira mensagem de uma janela.

Observação

running_sum($1) e funções similares retornam valores de mensagens previamente processadas. Para a mensagem atual, use $1.

Exemplos de regras de acionamento

Use estes exemplos para ver padrões comuns inputs e trigger em um objeto completo de configuração de gatilho:

  • Este exemplo mostra uma expressão de gatilho regular. A janela fecha quando o temperature atual é maior que 80.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature"],
      "trigger": "$1 > 80",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Este exemplo mostra uma expressão de gatilho que usa running_sum($1) + $1 para combinar mensagens anteriores na janela atual com a mensagem atual e, em seguida, fechar a janela quando o limite for excedido.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature"],
      "trigger": "running_sum($1) + $1 > 100",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Este exemplo mostra o tratamento seguro de entradas nulas com temperature ?? 0, além de messageInNext para colocar a mensagem de limite na janela seguinte.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature ?? 0"],
      "trigger": "running_avg($1) > 80",
      "boundaryMessage": "messageInNext"
    }
  ]
}
  • Este exemplo mostra um gatilho baseado em metadados em que a janela se fecha para um valor específico do tópico de $metadata.topic.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["$metadata.topic"],
      "trigger": "$1 == \"telemetry/high-priority\"",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Este exemplo mostra regras de gatilho usando enriquecimento de conjunto de dados: ele corresponde a mensagem factoryId a uma linha de armazenamento de estados, lê shiftId de $context(factory).shiftId, e fecha a janela quando esse valor de mudança muda (changed($1)).
{
  "type": "expression",
  "datasets": [
    {
      "key": "factory",
      "inputs": ["$source.factoryId", "$context.factoryId"],
      "expression": "$1 == $2"
    }
  ],
  "rules": [
    {
      "inputs": ["$context(factory).shiftId"],
      "trigger": "changed($1)",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}

Neste exemplo, espera-se que o conjunto de dados de armazenamento de estados representado por factory conter campos como factoryId e shiftId.

Comportamento de fronteira

A configuração boundaryMessage controla o que acontece com a mensagem que fez com que uma janela baseada em contagem, memória ou gatilho fosse fechada:

  • messageInCurrent: inclui a mensagem de fronteira na janela de fechamento.
  • messageInNext: fecha a janela atual primeiro, depois inicia a próxima janela com a mensagem de fronteira.

Se messageInNext for disparado na primeira mensagem de uma nova janela, o fechamento é suprimido para que uma janela vazia não seja emitida.

Observação

A janela baseada em duração não usa boundaryMessage. Limites de duração são baseados em tempo, não em mensagens, então não há uma mensagem de limite para colocar na janela atual ou próxima.

Combinar condições de fechamento

Você pode combinar duração, contagem, memória e condições de gatilho no mesmo gráfico.

  • A duração é guiada pelo tempo e avaliada pelo temporizador.
  • Para cada mensagem recebida, condições orientadas por mensagens são avaliadas nesta ordem: Memory > Count > Trigger.
  • Em triggers.rules, as regras são avaliadas sequencialmente e a primeira regra correspondente vence.

A ordem de avaliação em condições orientadas por mensagens é importante para os resultados de acumulação quando uma mensagem satisfaz múltiplas condições ao mesmo tempo. Por exemplo, se memory usa messageInCurrent e count usa messageInNext, uma mensagem que satisfaz ambas as condições segue a configuração de memória. A mensagem permanece na janela atual e contribui para a saída acumulada dessa janela.

Definir regras de acúmulo

Cada regra de acúmulo especifica como reduzir uma janela de mensagens em um único valor de saída. A chave de configuração é rules.

Na configuração da transformação de janela, adicione uma regra de acumulação com entrada temperature, saída avgTemperature e função de agregação average($1).

Propriedade Obrigatório Descrição
inputs Sim Lista de caminhos de campo para leitura de cada mensagem de entrada.
output Sim Campo de caminho para o resultado agregado. Cada regra deve ter uma saída exclusiva.
expression Sim Fórmula que reduz os valores de entrada em toda a janela para um único escalar. Deve conter pelo menos uma função de agregação.
description No Descrição legível por humanos.

Ao contrário das regras de mapa, expression é necessário para cada regra de acúmulo. Usar $1 sozinho não é válido porque faz referência a uma coleção de valores, não a um único escalar. Você deve encapsulá-lo em uma função de agregação como average($1).

Funções de agregação

Função Devoluções Comportamento de janela vazia
average Média de valores numéricos Erro
sum Soma de valores numéricos 0,0
min Valor numérico mínimo Erro
max Valor numérico máximo Erro
count Contagem de mensagens em que o campo existe 0
first Primeiro valor na janela Erro
last Último valor na janela Erro

Cada função usa uma única variável posicional como seu argumento ($1 para a primeira entrada, $2 para a segunda e assim por diante).

Valores não numéricos: As funções average, sum, min e max ignoram silenciosamente valores não numéricos.

Funções baseadas em presença: count, first e last operam na presença de campo, independentemente do tipo de valor.

Combinar agregações

Combine múltiplas funções de agregação em uma única expressão:

Adicione uma regra com entradas temperature e humidityexpressão average($1) + max($2).

Para converter um valor agregado, aplique a função de conversão fora da agregação. Por exemplo, cToF(average($1)) converte a temperatura média em Fahrenheit.

Cada função de agregação deve referenciar uma única variável posicional diretamente. average($1) + max($2) é válido, mas average($1 + $2) não é.

Diferenças em relação às regras de mapa

Capacidade Regras de mapa Regras de acúmulo
Expressão necessária No Sim
Entradas curinga Supported Sem suporte
$metadata Acesso Supported Sem suporte
$context Enriquecimento Supported Sem suporte
Diretiva ? $last Supported Sem suporte
Tipo de conteúdo de saída Corresponde à entrada Sempre application/json

Exemplo de configuração completa

Este exemplo mostra uma configuração completa de janela que fecha a janela após 30 segundos, 5 mensagens, 1.048.576 bytes em buffer, ou quando running_sum($1) + $1 > 100. O exemplo define o boundaryMessage valor para messageInCurrent as últimas três condições, e a janela calcula estatísticas de temperatura quando fecha.

Qual condição fecha a janela depende do tempo da mensagem, contagem, tamanho da carga útil e conteúdo. Os exemplos a seguir mostram a saída resultante para cada condição de fechamento.

Fim da duração

Se nenhuma outra condição disparar primeiro e a janela chegar a 30 segundos após receber estas três mensagens:

{ "temperature": 21.5 }
{ "temperature": 23.0 }
{ "temperature": 19.8 }

A mensagem de saída é:

{
  "avgTemperature": 21.433333333333334,
  "minTemperature": 19.8,
  "maxTemperature": 23.0,
  "readingCount": 3,
  "tempRange": 3.2
}

A contagem fecha

Se a janela receber estas cinco mensagens antes que qualquer outra condição seja disparada:

{ "temperature": 20.0 }
{ "temperature": 22.0 }
{ "temperature": 21.0 }
{ "temperature": 24.0 }
{ "temperature": 23.0 }

A mensagem de saída é:

{
  "avgTemperature": 22.0,
  "minTemperature": 20.0,
  "maxTemperature": 24.0,
  "readingCount": 5,
  "tempRange": 4.0
}

A memória se fecha

Se o tamanho da carga útil armazenada em buffer chegar a 1.048.576 bytes antes que qualquer outra condição seja disparada, por exemplo, após estas duas mensagens grandes:

{ "temperature": 21.0, "payloadPad": "<large string>" }
{ "temperature": 22.5, "payloadPad": "<large string>" }

A mensagem de saída é:

{
  "avgTemperature": 21.75,
  "minTemperature": 21.0,
  "maxTemperature": 22.5,
  "readingCount": 2,
  "tempRange": 1.5
}

Gatilho fecha

Se a expressão running_sum($1) + $1 > 100 do gatilho for acionada antes de qualquer outra condição, por exemplo, após estas três mensagens:

{ "temperature": 40.0 }
{ "temperature": 35.0 }
{ "temperature": 30.0 }

A mensagem de saída é:

{
  "avgTemperature": 35.0,
  "minTemperature": 30.0,
  "maxTemperature": 40.0,
  "readingCount": 3,
  "tempRange": 10.0
}

Na experiência de Operações, crie um gráfico de fluxo de dados com uma transformação de janela:

  1. Adicione uma origem que leia a partir de telemetry/temperature.
  2. Adicione uma transformação de janela. Configure uma janela com duração de 30 segundos, um limite de 5 mensagens, um limite de tamanho do buffer de 1,048,576 bytes e uma regra de acionamento em temperature com a expressão running_sum($1) + $1 > 100. Para as condições de contagem, memória e gatilho, defina o comportamento da mensagem de fronteira para messageInCurrent. Adicione regras de acumulação para média, mínimo, máximo, contagem e alcance no temperature campo.
  3. Adicionar um destino que envia para telemetry/aggregated.

Próximas Etapas