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

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_idadf_log_id
adf_idadf_id
e_namee_name
e_desce_desc
statusstatus
failure_msgfailure_msg

4. 运行验证

触发管道运行,检查Azure SQL表中是否成功写入数据。

关键注意事项

  • 确保Databricks和ADF都拥有Azure存储的访问权限(可通过服务主体或SAS令牌配置)。
  • Azure SQL表的e_desc字段建议设置为NVARCHAR(MAX)类型,以容纳完整的JSON内容。
  • 使用Delta格式时,ADF需确保连接配置正确,支持读取Delta表的元数据。

内容的提问来源于stack exchange,提问作者Developer Rajinikanth

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 01:20:31