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

如何配置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 17:20:27