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
相关产品推荐
相关产品推荐

