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

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

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

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.

Observação

O processamento em janelas não baseado na duração requer azureiotoperations/graph-dataflow-window:1.1.0 ou posterior.

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.

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

Pré-requisitos

  • Um endpoint de registo predefinido chamado default que aponta mcr.microsoft.com é criado automaticamente durante a implementação.

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 Descrição
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.

Limitação de escalamento para grafos com estado

Importante

As transformadas de janela e de acelerador são com estado. Cada instância mantém o seu próprio estado e as instâncias não partilham esse estado entre si. Quando o número de instâncias do perfil de fluxo de dados é superior a um, as subscrições partilhadas distribuem mensagens entre instâncias, de modo que cada instância vê apenas um subconjunto das mensagens. Uma transformada de janela calcula então 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 de forma independente 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 utilize uma transformação de janela ou de aceleração. Grafos de fluxo de dados sem estado que utilizam apenas transformadas de mapeamento, filtro, ramificação e concatenação podem usar com segurança contagens de instâncias mais elevadas para aumentar o rendimento.

Quando usar uma transformação de janela

Use uma transformação de janela quando receber dados de sensores de alta frequência e quiser reduzir o volume antes de os enviar a jusante. Cenários comuns incluem:

  • Calcule médias: Um sensor de temperatura publica a cada segundo, mas a sua aplicação cloud só precisa de uma média de 30 segundos.
  • Acompanhar extremos: Quer as leituras de pressão mínima e máxima em cada intervalo de um minuto.
  • Contar eventos: Precisa de saber quantos eventos de abertura de porta ocorreram nos últimos cinco minutos.
  • Crie lotes de produção: Quer calcular estatísticas para cada lote de tamanho fixo, como a cada 100 pacotes que saem de uma linha de enchimento.
  • Responder a alterações de estado: Quer saber sempre que um sinal de funcionamento muda, como quando um misturador muda de running para draining.

Como funciona a transformação de janelas

A transformação de janela possui dois passos internos que estão ligados em sequência:

  1. Janela: Armazena mensagens até que uma das condições de fecho configuradas dispare.
  2. Acumular: Aplica as suas regras de agregação quando a janela fecha. 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 fecho: delay, count, memory, ou triggers.

Configurar condições de fecho de janelas

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

Chave de configuração Tipo de janela Purpose
delay Janela com base na duração Fecha a janela após um período fixo.
count Janela baseada em contagem Fecha a janela após um número fixo de mensagens.
memory Janela baseada em memória Feche a janela quando a dimensão 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 com base na duração

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

Observação

O passo de atraso alinha os carimbos temporais das mensagens com os limites das janelas. Se uma mensagem chegar 7 segundos após o início de uma janela de 10 segundos, pertence ao limite de 10 segundos.

Observação

Se não fornecer delay, a janela usa um timeout 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 para 30 uma janela de rotação de 30 segundos.

Propriedade Tipo Descrição
type cadeia (de caracteres) 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 da transformação da janela, defina o número de mensagens para 5 e defina o comportamento da mensagem de fronteira para messageInCurrent.

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

Janela baseada em memória

Utilize a configuração memory para fechar a janela quando o tamanho do payload em memória intermédia atingir um limite.

Na configuração de transformação da janela, defina o tamanho do buffer para 1048576 bytes e defina o comportamento da mensagem de fronteira para messageInNext.

