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

PySpark解析API返回的复杂JSON为DataFrame遇阻求助

解决PySpark中API返回JSON转DataFrame并写入数据湖的方案

1. 解析JSON字符串,提取目标数据片段

首先将API返回的JSON字符串转换为Python字典,直接提取QueryResults.Records核心数据:

import json
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("APIJsonToDF").getOrCreate()

# 假设retjson是API返回的JSON字符串
json_dict = json.loads(retjson)
# 提取QueryResults下的Records列表,无数据时返回空列表
target_records = json_dict.get("QueryResults", {}).get("Records", [])

2. 扁平化数据,取出value子节点内容

Records中的字段都嵌套在value节点下,需要将每个Record的value展开为平级字典:

flattened_data = [record.get("value", {}) for record in target_records]

3. 创建DataFrame(两种可选方式)

方式一:自动推断Schema

扁平化后的数据结构简单,可让Spark自动推断Schema,避免手动定义的错误:

df = spark.createDataFrame(flattened_data)
# 验证数据结构
df.printSchema()
df.show(5)

方式二:使用自定义Schema

若自动推断不符合预期,或需要严格约束字段类型,需确保自定义Schema的字段名、类型与flattened_data的键完全匹配:

from pyspark.sql.types import StructType, StructField, StringType, IntegerType, TimestampType

# 替换为你的实际字段及类型
custom_schema = StructType([
    StructField("user_id", IntegerType(), nullable=True),
    StructField("user_name", StringType(), nullable=True),
    StructField("login_time", TimestampType(), nullable=True)
])

df = spark.createDataFrame(data=flattened_data, schema=custom_schema)

4. 持久化到数据湖

以Delta Lake(主流数据湖格式)为例,支持覆盖、追加等写入模式;若使用Parquet、ORC等格式,替换format参数即可:

# 写入数据湖路径(支持HDFS、S3、ADLS等)并注册为表
df.write.format("delta")\
    .mode("overwrite")\
    .option("path", "/datalake/query_results")\
    .saveAsTable("datalake_db.query_results")

# 若仅需写入路径,无需注册为表
# df.write.format("delta").mode("append").save("/datalake/query_results")

错误排查要点

  • 无法推断Schema:检查flattened_data是否为空,或原始JSON嵌套层级过深导致推断失败,扁平化后重试即可。
  • TypeError:确认flattened_data是字典组成的列表,同时检查自定义Schema的字段名、数据类型与实际数据完全匹配(比如不要将字符串字段定义为IntegerType)。

内容的提问来源于stack exchange,提问作者Scott S

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 19:19:58