Data flow graphs in Operaciones de IoT de Azure

Un grafo de flujo de datos es una tubería configurable que procesa los datos a medida que se mueven por Operaciones de IoT de Azure. Un flujo de datos estándar sigue una secuencia fija de enriquecimiento, filtro, mapeación, pero un gráfico de flujo de datos permite componer transformaciones en cualquier orden, ramificar en caminos paralelos y agregar datos a lo largo de ventanas temporales.

El DataflowGraph recurso personalizado Kubernetes define un grafo de flujo de datos. Dentro del recurso, conectas fuentes, transformaciones y destinos para construir pipelines de procesamiento que se ajusten a tu escenario.

Importante

Actualmente, los gráficos de flujo de datos solo admiten puntos de conexión MQTT, Kafka y OpenTelemetry. No se admiten otros tipos de punto de conexión, como Data Lake, Microsoft Fabric OneLake, Azure Data Explorer y Almacenamiento local.

Flujos de datos frente a gráficos de flujo de datos

Operaciones de IoT de Azure proporciona dos maneras de procesar datos en una canalización:

Capacidad Flujos de datos Gráficos de flujo de datos
Forma de canalización Corregido: enriquecer, filtrar, mapear Flexible: cualquier orden, bifurcación, combinación
Tipos de transformación Mapear, filtrar, enriquecer Asignar, filtrar, ramificar, concatenar, ventanilla, limitar, enriquecer
Agregación basada en tiempo No disponible Transformaciones de ventana con ventanas de saltos de tamaño constante
Enrutamiento condicional No disponible Transformaciones de bifurcación y concatenación
Soporte de puntos de conexión Todos los tipos de punto de conexión MQTT, Kafka y OpenTelemetry solo

Para los nuevos proyectos que usan tipos de punto de conexión admitidos, se recomiendan gráficos de flujo de datos. Los flujos de datos siguen siendo totalmente compatibles con todos los escenarios y admiten la gama completa de tipos de punto de conexión.

Transformaciones disponibles

Cada transformación es un paso de procesamiento preconstruido que configuras con reglas y encadenas con otras transformaciones dentro de un DataflowGraph recurso.

Transformación Artefacto Descripción
Mapa azureiotoperations/graph-dataflow-map:1.0.0 Cambiar el nombre, reestructurar, calcular y copiar campos.
Filter azureiotoperations/graph-dataflow-filter:1.0.0 Descartar los mensajes que cumplan una condición.
Rama azureiotoperations/graph-dataflow-branch:1.0.0 Dirija cada mensaje a una ruta true o false en función de una condición.
Concatenate azureiotoperations/graph-dataflow-concatenate:1.0.0 Combinar dos o más rutas en una sola.
Ventana azureiotoperations/graph-dataflow-window:1.0.0 Recopile mensajes a lo largo de un intervalo de tiempo y, a continuación, agregue.
Throttle azureiotoperations/graph-dataflow-throttle:1.0.0 Limita la tasa de mensajes por patrón de tema MQTT.

Todas las transformaciones comparten un lenguaje de expresión para operadores, funciones y referencias de campo. También puede enriquecer los mensajes con datos externos de un almacén de estado en transformaciones de mapa, filtro y rama.

Tip

Las expresiones usan variables posicionales, por lo que $1 la primera entrada es la segunda, $2 y así sucesivamente. La referencia Expressions lista funciones integradas como cToF y cubre todos los operadores, funciones y campos de metadatos disponibles para transformaciones.

Cómo se componen las transformaciones en un grafo de flujo de datos

Las transformaciones se conectan en secuencia dentro de un recurso DataflowGraph: Origen > Transformación A > Transformación B > ... > Destino.

Las transformaciones de rama dividen el flujo en rutas paralelas, y las transformaciones de concatenación las vuelven a unir.

Puede encadenar una cantidad cualquiera de transformaciones en el orden que desee. Una canalización con una sola transformación de mapeo es tan válida como una que filtra, bifurca, mapea cada ruta de forma diferente, combina y, a continuación, agrega a través de una ventana temporal.

Cómo funciona la configuración de grafos de flujo de datos

Cada transformación en un grafo de flujo de datos hace referencia a un artefacto preconstruido extraído de un registro de contenedores. Configure la transformación pasando reglas como JSON a través de la sección configuration del recurso gráfico.

Cuando despliegas Operaciones de IoT de Azure, se crea automáticamente un endpoint de registro predeterminado llamado default que apunta a mcr.microsoft.com. Las transformaciones integradas utilizan este endpoint para extraer artefactos del Microsoft Container Registry. No necesitas ninguna configuración extra de registro.

Un recurso de grafo de flujo de datos define tres tipos de elementos: una fuente, una o más transformaciones (cada una con nodeType: Graph), y un destino, y un conjunto de nodeConnections esos que describen cómo fluyen los datos entre ellos. Cada transformación configuration pasa sus reglas como una cadena JSON bajo la rules tecla.

