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

如何在Spark中将列数不同的多个JSON文件读取为一个DataFrame

解决Spark中嵌套结构不同的JSON文件合并问题

当两个JSON文件的嵌套结构存在差异时(比如第一个文件的a字段只有a1,第二个多了a2),直接使用union或unionByName会失败,核心原因是两个DataFrame的Schema不匹配——Spark要求合并的DataFrame必须具有完全一致的字段类型(包括嵌套Struct的内部结构)。

方法一:手动指定统一Schema(推荐)

提前定义包含所有可能字段的完整Schema,读取JSON时用该Schema约束,Spark会自动为缺失的字段填充null,确保两个DataFrame结构完全一致。

以Python为例:

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType

# 定义包含所有字段的完整Schema
full_schema = StructType([
    StructField("a", StructType([
        StructField("a1", StringType(), nullable=True),
        StructField("a2", StringType(), nullable=True)
    ]), nullable=True),
    StructField("b", StringType(), nullable=True)
])

# 初始化SparkSession
spark = SparkSession.builder.appName("MergeNestedJSON").getOrCreate()

# 用统一Schema读取两个JSON文件
df1 = spark.read.schema(full_schema).json("path/to/first.json")
df2 = spark.read.schema(full_schema).json("path/to/second.json")

# 合并DataFrame
combined_df = df1.unionByName(df2)

# 查看结果
combined_df.show(truncate=False)

方法二:动态合并Schema(灵活适配未知结构)

如果无法提前知晓完整结构,可以通过递归合并两个DataFrame的Schema,再将原DataFrame转换为合并后的Schema,最后完成合并。

Python实现示例:

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField

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

# 先读取两个JSON文件,获取各自的Schema
df1 = spark.read.json("path/to/first.json")
df2 = spark.read.json("path/to/second.json")

# 递归合并两个Schema的函数
def merge_two_schemas(schema1: StructType, schema2: StructType) -> StructType:
    field_map = {}
    # 先加入第一个Schema的所有字段
    for field in schema1.fields:
        field_map[field.name] = field
    # 遍历第二个Schema的字段,处理嵌套结构
    for field in schema2.fields:
        if field.name not in field_map:
            field_map[field.name] = field
        else:
            # 若字段是Struct类型,递归合并子Schema
            if isinstance(field.dataType, StructType) and isinstance(field_map[field.name].dataType, StructType):
                merged_sub_schema = merge_two_schemas(field_map[field.name].dataType, field.dataType)
                field_map[field.name] = StructField(
                    field.name, merged_sub_schema, field.nullable or field_map[field.name].nullable
                )
    return StructType(list(field_map.values()))

# 获取合并后的完整Schema
merged_schema = merge_two_schemas(df1.schema, df2.schema)

# 将两个DataFrame转换为合并后的Schema
df1_aligned = spark.createDataFrame(df1.rdd, merged_schema)
df2_aligned = spark.createDataFrame(df2.rdd, merged_schema)

# 合并DataFrame
combined_df = df1_aligned.unionByName(df2_aligned)

combined_df.printSchema()
combined_df.show(truncate=False)

关键说明

  • unionByName要求字段名完全匹配,但更严格的是字段的数据类型(包括嵌套结构)必须完全一致,这也是直接合并失败的核心原因。
  • 手动指定Schema的性能更优,适合结构固定的场景;动态合并Schema更灵活,适合结构不确定的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 23:55:22