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

基于嵌套对象数据类型过滤并展平Azure Synapse PySpark中的嵌套JSON对象

解决Azure Synapse PySpark中混合类型JSON字段的展平问题

我之前也碰到过类似的混合类型JSON字段处理场景,Spark自动推断schema时确实会因为字段类型不一致,统一选用最兼容的类型(比如这里的string),导致我们没法直接对struct类型的记录做展平操作。下面是具体的解决步骤和代码示例:

问题根源

Spark读取JSON数据时,默认会根据所有记录的字段类型推断schema。当同一个字段存在多种类型(比如你的cords既有struct又有数值),Spark会选择一个能容纳所有类型的通用类型(这里就是string),这就直接阻碍了我们对struct类型记录的嵌套字段展平。

解决方案思路

我们可以先把cords当作string类型读取,再通过尝试解析为struct的方式区分两种记录:

  1. 对能成功解析为struct的记录,展平其嵌套字段;
  2. 对解析失败的记录,保留原cords字段,同时将展平后的字段设为null(或者你需要的默认值);
  3. 最后合并两种处理后的结果。

具体代码实现

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, from_json, lit
from pyspark.sql.types import StructType, StructField, DoubleType

# Synapse环境中可以省略这一步,直接使用内置的spark对象
spark = SparkSession.builder.appName("FlattenMixedCords").getOrCreate()

# 模拟读取你的JSON数据(实际场景替换为spark.read.json("你的存储路径"))
json_records = [
    '{"dateTime":"2020-11-29T13:51:16.168659Z","cords":{"x_al":0.0191342489,"y_al":-0.1200904993}}',
    '{"dateTime":"2020-12-29T13:51:21.457739Z","cords":51.0}',
    '{"dateTime":"2021-10-29T13:51:26.634289Z","cords":{"x_al":0.01600042489,"y_al":-0.1200900993}}'
]
df = spark.read.json(spark.sparkContext.parallelize(json_records))

# 定义cords字段对应的struct schema
cords_struct_schema = StructType([
    StructField("x_al", DoubleType(), nullable=True),
    StructField("y_al", DoubleType(), nullable=True)
])

# 添加解析后的struct列:成功解析则为struct,失败则为null
df = df.withColumn("parsed_cords", from_json(col("cords"), cords_struct_schema))

# 处理struct类型的记录:展平嵌套字段
struct_df = df.filter(col("parsed_cords").isNotNull()) \
    .withColumn("x_al", col("parsed_cords.x_al")) \
    .withColumn("y_al", col("parsed_cords.y_al")) \
    .drop("parsed_cords")

# 处理非struct类型的记录:保留原cords,展平字段设为null
non_struct_df = df.filter(col("parsed_cords").isNull()) \
    .withColumn("x_al", lit(None).cast(DoubleType())) \
    .withColumn("y_al", lit(None).cast(DoubleType())) \
    .drop("parsed_cords")

# 合并两个结果DataFrame
final_df = struct_df.unionByName(non_struct_df)

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

代码说明

  • from_json是核心函数:它尝试将cords字符串解析为指定的struct schema,只有原始为JSON对象的记录能解析成功,其他类型(比如数值、普通字符串)会返回null;
  • 用filter分离两种记录后分别处理,逻辑更清晰;
  • unionByName确保合并时列名完全匹配,避免因列顺序不同导致的错误。

如果你的非struct类型记录需要特殊处理(比如把cords数值转成特定格式),可以在non_struct_df的处理逻辑中添加对应的转换操作。

内容的提问来源于stack exchange,提问作者Palaksha K S

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.01 02:48:12