如何让Azure流分析作业将Event Hub数据以JSON数组存入Blob?
Azure流分析输出JSON数组到Blob存储的解决方案
问题背景
多个程序向Azure Event Hub推送独立的JSON文档,通过Azure流分析(ASA)作业写入Azure Data Lake Storage Gen2后,数据以每行一个顶级JSON对象的格式存储(NDJSON),但需要将这些对象合并为标准的JSON数组格式。
可行解决方法
方法1:通过ASA窗口聚合生成JSON数组
这是最直接的实时处理方案,利用ASA的聚合函数将窗口内的事件打包为数组:
- 修改ASA作业的查询语句,使用窗口函数(根据业务需求选择滚动窗口、滑动窗口或跳跃窗口),示例查询:
SELECT CollectArray(input) AS jsonArray INTO [你的ADLS Gen2输出别名] FROM [你的Event Hub输入别名] TIMESTAMP BY EventEnqueuedUtcTime GROUP BY TumblingWindow(minute, 5) -- 按5分钟窗口聚合,可调整时间粒度
- 配置输出序列化:在ASA的输出设置中,选择JSON格式,ASA会自动将聚合后的数组序列化为标准JSON数组格式写入Blob。
该方案适合可以接受按时间批量聚合的场景,生成的每个Blob文件对应一个时间窗口内的事件数组。
方法2:Blob写入后批量转换
如果无法调整ASA的实时处理逻辑,可在数据写入Blob后通过Azure函数或Azure Data Factory做后续转换:
- 触发机制:设置Blob存储的事件触发Azure函数,当新Blob写入时自动执行转换
- 处理逻辑:读取Blob中的NDJSON内容,将每行JSON对象合并为数组,覆盖原Blob或写入新路径
示例Azure函数Python代码:
import json import azure.functions as func def main(blob_trigger: func.InputStream, output_blob: func.Out[str]): # 读取Blob内容并按行分割 raw_lines = blob_trigger.read().decode('utf-8').splitlines() # 过滤空行并解析为JSON对象 json_objects = [json.loads(line) for line in raw_lines if line.strip()] # 转换为JSON数组并写入目标Blob output_blob.set(json.dumps(json_objects, indent=2))
方法3:单条事件包装为数组(特殊场景)
如果需要每条事件单独生成为单元素数组,可直接在ASA查询中包装:
SELECT Array(input) AS jsonArray INTO [你的ADLS Gen2输出别名] FROM [你的Event Hub输入别名]
此方案仅适合必须单条生成数组的特殊场景,一般不推荐用于批量数据存储。
总结
Azure流分析完全支持输出JSON数组格式,无需更换工具,只需调整查询逻辑或配合后续处理即可满足需求,并非对ASA要求过高。
内容的提问来源于stack exchange,提问作者Stuart Hemming
相关产品推荐
相关产品推荐

