如何按id拆分IoT Hub流数据并分层存储至Blob Storage?
你之前用GetArrayElement没得到预期结果,大概率是用错了函数——应该用**GetArrayElements(复数)**配合CROSS APPLY来展平整个values数组,而不是单数版本(单数只能提取数组中指定索引的单个元素)。下面是完整的实现步骤:
1. 先展平Values数组
首先把输入的JSON数组拆成单条记录,这样才能对每个values元素单独处理:
SELECT input.timestamp AS eventTimestamp, flattenedValues.ArrayValue.id AS metricId, flattenedValues.ArrayValue.v AS value, flattenedValues.ArrayValue.q AS quality, flattenedValues.ArrayValue.t AS metricTimestamp INTO [BlobOutput] FROM [IoTHubInput] AS input CROSS APPLY GetArrayElements(input.values) AS flattenedValues
CROSS APPLY会把原输入的每条带数组的记录,拆分成N条独立记录(N是values数组的元素个数),每条对应一个values元素。
2. 拆分ID的层级字段
假设你的ID格式是ChannelX/DeviceY或ChannelX/Functions/AlertZ,我们需要提取前两层作为分组和路径依据。用CHARINDEX和SUBSTRING拆分:
SELECT input.timestamp AS eventTimestamp, -- 提取第一层级:Channel SUBSTRING(flattenedValues.ArrayValue.id, 1, CHARINDEX('/', flattenedValues.ArrayValue.id) - 1) AS Channel, -- 提取第二层级:Device或Functions CASE WHEN CHARINDEX('/', flattenedValues.ArrayValue.id, CHARINDEX('/', flattenedValues.ArrayValue.id) + 1) > 0 THEN SUBSTRING(flattenedValues.ArrayValue.id, CHARINDEX('/', flattenedValues.ArrayValue.id) + 1, CHARINDEX('/', flattenedValues.ArrayValue.id, CHARINDEX('/', flattenedValues.ArrayValue.id) + 1) - CHARINDEX('/', flattenedValues.ArrayValue.id) - 1) ELSE SUBSTRING(flattenedValues.ArrayValue.id, CHARINDEX('/', flattenedValues.ArrayValue.id) + 1, LEN(flattenedValues.ArrayValue.id)) END AS DeviceOrFunction, flattenedValues.ArrayValue.v AS value, flattenedValues.ArrayValue.q AS quality, flattenedValues.ArrayValue.t AS metricTimestamp INTO [BlobOutput] FROM [IoTHubInput] AS input CROSS APPLY GetArrayElements(input.values) AS flattenedValues
这段代码会自动处理ID是否有第三层级的情况,确保第二层级能正确提取。
3. 按层级分组重组数据
如果需要把同一层级(同一Channel+DeviceOrFunction)的指标数据重新组合成数组(和原输入结构一致),用COLLECT函数聚合:
SELECT eventTimestamp, Channel, DeviceOrFunction, -- 把同层级的指标聚合成数组 COLLECT( { "v": value, "q": quality, "t": metricTimestamp } ) AS values INTO [BlobOutput] FROM ( -- 子查询:先展平数组并拆分层级 SELECT input.timestamp AS eventTimestamp, SUBSTRING(flattenedValues.ArrayValue.id, 1, CHARINDEX('/', flattenedValues.ArrayValue.id) - 1) AS Channel, CASE WHEN CHARINDEX('/', flattenedValues.ArrayValue.id, CHARINDEX('/', flattenedValues.ArrayValue.id) + 1) > 0 THEN SUBSTRING(flattenedValues.ArrayValue.id, CHARINDEX('/', flattenedValues.ArrayValue.id) + 1, CHARINDEX('/', flattenedValues.ArrayValue.id, CHARINDEX('/', flattenedValues.ArrayValue.id) + 1) - CHARINDEX('/', flattenedValues.ArrayValue.id) - 1) ELSE SUBSTRING(flattenedValues.ArrayValue.id, CHARINDEX('/', flattenedValues.ArrayValue.id) + 1, LEN(flattenedValues.ArrayValue.id)) END AS DeviceOrFunction, flattenedValues.ArrayValue.v AS value, flattenedValues.ArrayValue.q AS quality, flattenedValues.ArrayValue.t AS metricTimestamp FROM [IoTHubInput] AS input CROSS APPLY GetArrayElements(input.values) AS flattenedValues ) AS flattenedData GROUP BY eventTimestamp, Channel, DeviceOrFunction
COLLECT会把分组内的所有指标打包成一个JSON数组,输出结构和原输入匹配,但已经按层级拆分完成。
4. 配置Blob输出的路径规则
在Stream Analytics的Blob输出配置里,设置路径模式为:
Channel/{Channel}/{DeviceOrFunction}/{date}/{time}
这里的{Channel}和{DeviceOrFunction}会自动替换为查询中对应的字段值,自动创建主文件夹(Channel命名)和子文件夹(Device/Functions命名)。{date}和{time}是内置变量,用于按时间分片存储,不需要的话可以去掉。
输出格式选JSON,序列化方式根据需求选Array(整个文件是一个JSON数组)或Line separated(每行一条JSON记录)。
常见坑点
- 用错函数:一定要用
GetArrayElements(复数)+CROSS APPLY来展平数组,单数GetArrayElement只能取单个元素,无法遍历整个数组。 - ID格式不统一:如果存在没有
/的ID,需要在CASE语句里加默认逻辑(比如归到DefaultChannel),避免拆分报错。
内容的提问来源于stack exchange,提问作者Nilay

