在Databricks中从REST API获取无Schema嵌套JSON转DataFrame的典型方案
在Databricks里处理这种无固定Schema的嵌套JSON数据我熟得很,毕竟经常遇到API返回结构变来变去的情况。下面分享几个最实用的方法,一步步来:
1. 先从REST API拿到JSON数据
首先得把API的数据拉下来,用Python的requests库就行,Databricks环境默认已经装了这个库,直接用:
import requests # 替换成你的API地址,需要认证的话加headers api_url = "https://your-api-endpoint.com/data" headers = {"Authorization": "Bearer YOUR_AUTH_TOKEN"} # 如果需要认证的话 response = requests.get(api_url, headers=headers) # 把响应转成JSON对象 raw_json = response.json()
这里要注意:如果API返回的是JSON数组(比如[{"a":1}, {"b":2}]),直接用就行;如果是单个JSON对象(比如{"a":1, "nested":{"c":3}}),后面创建DataFrame的时候要把它包成列表。
2. 动态推断Schema创建DataFrame(核心方法)
因为Schema不固定还会变,不能提前写死StructType,所以用Spark的自动推断Schema功能最省心。
方法一:从JSON字符串RDD创建DataFrame
先把内存里的JSON转成字符串列表,再转成RDD,最后用spark.read.json自动推断Schema:
import json from pyspark.sql import SparkSession # 处理单个对象/数组的情况 if isinstance(raw_json, list): json_strings = [json.dumps(item) for item in raw_json] else: json_strings = [json.dumps(raw_json)] # 转成RDD json_rdd = spark.sparkContext.parallelize(json_strings) # 自动推断Schema,开启mergeSchema应对后续结构变化 df = spark.read.option("mergeSchema", "true").json(json_rdd) # 看看结果 df.show(truncate=False)
mergeSchema=true这个参数很重要!如果后续API返回的JSON多了新字段或者改了结构,重新运行的时候Spark会自动合并新旧Schema,不会因为结构变化报错。
方法二:用from_json解析字符串列(适合分批/流式场景)
如果是分批拉取API数据,或者要把JSON字符串存在DataFrame里再解析,用这个方法更灵活:
from pyspark.sql.functions import from_json, col # 先把原始JSON转成只有一个字符串列的DataFrame single_row_df = spark.createDataFrame([(json.dumps(raw_json),)], ["raw_json_str"]) # 先从字符串列推断Schema inferred_schema = spark.read.json(single_row_df.select("raw_json_str").rdd.map(lambda x: x[0])).schema # 解析JSON字符串列,展开所有字段 parsed_df = single_row_df.withColumn("parsed_data", from_json(col("raw_json_str"), inferred_schema)) \ .select("parsed_data.*") parsed_df.show(truncate=False)
这个方法的好处是可以单独获取Schema,后续如果要处理大量同结构数据,可以复用这个Schema,提升性能。
3. 处理嵌套结构的小技巧
如果JSON嵌套很深,想要扁平化方便分析,可以试试这些方法:
手动展开嵌套字段
如果知道常用的嵌套字段,直接用selectExpr展开:
flattened_df = df.selectExpr( "id", "user_info.name as user_name", "user_info.email as user_email", "order_details[0].product as first_product" # 处理数组字段 )
通用扁平化函数(自动处理所有嵌套Struct)
如果嵌套层级多且结构不固定,可以写个递归函数自动扁平化:
def flatten_dataframe(df): # 区分普通列和嵌套Struct列 flat_columns = [col_name for col_name, col_type in df.dtypes if col_type[:6] != "struct"] nested_columns = [col_name for col_name, col_type in df.dtypes if col_type[:6] == "struct"] # 展开每个嵌套Struct列 for nested_col in nested_columns: # 获取嵌套列里的所有子字段 sub_fields = df.select(f"{nested_col}.*").columns # 重命名子字段避免冲突 expanded_fields = [f"{nested_col}.{sub_field} as {nested_col}_{sub_field}" for sub_field in sub_fields] # 替换原嵌套列为展开后的字段 df = df.select(flat_columns + expanded_fields) # 递归处理新的嵌套列(如果有的话) if len(nested_columns) > 0: return flatten_dataframe(df) else: return df # 用这个函数扁平化DataFrame flat_df = flatten_dataframe(df) flat_df.show(truncate=False)
注意:这个函数只处理Struct类型的嵌套,如果有数组(ArrayType),需要先用explode函数展开数组,再进行扁平化。
4. 注意事项
- 大数据量处理:如果API返回的数据特别大,别一次性加载到内存,用API的分页参数分批拉取,或者用Spark的流式处理来分批解析。
- Schema推断性能:数据量很大时,全量推断Schema会慢,可以先采样一部分数据(比如取前1000条)来推断Schema,再用这个Schema解析全量数据。
- Schema冲突处理:如果新旧Schema的数据类型冲突(比如原来的int变成了string),
mergeSchema可能会报错,这时候需要手动指定兼容的Schema,或者在解析前做数据类型转换。
内容的提问来源于stack exchange,提问作者Kertis van Kertis