Propriedade Tipo Descrição
type cadeia (de caracteres) Deve ser "bufferSize".
maxBufferBytes uint64 Número máximo acumulado de bytes de carga útil antes do fecho da janela. Deve ser maior que 0.
boundaryMessage cadeia (de caracteres) Se a mensagem que fecha a janela permanece na janela atual (messageInCurrent) ou inicia a janela seguinte (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 de transformação de janela, adicione uma regra de acionamento com o campo de entrada temperature, a expressão running_sum($1) + $1 > 100 e o comportamento da mensagem de limite messageInCurrent.

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

Cada regra de gatilho suporta estes campos:

Propriedade Obrigatório Descrição
inputs Sim Conjunto de referências para campos de entrada. A expressão liga a $1, $2, e assim sucessivamente.
trigger Sim Expressão booleana que fecha a janela quando é avaliada como true.
boundaryMessage Sim Se a mensagem que fecha a janela permanece na janela atual (messageInCurrent) ou inicia a janela seguinte (messageInNext).

O inputs campo suporta a mesma sintaxe de entrada usada noutros grafos de fluxo de dados, incluindo campos simples, ?? predefinidos, ? $last, $context(key).field, e $metadata.*. Para mais detalhes sobre a utilização de $context(key), consulte Enriquecer com dados externos.

As expressões de gatilho podem usar as funções regulares de expressão do grafo e as seguintes funções de estado de execução que se reiniciam quando a janela fecha:

Função Descrição
running_sum($1) Soma cumulativa de $1 entre as mensagens anteriores na janela atual.
running_avg($1) Média acumulada de $1 ao longo das mensagens anteriores.
running_min($1) Valor mínimo de $1 visto nas mensagens anteriores. Devolve $1 na primeira mensagem (o mínimo de um elemento é ele próprio).
running_max($1) Valor máximo de $1 visto em mensagens anteriores. Retorna $1 na primeira mensagem (o máximo de um elemento é ele próprio).
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. Devolve $1 na primeira mensagem.
changed($1) true se $1 difere do 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 ignoradas (o valor armazenado não é sobrescrito). Retorna $1 na primeira mensagem de uma janela.

Observação

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

Exemplos de regras de acionamento

Utilize estes exemplos para ver padrões comuns de inputs e trigger num objeto completo de configuração de acionador:

  • Este exemplo mostra uma expressão de gatilho regular. A janela fecha-se quando o valor atual temperature for superior a 80.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature"],
      "trigger": "$1 > 80",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Este exemplo mostra uma expressão de acionamento 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 limiar for ultrapassado.
{
  "type": "expression",
  "rules": [
    {
      "inputs": ["temperature"],
      "trigger": "running_sum($1) + $1 > 100",
      "boundaryMessage": "messageInCurrent"
    }
  ]
}
  • Este exemplo mostra o tratamento da entrada de forma segura em caso de valor nulo 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 acionamento com enriquecimento de dados: faz corresponder a mensagem factoryId a uma linha de armazenamento de estado, lê shiftId de $context(factory).shiftId e fecha a janela quando esse valor de deslocamento 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 definição boundaryMessage controla o que acontece à mensagem que levou ao fecho de uma janela baseada na contagem, na memória ou num acionador:

  • messageInCurrent: inclui a mensagem de fronteira na janela de encerramento.
  • messageInNext: fecha primeiro a janela atual, depois inicia a janela seguinte com a mensagem de fronteira.

Se messageInNext for acionado com a primeira mensagem numa nova janela, o fecho é suprimido para que não seja gerada uma janela vazia.

Observação

A janela baseada na duração não usa boundaryMessage. Os limites de duração são baseados em tempo, não em mensagens, por isso não há uma mensagem de fronteira para colocar na janela atual ou na próxima.

Condições de fecho combinadas

Podes combinar duração, contagem, memória e condições de disparo no mesmo gráfico.

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

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 usar messageInCurrent e count usar messageInNext, uma mensagem que satisfaz ambas as condições segue a configuração de memória. A mensagem mantém-se na janela atual e contribui para a produção de acumulação dessa janela.

Defina regras de acumulação

Cada regra de acumulação especifica como reduzir uma janela de mensagens num ú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 a ler de cada mensagem recebida.
output Sim Caminho do campo para o resultado agregado. Cada regra deve ter uma saída única.
expression Sim Fórmula que reduz os valores de entrada ao longo da janela para um único escalar. Deve conter pelo menos uma função de agregação.
description No Descrição legível para humanos.

Ao contrário das regras de mapa, expression é obrigatório para todas as regras de acumulação. Usar $1 sozinho não é válido porque faz referência a um conjunto de valores, não a um único escalar. Deve envolvê-lo numa função de agregação como average($1).

Funções de agregação

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

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

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

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

Agregações combinadas

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

Adicione uma regra com entradas temperature e humidity, e expressã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 para Fahrenheit.

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

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

Capacidade Regras do mapa Regras de acumulação
Expressão necessária No Sim
Entradas wildcard Suportado Não suportado
$metadata Acesso Suportado Não suportado
$context Enriquecimento Suportado Não suportado
? $last diretiva Suportado Não suportado
Tipo de conteúdo de saída Entrada de combates Sempre application/json

Exemplo completo de configuração

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.

A condição que fecha a janela depende da temporização da mensagem, da contagem, do tamanho da carga útil e do conteúdo. Os exemplos seguintes mostram a saída resultante para cada condição de fecho.

A duração termina

Se nenhuma outra condição for acionada primeiro e a janela completar 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
}

Encerramento da contagem

Se a janela receber estas cinco mensagens antes de qualquer outra condição ser acionada:

{ "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 encerra-se

Se o tamanho da carga útil armazenada em memória intermédia atingir 1.048.576 bytes antes de qualquer outra condição ser ativada, 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
}

O gatilho fecha

Se a expressão running_sum($1) + $1 > 100 do gatilho for ativada 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 fonte que leia de telemetry/temperature.
  2. Adicione uma transformação de janela. Configure uma janela de duração de 30 segundos, um limite de contagem 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 acionamento, defina o comportamento da mensagem de limite para messageInCurrent. Adiciona regras de acumulação para média, mínimo, máximo, contagem e alcance no temperature campo.
  3. Adicione um destino que envie para telemetry/aggregated.

Passos seguintes