基于嵌套对象数据类型过滤并展平Azure Synapse PySpark中的嵌套JSON对象
解决Azure Synapse PySpark中混合类型JSON字段的展平问题
我之前也碰到过类似的混合类型JSON字段处理场景,Spark自动推断schema时确实会因为字段类型不一致,统一选用最兼容的类型(比如这里的string),导致我们没法直接对struct类型的记录做展平操作。下面是具体的解决步骤和代码示例:
问题根源
Spark读取JSON数据时,默认会根据所有记录的字段类型推断schema。当同一个字段存在多种类型(比如你的cords既有struct又有数值),Spark会选择一个能容纳所有类型的通用类型(这里就是string),这就直接阻碍了我们对struct类型记录的嵌套字段展平。
解决方案思路
我们可以先把cords当作string类型读取,再通过尝试解析为struct的方式区分两种记录:
- 对能成功解析为struct的记录,展平其嵌套字段;
- 对解析失败的记录,保留原
cords字段,同时将展平后的字段设为null(或者你需要的默认值); - 最后合并两种处理后的结果。
具体代码实现
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
相关产品推荐
相关产品推荐

