Azure Databricks导出JSON数组至存储并通过ADF写入Azure SQL的方案问询
解决方案:Azure Databricks导出JSON至存储 + ADF写入Azure SQL
一、Azure Databricks:JSON数组转DataFrame并写入存储
1. 加载JSON数组并转换为目标结构DataFrame
直接将示例JSON数组转为Spark DataFrame,再通过字段映射、类型转换匹配Azure SQL表结构:
# 示例JSON数组 json_data = [ {'Details': {'Input': {'id': '1', 'name': 'asdsdasd', 'a1': None, 'a2': None, 'c': None, 's': None, 'c1': None, 'z': None}, 'Output': '{"msg":"some error"}'}, 'Failure': '{"msg":"error"}', 's': 'f'}, {'Details': {'Input': {'id': '2', 'name': 'sadsadsad', 'a1': 'adsadsad', 'a2': 'sssssss', 'c': 'cccc', 's': 'test', 'c1': 'ind', 'z': '22222'}, 'Output': '{"s":"2"}'}, 'Failure': '', 's': 's'} ] # 转为基础DataFrame df = spark.createDataFrame(json_data) # 导入必要函数处理字段映射 from pyspark.sql.functions import col, struct, to_json, monotonically_increasing_id, when # 转换为匹配SQL表的结构 transformed_df = df \ .withColumn("adf_log_id", monotonically_increasing_id() + 1) # 生成自增ID .withColumn("adf_id", col("Details.Input.id")) .withColumn("e_name", col("Details.Input.name")) # 将Input和Output合并为JSON字符串作为e_desc .withColumn("e_desc", to_json(struct(col("Details.Input").alias("Input"), col("Details.Output").alias("Output")))) # 根据`s`字段映射status .withColumn("status", when(col("s") == "s", "success").when(col("s") == "f", "failure").otherwise("unknown")) .withColumn("failure_msg", col("Failure")) # 保留目标字段 final_df = transformed_df.select("adf_log_id", "adf_id", "e_name", "e_desc", "status", "failure_msg") # 预览结果 final_df.show(truncate=False)
2. 写入Azure存储(ADLS Gen2/Blob)
推荐使用Delta Lake格式(支持ACID事务、版本控制),也可选择Parquet:
写入Delta格式
# 替换为你的Azure存储路径(ADLS Gen2示例:abfss://容器名@存储账户名.dfs.core.windows.net/路径) delta_path = "abfss://your-container@your-storage-account.dfs.core.windows.net/adf-logs/delta" # 写入Delta表(append模式支持增量写入) final_df.write.format("delta").mode("append").save(delta_path) # 可选:创建Spark外部表方便后续查询 spark.sql(f""" CREATE TABLE IF NOT EXISTS adf_logs USING delta LOCATION '{delta_path}' """)
写入Parquet格式
parquet_path = "abfss://your-container@your-storage-account.dfs.core.windows.net/adf-logs/parquet" final_df.write.format("parquet").mode("append").save(parquet_path)
二、Azure Data Factory:读取存储数据写入Azure SQL
1. 创建数据集
- 源数据集:选择
Delta Lake(用Delta格式时)或Parquet(用Parquet格式时),配置Azure存储账户连接,指定文件/表路径。 - 目标数据集:选择
Azure SQL Database,配置数据库连接,指定目标表(需提前创建好结构匹配的表)。
2. 配置Copy Data活动
在管道中添加Copy Data活动,完成以下配置:
- 源端:选择创建好的源数据集,设置读取模式(全量/增量,Delta支持增量读取)。
- 目标端:选择目标SQL数据集,设置写入行为(推荐
Append模式,避免覆盖已有数据)。
3. 字段映射
在活动的Mapping标签页,确保字段一一对应:
| DataFrame字段 | SQL表字段 |
|---|---|
| adf_log_id | adf_log_id |
| adf_id | adf_id |
| e_name | e_name |
| e_desc | e_desc |
| status | status |
| failure_msg | failure_msg |
4. 运行验证
触发管道运行,检查Azure SQL表中是否成功写入数据。
关键注意事项
- 确保Databricks和ADF都拥有Azure存储的访问权限(可通过服务主体或SAS令牌配置)。
- Azure SQL表的
e_desc字段建议设置为NVARCHAR(MAX)类型,以容纳完整的JSON内容。 - 使用Delta格式时,ADF需确保连接配置正确,支持读取Delta表的元数据。
内容的提问来源于stack exchange,提问作者Developer Rajinikanth
相关产品推荐
相关产品推荐

