You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

寻求用Azure Functions替代ASA分发IoT Hub数据至Azure SQL DB的示例及资源

替换Azure Stream Analytics为Azure Functions(Python/C#实现)

替换后的架构为 IoT edge devices > IoT hub > Azure Functions > SQL database,核心是用IoT Hub触发的Azure Functions替代Stream Analytics,完成JSON消息解析与SQL写入。

Python 实现步骤

1. 创建Azure Function(IoT Hub触发器)

选择事件中心触发器(IoT Hub的内置事件端点兼容事件中心协议),在Function App配置中添加以下应用设置:

  • IoTHubConnectionString: IoT Hub的Event Hub-compatible endpoint连接字符串
  • SQL_SERVER/SQL_DATABASE/SQL_USER/SQL_PASSWORD: SQL数据库的连接参数

2. 核心代码实现

使用pyodbc库连接SQL数据库,批量处理IoT Hub传入的JSON消息:

import json
import pyodbc
import os
from azure.functions import EventHubEvent

# 从应用设置加载SQL连接串
SQL_CONNECTION_STRING = (
    f"DRIVER={{ODBC Driver 17 for SQL Server}};"
    f"SERVER={os.environ['SQL_SERVER']};"
    f"DATABASE={os.environ['SQL_DATABASE']};"
    f"UID={os.environ['SQL_USER']};"
    f"PWD={os.environ['SQL_PASSWORD']}"
)

def main(events: list[EventHubEvent]):
    # 定义插入SQL语句(根据你的表结构调整字段)
    insert_query = """
        INSERT INTO DeviceTelemetry (DeviceId, Timestamp, Temperature, Humidity)
        VALUES (?, ?, ?, ?)
    """
    
    try:
        # 使用连接池建立SQL连接
        with pyodbc.connect(SQL_CONNECTION_STRING) as conn:
            with conn.cursor() as cursor:
                for event in events:
                    try:
                        # 解析JSON消息体
                        msg_body = json.loads(event.get_body().decode('utf-8'))
                        # 从消息元数据获取设备ID
                        device_id = event.metadata["iothub-connection-device-id"]
                        
                        # 提取字段(根据你的JSON结构调整)
                        timestamp = msg_body["timestamp"]
                        temperature = msg_body["temperature"]
                        humidity = msg_body["humidity"]
                        
                        # 执行单条插入
                        cursor.execute(insert_query, (device_id, timestamp, temperature, humidity))
                    
                    except json.JSONDecodeError:
                        print(f"无效的JSON消息: {event.get_body()}")
                        continue
                    except KeyError as e:
                        print(f"消息缺少必填字段: {str(e)}")
                        continue
                
                conn.commit()
                print(f"成功处理 {len(events)} 条消息")
    
    except pyodbc.Error as e:
        print(f"SQL操作错误: {str(e)}")

3. 依赖配置

在requirements.txt中添加依赖包:

azure-functions
pyodbc

C# 实现示例

使用ServiceBusTrigger(IoT Hub兼容服务总线端点),结合SqlClient完成写入:

using System;
using System.Data.SqlClient;
using System.Text.Json;
using Microsoft.Azure.WebJobs;
using Microsoft.Extensions.Logging;

namespace IoTHubToSqlFunction
{
    public static class TelemetryProcessor
    {
        [FunctionName("TelemetryProcessor")]
        public static void Run(
            [ServiceBusTrigger("%EventHubName%", Connection = "IoTHubConnectionString")] string[] messages,
            ILogger log)
        {
            var sqlConnString = Environment.GetEnvironmentVariable("SqlConnectionString");
            var insertQuery = @"INSERT INTO DeviceTelemetry (DeviceId, Timestamp, Temperature, Humidity) 
                               VALUES (@DeviceId, @Timestamp, @Temperature, @Humidity)";

            using (var conn = new SqlConnection(sqlConnString))
            {
                conn.Open();
                using (var cmd = new SqlCommand(insertQuery, conn))
                {
                    foreach (var message in messages)
                    {
                        try
                        {
                            var telemetry = JsonSerializer.Deserialize<TelemetryData>(message);
                            // 从触发器元数据获取设备ID
                            var deviceId = System.Environment.GetEnvironmentVariable("IotHubDeviceId");

                            cmd.Parameters.Clear();
                            cmd.Parameters.AddWithValue("@DeviceId", deviceId);
                            cmd.Parameters.AddWithValue("@Timestamp", telemetry.Timestamp);
                            cmd.Parameters.AddWithValue("@Temperature", telemetry.Temperature);
                            cmd.Parameters.AddWithValue("@Humidity", telemetry.Humidity);

                            cmd.ExecuteNonQuery();
                        }
                        catch (JsonException ex)
                        {
                            log.LogError(ex, "解析JSON消息失败");
                            continue;
                        }
                        catch (SqlException ex)
                        {
                            log.LogError(ex, "SQL写入失败");
                            continue;
                        }
                    }
                }
            }
            log.LogInformation($"已处理 {messages.Length} 条消息");
        }
    }

    public class TelemetryData
    {
        public DateTime Timestamp { get; set; }
        public double Temperature { get; set; }
        public double Humidity { get; set; }
    }
}

关键优化点

  • 批量处理: 调整Function的批量消息阈值,减少SQL连接次数
  • 连接池: 依赖数据库驱动默认的连接池机制,避免频繁创建连接
  • 错误重试: 可引入Polly库实现SQL操作的重试逻辑
  • 监控: 启用Application Insights追踪函数执行状态与错误
  • 成本控制: 使用消费计划时,平衡批量大小与函数执行频率,降低冷启动开销

内容的提问来源于stack exchange,提问作者Arseniy Sleptsov

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 17:17:32