Pyspark将Struct结构体转数组 解析Azure Data Factory JSON文件
解决方案
1. 解决JSON读取时的损坏记录问题
出现损坏记录的核心原因是Spark默认推断Schema时,只会基于首条记录的properties.parameters结构体生成固定字段的Schema,后续记录参数名和首条不一致时就会被判定为格式错误。你可以直接手动指定Schema,将parameters字段声明为Map类型适配动态参数名,无需依赖自动推断:
from pyspark.sql import SparkSession from pyspark.sql.types import * spark = SparkSession.builder.appName("ADFJsonParse").getOrCreate() # 手动定义Schema,parameters设为键为参数名、值为参数属性的Map结构 adf_schema = StructType([ StructField("name", StringType(), nullable=False), StructField("properties", StructType([ StructField("annotations", ArrayType(StringType()), nullable=True), StructField("parameters", MapType( StringType(), StructType([ StructField("type", StringType(), nullable=True), StructField("defaultValue", StringType(), nullable=True) ]) ), nullable=True) ]), nullable=False) ]) # 按指定Schema读取JSON文件 raw_df = spark.read.schema(adf_schema).json("/path/to/your/adf_json_files/")
如果使用Scala开发,逻辑完全一致,仅需调整Schema定义和API的对应语法即可。
2. 生成注释表
用内置explode函数炸开注释数组即可:
from pyspark.sql.functions import explode, col, explode_outer # 不需要保留无注释的pipeline用explode,需要保留则替换为explode_outer annotation_df = raw_df.select( col("name").alias("pipeline_name"), explode(col("properties.annotations")).alias("annotation") )
3. 生成参数表
直接炸开parameters的Map结构,每行对应一个参数,提取对应属性即可:
# 不需要保留无参数的pipeline用explode,需要保留则替换为explode_outer parameter_df = raw_df.select( col("name").alias("pipeline_name"), explode(col("properties.parameters")).alias("parameter_name", "parameter_attr") ).select( "pipeline_name", "parameter_name", col("parameter_attr.type").alias("parameter_type"), col("parameter_attr.defaultValue").alias("parameter_default") )
报错说明
你之前自定义UDF报错大概率是因为读取阶段生成的parameters字段本身已是损坏结构,或是UDF返回值类型定义和实际输出不匹配。上述方案全程使用Spark内置函数实现,无需自定义UDF,执行性能更优,也不会出现类型转换类错误。
内容的提问来源于stack exchange,提问作者Alfred G
相关产品推荐
相关产品推荐

