如何提取JSON API响应的result数组记录并写入Databricks Delta Table
Databricks解析API响应JSON并存为可追加Delta表方案
问题根因
你遇到的两个异常均为参数配置错误导致:
- 读取后所有字段为null:默认
spark.read.json按单行解析JSON,你存储的API响应是跨行的完整JSON对象,未开启跨行解析参数会解析失败返回空值。 - 仅能读取result数组单条记录:未对result数组做爆炸展开操作,Spark默认将整个数组识别为单个行的单个字段值。
具体实现步骤
1 读取JSON文件
开启multiLine参数支持跨行JSON解析:
# 路径替换为你datalake中存放API响应JSON的目录路径 raw_df = spark.read.option("multiLine", "true").json("abfss://容器@存储账号.dfs.core.chinacloudapi.cn/json_path/*")
2 提取并展开result数组
用explode函数将数组中的每条record拆为独立行,再展开record内部的所有字段为结构化列:
from pyspark.sql.functions import explode, col structured_df = raw_df.select(explode(col("result")).alias("record")) \ .select("record.*")
执行后structured_df的列即为每个record中的field1、field2等字段,可直接用于分析。
3 写入Delta表(支持增量追加)
首次创建表可使用覆盖模式,后续新增API响应数据执行时替换为追加模式即可:
# 写入Delta表,开启schema自动演进兼容后续record新增字段 structured_df.write.format("delta") \ .mode("append") # 首次创建可替换为"overwrite" .option("mergeSchema", "true") \ .saveAsTable("api_result_delta") # 也可填路径存为外部表 save("你的Delta存储路径")
自动化流程配置
要实现新API响应自动追加,可选两种方案:
- 定时调度:在Databricks工作流中配置定时任务,定期扫描JSON目录新增文件,执行上述处理逻辑写入Delta表。
- 流处理自动触发:用Auto Loader监控JSON目录,新文件上传时自动解析写入:
stream_raw_df = spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "json") \ .option("multiLine", "true") \ .option("cloudFiles.schemaLocation", "你的schema校验点路径") \ .load("你的JSON存放目录路径") stream_structured_df = stream_raw_df.select(explode(col("result")).alias("record")) \ .select("record.*") query = stream_structured_df.writeStream.format("delta") \ .option("checkpointLocation", "你的写入校验点路径") \ .option("mergeSchema", "true") \ .trigger(availableNow=True) # 按需执行批次,也可配置为持续触发 .toTable("api_result_delta") query.awaitTermination()
内容的提问来源于stack exchange,提问作者Davio
相关产品推荐
相关产品推荐

