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

如何为Parquet动态嵌套列创建静态Schema并转换M字段?

问题描述

我正在处理包含复杂嵌套列结构的Parquet文件,其中某列结构如下:

{
   "dynamodb": {
      "NewImage": {
         "DiscountData": {
            "M": {
               "TEST1": {"N": "0"},
               "TEST2": {"N": "0"},
               "TEST3": {"N": "0"},
               "TEST4": {"N": "0"},
               "TEST5": {"N": "1"}
            }
         }
      }
   }
}

核心问题在于M字段是动态的:其包含的键(如TEST1、TEST2等)会随数据集变化,且不同Parquet文件的M结构可能不同。

我尝试使用Schema自动推断功能,却遇到如下错误:

AnalysisException: Found duplicate column(s) in the data schema 

为保证处理一致性,我希望:将M字段的全部内容转为JSON字符串,或者将M转换为键值均为字符串类型的MapType。

我尝试定义如下Schema:

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

schema = StructType([
    StructField("dynamodb", StructType([
        StructField("NewImage", StructType([
            StructField("DiscountData", StructType([
                StructField("M", MapType(StringType(), StringType()), True)
            ]), True)
        ]), True)
    ]), True)
])

并通过以下代码创建DataFrame:

df_disscountcoderule = spark.read \
    .schema(schema) \
    .option("mergeSchema", "true") \
    .option("recursiveFileLookup", "true") \
    .parquet(f"s3://{source_bucket}/{source_discountcodesrules_prefix}")\
    .filter(input_file_name().endswith(".parquet"))\
    .filter(~input_file_name().contains("processing-failed"))

但应用该Schema时出现如下错误:

java.lang.ClassCastException: class org.apache.spark.sql.types.MapType cannot be cast to class org.apache.spark.sql.types.StructType (org.apache.spark.sql.types.MapType and org.apache.spark.sql.types.StructType are in unnamed module of loader 'app')

我的问题:

  1. 是否可以创建能处理该动态结构的静态Schema?
  2. 如何将M字段转换为JSON字符串或MapType以实现一致性处理?

解决方案

问题1:是否可以创建处理动态结构的静态Schema?

可以,但不能直接把M定义为MapType——因为原始Parquet文件中M的实际存储类型是StructType(动态键会被解析为Struct的字段),直接用MapType绑定会触发类型转换错误。

正确的做法是先用一个兼容任意子结构的静态Schema读取数据,再通过后续转换将动态Struct转为Map或JSON字符串。比如定义Schema时把M设为空StructType(),Spark会自动兼容其下的任意动态字段。

问题2:转换为JSON字符串或MapType的实现方案

方案一:将M字段转为JSON字符串

读取数据时先使用兼容动态结构的Schema,再用to_json函数把M字段序列化为JSON字符串:

from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField

# 定义兼容动态结构的静态Schema:空Struct接受任意子字段
schema = StructType([
    StructField("dynamodb", StructType([
        StructField("NewImage", StructType([
            StructField("DiscountData", StructType([
                StructField("M", StructType(), True)
            ]), True)
        ]), True)
    ]), True)
])

# 读取数据(关闭mergeSchema避免自动推断冲突)
df = spark.read \
    .schema(schema) \
    .option("recursiveFileLookup", "true") \
    .parquet(f"s3://{source_bucket}/{source_discountcodesrules_prefix}")\
    .filter(F.input_file_name().endswith(".parquet"))\
    .filter(~F.input_file_name().contains("processing-failed"))

# 将M字段转为JSON字符串
df_with_json = df.withColumn(
    "DiscountData_M_JSON",
    F.to_json("dynamodb.NewImage.DiscountData.M")
)

# 可选:保留目标字段,丢弃原始嵌套结构
df_final = df_with_json.select("DiscountData_M_JSON")

方案二:将M字段转为键值均为字符串的MapType

利用map_from_entries、map_entries和transform函数,把动态Struct的字段转为键值对,再聚合为Map。如果原始M中的值是嵌套结构(如{"N": "0"}),还需要先提取字符串值:

from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField

# 同样使用兼容动态结构的Schema读取数据
schema = StructType([
    StructField("dynamodb", StructType([
        StructField("NewImage", StructType([
            StructField("DiscountData", StructType([
                StructField("M", StructType(), True)
            ]), True)
        ]), True)
    ]), True)
])

df = spark.read \
    .schema(schema) \
    .option("recursiveFileLookup", "true") \
    .parquet(f"s3://{source_bucket}/{source_discountcodesrules_prefix}")\
    .filter(F.input_file_name().endswith(".parquet"))\
    .filter(~F.input_file_name().contains("processing-failed"))

# 将动态Struct转为Map<String, String>,同时提取嵌套值中的字符串内容
df_with_map = df.withColumn(
    "DiscountData_M_Map",
    F.map_from_entries(
        F.transform(
            F.map_entries("dynamodb.NewImage.DiscountData.M"),
            lambda x: F.struct(x.key, x.value.getItem("N").cast("string"))
        )
    )
)

# 可选:保留目标字段
df_final = df_with_map.select("DiscountData_M_Map")

关键注意事项

  • 读取阶段不要直接用MapType绑定Schema:Parquet中动态键是以Struct字段存储的,Spark无法直接将Struct转为Map,必须在读取后通过转换函数处理。
  • 关闭mergeSchema:动态结构会导致Schema冲突,提前用静态Schema限定结构可以避免自动推断时的重复列错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 01:47:28