Para un ejemplo completo y ejecutable que lee datos de temperatura, convierte Celsius a Fahrenheit con una transformación de mapa y publica el resultado —en la experiencia de Operaciones, CLI de Azure, Bicep y Kubernetes— véase Crear un gráfico de flujo de datos. En los artículos de procedimientos siguientes, los ejemplos se centran en las propias reglas de transformación.

Configurar esquemas en conexiones de nodos

Los grafos de flujo de datos gestionan los esquemas de forma diferente a los flujos de datos. En lugar de establecer el esquema en el origen o la transformación, configure esquemas en las conexiones de nodo entre los nodos del gráfico. Las transformaciones de bifurcación y filtro pueden validar opcionalmente los datos en tiempo de ejecución contra esquemas adjuntos a conexiones de nodos.

Cada entrada en el nodeConnections array puede incluir un schema en el from lateral de una conexión. Este esquema describe el formato esperado de los datos que fluyen entre esos dos nodos:

nodeConnections: [
  {
    from: {
      name: 'source'
      schema: {
        schemaRef: 'aio-sr://my-namespace/sensor-data:1'
        serializationFormat: 'Json'
      }
    }
    to: {
      name: 'transform'
    }
  }
]

El schemaRef valor utiliza el formato aio-sr://<namespace>/<name>:<version> y apunta a un esquema almacenado en el registro de esquemas. Dado que los grafos de flujo de datos solo soportan extremos MQTT, Kafka y OpenTelemetry, el formato de serialización soportado es Json.

La siguiente tabla resume cómo difiere la configuración del esquema entre flujos de datos y grafos de flujo de datos:

Aspecto Flujos de datos Gráficos de flujo de datos
Ubicación del esquema En origen (sourceSettings.schemaRef) y transformación (builtInTransformationSettings.schemaRef) En las conexiones de nodo (nodeConnections[].from.schema)
Formatos de destino admitidos JSON, Parquet, Delta JSON
Validación en tiempo de ejecución No soportado para esquemas fuente Opcional en conexiones de nodos mediante transformaciones de ramificación y filtro

Nota:

Para grafos de flujo de datos, JSON es actualmente el único formato de destino soportado, a pesar de los formatos listados en la documentación de referencia de la API REST.

Para definiciones de esquemas de mensajes, formatos y cómo subir esquemas, consulte Entender esquemas de mensajes.

Transformaciones integradas frente a transformaciones WASM

Los gráficos de flujo de datos admiten dos tipos de transformaciones:

  • Las transformaciones integradas están predefinidas por Microsoft (asignar, filtrar, ramificar, concatenar, ventanilla, limitar). Tú los configuras con reglas. No se requiere codificación.
  • Las transformaciones WASM son módulos de WebAssembly personalizados que los desarrolladores crean e implementan. Úselas cuando necesite lógica que las transformaciones integradas no cubran.

Ambos tipos de transformaciones se ejecutan dentro del mismo DataflowGraph recurso y puedes mezclarlas en una sola pipeline. Para obtener información sobre cómo crear e implementar transformaciones personalizadas, consulte Uso de transformaciones WASM en gráficos de flujo de datos.

Gestión de errores en grafos de flujo de datos

Cuando una transformación encuentra un error mientras procesa un mensaje (por ejemplo, un campo faltante o una expresión inválida), la transformación elimina el mensaje y registra un error. La canalización continúa procesando los mensajes posteriores.

Causas comunes de errores de procesamiento:

  • Un campo al que se hace referencia en la inputs regla no existe en el mensaje.
  • Una expresión de filtro o rama devuelve un valor no booleano.
  • Una expresión hace referencia a un tipo de dato incompatible (por ejemplo, un objeto JSON en aritmética).
  • No se puede acceder a un almacén de estado usado para el enriquecimiento.

Para supervisar los errores de procesamiento, compruebe los registros de pods del gráfico de flujo de datos o use los puntos de conexión de métricas. Para obtener más información, consulte Configuración de la observabilidad y la supervisión.

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.

Orientación de rendimiento para grafos de flujo de datos

Cada transformación de la canalización agrega sobrecarga de procesamiento. Tenga en cuenta estas directrices:

  • Prefiere menos transformaciones con más reglas. Si tiene muchas reglas de transformación que operan en la misma estructura, colóquelas en una sola transformación de mapa en lugar de crear transformaciones independientes para cada regla.
  • Use varias transformaciones cuando la lógica sea distinta. Las transformaciones independientes tienen sentido cuando los distintos pasos de procesamiento son fundamentalmente diferentes (filtrado frente a asignación frente a agregación).
  • Mantenga juntas las reglas relacionadas. Una sola transformación de mapa puede controlar el cambio de nombre de campo, la reestructuración, los campos calculados y las transformaciones de metadatos a la vez.