Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
Una transformación de ventana agrupa los mensajes entrantes y produce un único mensaje de salida con valores agregados cuando la ventana se cierra. En lugar de reenviar cada lectura individualmente, puedes calcular estadísticas como promedios, mínimos o conteos y enviar un resultado consolidado aguas abajo.
Actualmente, una ventana puede cerrarse en función de la duración, el recuento, la memoria o las condiciones del disparador.
Para obtener información general sobre los gráficos de flujo de datos y cómo las transformaciones se componen en una canalización, consulte Introducción a los gráficos de flujo de datos.
Nota:
Las ventanas no basadas en duración requieren azureiotoperations/graph-dataflow-window:1.1.0 o versiones posteriores.
Las transformadas utilizan un lenguaje de expresiones para calcular valores, condiciones de prueba y campos de referencia. Las expresiones se refieren a las entradas por posición, no por nombre: la primera entrada en la inputs lista es $1, la segunda es $2, y así sucesivamente. Las funciones integradas, como cToF, convierten y manipulan esos valores.
Para la lista completa de operadores, funciones, tipos de datos y campos de metadatos, consulte la referencia Expressions.
Las transformadas de ventana añaden funciones de agregación como average, min, y max, que solo están disponibles en las reglas de acumulación. Para la lista completa, véase Funciones de agregación.
Prerrequisitos
- Instancia de Operaciones de IoT de Azure implementada en un clúster de Kubernetes. Para obtener más información, consulte Deploy Operaciones de IoT de Azure.
- Un punto de conexión del Registro predeterminado denominado
defaultque apunta amcr.microsoft.comse crea automáticamente durante la implementación.
Los CLI de Azure ejemplos de este artículo usan variables de entorno para que puedas establecer cada valor una vez y luego copiar y pegar los comandos as-is. Si usas el entorno Operaciones de IoT de Azure Codespaces del quickstart, estas variables ya están configuradas para ti y puedes saltarte este paso. De lo contrario, configura las siguientes variables de entorno en tu shell antes de ejecutar los comandos.
Los siguientes scripts establecen las variables de entorno más usadas:
| Variable del entorno | Descripción |
|---|---|
SUBSCRIPTION_ID |
El ID de la suscripción que contiene tu instancia de Operaciones de IoT de Azure. |
RESOURCE_GROUP |
El nombre del grupo de recursos que contiene tu instancia de Operaciones de IoT de Azure. |
AIO_INSTANCE_NAME |
El nombre de tu instancia de Operaciones de IoT de Azure. Para listar tus instancias, ejecuta az iot ops list -o table. |
CLUSTER_NAME |
El nombre del clúster Kubernetes habilitado para Azure Arc que aloja tu instancia. |
LOCATION |
La región de Azure para usar para nuevos recursos, por ejemplo eastus. |
SUBSCRIPTION_ID=<subscription-id>
RESOURCE_GROUP=<resource-group-name>
AIO_INSTANCE_NAME=<instance-name>
CLUSTER_NAME=<cluster-name>
LOCATION=<region>
Solo necesitas establecer las variables que utiliza este artículo. Este artículo podría utilizar variables de entorno adicionales para los nombres de recursos que elijas. El artículo explica cómo situarlos donde se introducen.
Limitación de escalado de grafos con estado
Importante
Las transformadas de ventana y acelerador son con estado. Cada instancia mantiene su propio estado y las instancias no comparten ese estado entre sí. Cuando el conteo de instancias del perfil de flujo de datos es mayor que uno, las suscripciones compartidas distribuyen los mensajes entre instancias, de modo que cada instancia solo ve un subconjunto de los mensajes. Una transformada de ventana calcula entonces agregaciones como promedios, sumas y conteos sobre un conjunto de datos parcial, y una transformada de aceleración impone el límite de tasa configurado de forma independiente en cada instancia en lugar de a lo largo de toda la pipeline.
Establece el recuento de instancias del perfil de flujo de datos a 1 para cualquier grafo de flujo de datos que utilice una transformación de ventana o de aceleración. Los grafos de flujo de datos sin estado que usan únicamente transformadas de mapeo, filtro, bifurcación y concatenación pueden emplear con seguridad mayores conteos de instancias para aumentar el rendimiento.
Cuándo usar una transformación de ventana
Utilice una transformación de ventana cuando reciba datos de sensor de alta frecuencia y quiera reducir el volumen antes de enviarlos a un nivel inferior. Entre los escenarios habituales se incluyen los siguientes:
- Calcular promedios: un sensor de temperatura publica cada segundo, pero tu aplicación en la nube solo necesita una media de 30 segundos.
- Rastrea los extremos: desea obtener las lecturas de presión mínima y máxima en cada intervalo de un minuto.
- Recuento de eventos: necesita saber cuántos eventos de apertura de puerta se produjeron en los últimos cinco minutos.
- Crear lotes de producción: Quieres calcular estadísticas para cada lote de tamaño fijo, como cada 100 paquetes que salen de una línea de relleno.
-
Responder a los cambios de estado: Quieres saber cuándo cambia una señal operativa, como cuando un mezclador cambia de
runningadraining.
Funcionamiento de la transformación de ventana
La transformación de ventana tiene dos pasos internos conectados en secuencia:
- Ventana: Almacena los mensajes en búfer hasta que se activa una de las condiciones de cierre configuradas.
- Acumular: Aplica las reglas de agregación cuando se cierra la ventana. La transformación reduce todos los mensajes de la ventana a un solo mensaje de salida.
Nota:
Una transformada de ventana debe configurar al menos una condición de cierre: delay, count, memory, o triggers.
Configurar las condiciones de cierre de ventanas
A partir de la versión 1.1.0, la transformación de ventana añade tres claves de configuración junto con la clave existente delay :
| Clave de configuración | Tipo de ventana | propósito |
|---|---|---|
delay |
Ventana basada en la duración | Cierra la ventana tras un periodo fijo. |
count |
Ventana basada en conteo | Cierra la ventana tras un número fijo de mensajes. |
memory |
Ventana basada en memoria | Cierra la ventana cuando el tamaño de la carga útil almacenada en búfer alcance un límite. |
triggers |
Ventana basada en activadores | Cierra la ventana cuando una expresión personalizada dé como resultado true. |
Ventana basada en la duración
Usa la delay configuración para cerrar la ventana tras un periodo fijo. Esta configuración controla cuánto dura cada ventana de caída.
Nota:
El paso de retraso alinea las marcas de tiempo del mensaje con los límites de la ventana. Si un mensaje llega 7 segundos después de iniciar una ventana de 10 segundos, pertenece al límite de 10 segundos.
Nota:
Si no lo proporcionas delay, la ventana usa un tiempo de espera por defecto de 60 segundos como válvula de escape.
En la configuración de transformación de ventana, establezca la duración de la ventana en segundos. Por ejemplo, configúrelo en 30 para una ventana deslizante de 30 segundos.
| Propiedad | Tipo | Descripción |
|---|---|---|
type |
cuerda / cadena | Debe ser "duration". |
delaySeconds |
uint64 | Número de segundos antes de que se cierre la ventana. Debe ser mayor que 0. |
Ventana basada en conteo
Utiliza la count configuración para cerrar la ventana tras un número fijo de mensajes.
En la configuración de la transformación de ventana, establece el conteo de mensajes en 5 y el comportamiento de los mensajes de frontera en mensajeInCurrent.
| Propiedad | Tipo | Descripción |
|---|---|---|
type |
cuerda / cadena | Debe ser "messageCount". |
maxMessageCount |
uint64 | Número de mensajes para almacenar antes de que se cierre la ventana. Debe ser mayor que 0. |
boundaryMessage |
cuerda / cadena | Si el mensaje que cierra la ventana permanece en la ventana actual (messageInCurrent) o inicia la siguiente ventana (messageInNext). |
Ventana basada en memoria
Utiliza la memory configuración para cerrar la ventana cuando el tamaño de la carga útil almacenada en búfer alcanza un límite.
En la configuración de transformación de ventanas, establece Tamaño del búfer en 1048576 bytes y establece el comportamiento del mensaje delimitador en messageInNext.
| Propiedad | Tipo | Descripción |
|---|---|---|
type |
cuerda / cadena | Debe ser "bufferSize". |
maxBufferBytes |
uint64 | Bytes acumulados máximos de carga útil antes de que se cierre la ventana. Debe ser mayor que 0. |
boundaryMessage |
cuerda / cadena | Si el mensaje que cierra la ventana permanece en la ventana actual (messageInCurrent) o inicia la siguiente ventana (messageInNext). |
Ventana basada en activadores
Utiliza la triggers configuración cuando la ventana debería cerrarse según el contenido del mensaje o el estado en ejecución dentro de la ventana actual.
En la configuración de la transformación de ventanas, añade una regla de activación con el campo de entrada temperature, la expresión running_sum($1) + $1 > 100 y el comportamiento del mensaje de límite messageInCurrent.
| Propiedad | Obligatorio | Descripción |
|---|---|---|
type |
Sí | Debe ser "expression". |
rules |
Sí | Conjunto de reglas de activación. Las reglas se evalúan secuencialmente por mensaje; La primera regla de coincidencia cierra la ventana. |
datasets |
No | Conjuntos de datos opcionales en el almacén de estados que hacen referencia al almacén de estados. |
Cada regla de desencadenante soporta estos campos:
| Propiedad | Obligatorio | Descripción |
|---|---|---|
inputs |
Sí | Array de referencias de campos de entrada. La expresión vincula a $1, $2, y así sucesivamente. |
trigger |
Sí | Expresión booleana que cierra la ventana cuando evalúa a true. |
boundaryMessage |
Sí | Si el mensaje que cierra la ventana permanece en la ventana actual (messageInCurrent) o inicia la siguiente ventana (messageInNext). |
El inputs campo soporta la misma sintaxis de entrada usada en otros grafos de flujo de datos, incluyendo campos simples, ?? valores predeterminados, ? $last, $context(key).field, y $metadata.*. Para obtener más información sobre cómo usar $context(key), consulte Enriquecimiento con datos externos.
Las expresiones desencadenadoras pueden usar las funciones habituales de expresiones de grafo y las siguientes funciones de estado en ejecución, que se restablecen cuando se cierra la ventana:
| Function | Descripción |
|---|---|
running_sum($1) |
Suma acumulada de $1 entre mensajes anteriores en la ventana actual. |
running_avg($1) |
Promedio acumulado de $1 a lo largo de los mensajes anteriores. |
running_min($1) |
Valor mínimo de $1 visto en mensajes anteriores. Devuelve $1 en el primer mensaje (el mínimo de un único elemento es el propio elemento). |
running_max($1) |
Valor máximo de $1 visto en mensajes anteriores. Devuelve $1 en el primer mensaje (el máximo de un solo elemento es el propio elemento). |
running_count($1) |
Recuento de mensajes en los que $1 estaba presente. |
running_count() |
Recuento total de mensajes (sin filtro de campo). |
first($1) |
Primer valor no vacío de $1 en la ventana actual. Devuelve $1 en el primer mensaje. |
changed($1) |
true si $1 difiere de su valor en el mensaje anterior.
false en el primer mensaje de una ventana (sin valor previo con el que comparar). |
prev($1) |
Valor no vacío más reciente de $1 procedente de un mensaje anterior en la ventana actual. Los mensajes donde $1 estaba Vacío se omiten (el valor almacenado no se sobrescribe). Regresa $1 en el primer mensaje de una ventana. |
Nota:
running_sum($1) y funciones similares devolven valores de mensajes previamente procesados. Para el mensaje actual, usa $1.
Ejemplos de reglas de activación
Utiliza estos ejemplos para ver patrones comunes inputs y trigger en un objeto completo de configuración de disparador:
- Este ejemplo muestra una expresión de disparador regular. La ventana se cierra cuando la corriente
temperaturesupera los 80.
{
"type": "expression",
"rules": [
{
"inputs": ["temperature"],
"trigger": "$1 > 80",
"boundaryMessage": "messageInCurrent"
}
]
}
- Este ejemplo muestra una expresión desencadenadora que usa
running_sum($1) + $1para combinar los mensajes anteriores de la ventana actual con el mensaje actual, y luego cierra la ventana cuando se supera el umbral.
{
"type": "expression",
"rules": [
{
"inputs": ["temperature"],
"trigger": "running_sum($1) + $1 > 100",
"boundaryMessage": "messageInCurrent"
}
]
}
- Este ejemplo muestra la gestión segura de entradas nulas con
temperature ?? 0, además demessageInNextpara colocar el mensaje de límite en la ventana siguiente.
{
"type": "expression",
"rules": [
{
"inputs": ["temperature ?? 0"],
"trigger": "running_avg($1) > 80",
"boundaryMessage": "messageInNext"
}
]
}
- Este ejemplo muestra un disparador basado en metadatos donde la ventana se cierra para un valor específico de tema de
$metadata.topic.
{
"type": "expression",
"rules": [
{
"inputs": ["$metadata.topic"],
"trigger": "$1 == \"telemetry/high-priority\"",
"boundaryMessage": "messageInCurrent"
}
]
}
- Este ejemplo muestra reglas de disparo usando enriquecimiento de conjuntos de datos: empareja el mensaje
factoryIdcon una fila de almacenamiento de estados, leeshiftIddesde$context(factory).shiftId, y cierra la ventana cuando ese valor de desplazamiento cambia (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"
}
]
}
En este ejemplo, se espera que el conjunto de datos de la tienda de estados representado por factory contenga campos como factoryId y shiftId.
Comportamiento de frontera
La opción boundaryMessage controla qué ocurre con el mensaje que hizo que se cerrara una ventana por recuento, memoria o activador:
-
messageInCurrent: incluye el mensaje de límite en la ventana de cierre. -
messageInNext: cierra primero la ventana actual, luego inicia la siguiente ventana con el mensaje de límite.
Si messageInNext se activa en el primer mensaje de una nueva ventana, el cierre se suprime para que no se emita una ventana vacía.
Nota:
La ventana basada en duración no usa boundaryMessage. Los límites de duración se basan en tiempos, no en mensajes, por lo que no hay un mensaje límite que colocar en la ventana actual o en la siguiente.
Condiciones de cierre combinadas
Puedes combinar duración, conteo, memoria y condiciones de disparo en el mismo gráfico.
- La duración se determina por el tiempo y se evalúa mediante el temporizador.
- Para cada mensaje entrante, las condiciones basadas en mensajes se evalúan en este orden:
Memory > Count > Trigger. - Dentro de
triggers.rules, las reglas se evalúan secuencialmente y gana la primera regla que coincida.
El orden de evaluación en condiciones basadas en mensajes es importante para los resultados de acumulación cuando un mensaje cumple varias condiciones al mismo tiempo. Por ejemplo, si memory usa messageInCurrent y count usa messageInNext, un mensaje que cumple ambas condiciones sigue la configuración de memoria. El mensaje se mantiene en la ventana actual y contribuye a la salida acumulada de esa ventana.
Definición de reglas de acumulación
Cada regla de acumulación especifica cómo reducir una ventana de mensajes en un único valor de salida. La clave de configuración es rules.
En la configuración de transformación de ventana, se añade una regla de acumulación con entrada temperature, salida avgTemperature, y función average($1)de agregación .
| Propiedad | Obligatorio | Descripción |
|---|---|---|
inputs |
Sí | Lista de rutas de campo para leer de cada mensaje entrante. |
output |
Sí | Ruta del campo para el resultado agregado. Cada regla debe tener una salida única. |
expression |
Sí | Fórmula que reduce los valores de entrada en la ventana a un solo escalar. Debe contener al menos una función de agregación. |
description |
No | Descripción legible para humanos. |
A diferencia de las reglas de mapeo, expression es necesario para cada regla de acumulación. El uso $1 solo no es válido porque hace referencia a una colección de valores, no a una sola escalar. Debe encapsularlo en una función de agregación como average($1).
Funciones de agregación
| Function | Devoluciones | Comportamiento de ventana vacía |
|---|---|---|
average |
Media de valores numéricos | Error |
sum |
Suma de valores numéricos | 0,0 |
min |
Valor numérico mínimo | Error |
max |
Valor numérico máximo | Error |
count |
Recuento de mensajes en los que existe el campo | 0 |
first |
Primer valor de la ventana | Error |
last |
Último valor de la ventana | Error |
Cada función toma una sola variable posicional como argumento ($1 para la primera entrada, $2 para el segundo, etc.).
Valores no numéricos: las averagefunciones , sum, miny max omiten silenciosamente los valores no numéricos.
Funciones basadas en presencia: count, firsty last operan en presencia de campo independientemente del tipo de valor.
Combinar agregaciones
Combina múltiples funciones de agregación en una sola expresión:
Agregue una regla con entradas temperature y humidity, y expresión average($1) + max($2).
Para convertir un valor agregado, aplique la función de conversión fuera de la agregación. Por ejemplo, cToF(average($1)) convierte la temperatura media en Fahrenheit.
Cada función de agregación debe hacer referencia directamente a una única variable posicional.
average($1) + max($2) es válido, pero average($1 + $2) no lo es.
Diferencias con las reglas del mapa
| Capacidad | Reglas del mapa | Reglas de acumulación |
|---|---|---|
| Expresión necesaria | No | Sí |
| Entradas con caracteres comodín | Soportado | No soportado |
$metadata acceso |
Soportado | No soportado |
$context enriquecimiento |
Soportado | No soportado |
directiva ? $last |
Soportado | No soportado |
| Tipo de contenido de salida | Coincide con la entrada | Siempre application/json |
Ejemplo de configuración completa
Este ejemplo muestra una configuración completa de ventana que cierra la ventana tras 30 segundos, 5 mensajes, 1.048.576 bytes en búfer, o cuando running_sum($1) + $1 > 100. El ejemplo establece el boundaryMessage valor en messageInCurrent para las tres últimas condiciones, y la ventana calcula las estadísticas de temperatura cuando se cierra.
Qué condición cierra la ventana depende del momento del mensaje, el recuento, el tamaño de la carga útil y el contenido. Los siguientes ejemplos muestran la salida resultante para cada condición de cierre.
Fin de la duración
Si no se cumple ninguna otra condición antes y la ventana llega a los 30 segundos después de recibir estos tres mensajes:
{ "temperature": 21.5 }
{ "temperature": 23.0 }
{ "temperature": 19.8 }
El mensaje de salida es:
{
"avgTemperature": 21.433333333333334,
"minTemperature": 19.8,
"maxTemperature": 23.0,
"readingCount": 3,
"tempRange": 3.2
}
Finaliza el recuento
Si la ventana recibe estos cinco mensajes antes de que se active cualquier otra condición:
{ "temperature": 20.0 }
{ "temperature": 22.0 }
{ "temperature": 21.0 }
{ "temperature": 24.0 }
{ "temperature": 23.0 }
El mensaje de salida es:
{
"avgTemperature": 22.0,
"minTemperature": 20.0,
"maxTemperature": 24.0,
"readingCount": 5,
"tempRange": 4.0
}
La memoria se cierra
Si el tamaño de la carga útil almacenada en búfer alcanza 1.048.576 bytes antes de que se active cualquier otra condición, por ejemplo después de estos dos mensajes grandes:
{ "temperature": 21.0, "payloadPad": "<large string>" }
{ "temperature": 22.5, "payloadPad": "<large string>" }
El mensaje de salida es:
{
"avgTemperature": 21.75,
"minTemperature": 21.0,
"maxTemperature": 22.5,
"readingCount": 2,
"tempRange": 1.5
}
El gatillo se cierra
Si la expresión running_sum($1) + $1 > 100 del disparador se activa antes de cualquier otra condición, por ejemplo después de estos tres mensajes:
{ "temperature": 40.0 }
{ "temperature": 35.0 }
{ "temperature": 30.0 }
El mensaje de salida es:
{
"avgTemperature": 35.0,
"minTemperature": 30.0,
"maxTemperature": 40.0,
"readingCount": 3,
"tempRange": 10.0
}
En la experiencia de operaciones, cree un diagrama de flujo de datos con una transformación de ventana:
- Agregue un origen que lea de
telemetry/temperature. - Agregue una transformación de ventana. Configura una ventana de 30 segundos, un límite de recuento de 5 mensajes, un límite de tamaño de búfer de 1.048.576 bytes y una regla de desencadenamiento en
temperaturecon la expresiónrunning_sum($1) + $1 > 100. Para las condiciones de recuento, memoria y activación, establece el comportamiento del mensaje de borde enmessageInCurrent. Añade reglas de acumulación para el promedio, mínimo, máximo, conteo y rango en eltemperaturecampo. - Agregue un destino que envíe a
telemetry/aggregated.