Azure Stream Analytics作业无IoT Hub输入数据故障排查问询
问题根因分析
- 消息体格式错误:创建Message时传入的普通字符串被序列化为字节数组存储在IoT Hub消息的body字段中,你在遥测数据里看到的数字键值对就是ASCII编码的字节值,Stream Analytics默认按JSON格式解析时无法识别字节数组,会直接判定为无效数据丢弃。
- 自定义属性格式不规范:写入custom_properties的payload是Python字典类型,存储时会被转换为单引号包裹的字符串,不符合双引号格式的标准JSON规范,即便要从属性字段提取数据也无法直接解析。
- Stream Analytics输入配置不匹配:如果作业输入的序列化配置设为JSON,但实际消息体为字节数组,或编码配置与消息实际编码不符,也会导致作业读取不到有效输入。
排查步骤
- 检查Stream Analytics输入配置:进入作业「输入」页面,查看IoT Hub输入的序列化设置,确认序列化格式、编码配置是否与发送的消息匹配,默认配置为JSON、UTF-8,你当前的消息体不符合该要求。
- 查看作业无效数据日志:进入作业「监视」-「日志」,运行如下查询确认是否存在数据解析错误:
AzureDiagnostics | where ResourceType == "STREAMANALYTICSJOBS" | where Category == "DataErrors" | project TimeGenerated, Resource, Message, ResultDescription
- 验证IoT Hub路由规则:进入IoT Hub「消息路由」-「测试路由」,上传你抓取的遥测数据样例,确认是否能匹配到Stream Analytics读取的内置事件中心端点(删除自定义路由后,IoT Hub默认会将消息转发到该端点)。
解决方法
修正Python发送代码
直接构造标准JSON格式的消息体传入Message,确保Stream Analytics可以正常解析:
import json import uuid import time from azure.iot.device import IoTHubDeviceClient, Message device_client = IoTHubDeviceClient.create_from_connection_string(conn_str) device_client.connect() for i in range(1, 180): print(f"sending message #{i}") # 构造结构化JSON消息体 msg_body = { "wind_speed": i, "temp": 1, "speed": 1 } # 转换为标准JSON字符串再封装为消息 msg = Message(json.dumps(msg_body)) msg.message_id = str(uuid.uuid4()) msg.correlation_id = "correlation-1234" msg.content_encoding = "utf-8" msg.content_type = "application/json" device_client.send_message(msg) time.sleep(3)
适配存量字节数组格式消息
如果需要解析之前发送的字节数组格式消息,可以在Stream Analytics中创建自定义UDF转换字节数组,查询语句如下:
SELECT UDF.ByteArrayToString(body) AS raw_body, properties.payload, systemProperties.[iothub-connection-device-id] AS device_id INTO streamdataoutput FROM streamdatainput
对应的UDF函数ByteArrayToString代码为:
function main(input) { return String.fromCharCode.apply(null, Object.values(input)); }
修改完成后重启Stream Analytics作业,测试查询时需同步运行Python脚本发送数据,测试查询仅读取最近几分钟的实时数据,停止发送时会提示无输入。
内容的提问来源于stack exchange,提问作者Luis Ramon Ramirez Rodriguez
相关产品推荐
相关产品推荐

