如何配置Spark Read Stream Schema读取多Schema JSON并写入Delta表
用Databricks Autoloader处理异构Schema JSON并生成扁平Delta表
核心思路
两种JSON的差异在于Col2字段:一种是单个结构体,另一种是结构体数组。我们需要先将Col2统一转为数组格式,再展开数组提取subcol1,同时统一col1的大小写命名。
具体实现步骤
1. 读取异构JSON数据(Autoloader配置)
使用Autoloader读取数据时,开启Schema自动合并以兼容两种结构,同时保留字段的大小写敏感性:
df = (spark.readStream .format("cloudFiles") .option("cloudFiles.format", "json") .option("cloudFiles.schemaLocation", "/path/to/schema-storage") # 存储自动推断的Schema .option("mergeSchema", "true") # 合并不同文件的Schema .option("caseSensitive", "true") # 适配Col1和col1的大小写差异 .load("/path/to/json-source"))
2. 统一col1字段命名
将大小写不同的Col1和col1合并为统一的col1字段:
from pyspark.sql.functions import coalesce df = df.withColumn("col1", coalesce(df.Col1, df.col1)) .drop("Col1") # 移除原大写字段
3. 统一Col2为数组格式
判断Col2的类型,将单个结构体转为包含一个元素的数组:
from pyspark.sql.functions import when, array, col df = df.withColumn( "Col2_array", when( col("Col2").isinstance("struct"), array(col("Col2")) # 单个结构体转数组 ).otherwise(col("Col2")) # 已为数组则保持不变 )
4. 展开数组并提取子字段
使用explode展开数组,提取subcol1并重命名为目标字段:
from pyspark.sql.functions import explode df_flatten = df.select( "col1", explode("Col2_array").alias("Col2_struct") ).select( "col1", col("Col2_struct.subcol1").alias("col2_subcol1") )
5. 写入Delta表
将扁平后的数据流写入Delta表:
(df_flatten.writeStream .format("delta") .option("checkpointLocation", "/path/to/checkpoint-folder") .start("/path/to/target-delta-table"))
关键说明
schemaLocation必须指定,Autoloader会自动管理Schema演化,后续新增结构也能兼容mergeSchema开启后,Spark会自动合并不同文件的Schema,避免因结构差异导致读取失败- 类型判断时,
isinstance方法用于检测字段的结构类型,确保单个结构体被正确转为数组
内容的提问来源于stack exchange,提问作者mytabi
相关产品推荐
相关产品推荐

