Spark如何为DataFrame多列复用已定义全字段的Array结构体Schema
需求说明
我有上百个列a、b、c……,需要修改DataFrame的Schema,让每个数组类型字段的元素结构体都具备相同的结构,包含date、num和val三个字段。由于存在数千个id,我只希望修改Schema而非改动DataFrame本身,修改后的Schema将用于后续高效加载数据到DataFrame,避免使用UDF修改全量DataFrame。
输入Schema(执行df.printSchema()输出)
root |-- a: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- date: long (nullable = true) | | |-- num: long (nullable = true) !!! 注意 : `num` 字段存在 !!! | | |-- val: long (nullable = true) |-- b: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- date: long (nullable = true) | | |-- val: long (nullable = true) |-- c: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- date: long (nullable = true) | | |-- val: long (nullable = true) |-- d: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- date: long (nullable = true) | | |-- val: long (nullable = true) |-- id: long (nullable = true)
期望输出Schema
root |-- a: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- date: long (nullable = true) | | |-- num: long (nullable = true) | | |-- val: long (nullable = true) |-- b: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- date: long (nullable = true) | | |-- num: long (nullable = true) | | |-- val: long (nullable = true) |-- c: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- date: long (nullable = true) | | |-- num: long (nullable = true) | | |-- val: long (nullable = true) |-- d: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- date: long (nullable = true) | | |-- num: long (nullable = true) | | |-- val: long (nullable = true) |-- id: long (nullable = true)
输入Schema复现代码
df = spark.read.json(sc.parallelize([ """{"id":1,"a":[{"date":2001,"num":1},{"date":2002,},{"date":2003,}],"b":[{"date":2001,"val":4},{"date":2002,"val":5},{"date":2003,"val":6}],"d":[{"date":2001,"val":21},{"date":2002,"val":22},{"date":2003,"val":23}],"c":[{"date":2000,"val":30},{"date":2001,"val":31},{"date":2002,"val":32},{"date":2003,"val":33}]}""", """{"id":2,"a":[{"date":2001,"num":2},{"date":2002},{"date":2003}],"b":[{"date":2001,"val":4},{"date":2002,"val":5},{"date":2003,"val":6}],"d":[{"date":2001,"val":21},{"date":2002,"val":22},{"date":2003,"val":23}],"c":[{"date":1990,"val":39},{"date":2000,"val":30},{"date":2001,"val":31},{"date":2002,"val":32},{"date":2003,"val":33},{"date":2004,"val":34}]}}""" ])) for field in df.schema: print(field)
打印输出
StructField(a,ArrayType(StructType(List(StructField(date,LongType,true),StructField(num,LongType,true),StructField(val,LongType,true))),true),true) StructField(b,ArrayType(StructType(List(StructField(date,LongType,true),StructField(val,LongType,true))),true),true) StructField(c,ArrayType(StructType(List(StructField(date,LongType,true),StructField(val,LongType,true))),true),true) StructField(d,ArrayType(StructType(List(StructField(date,LongType,true),StructField(val,LongType,true))),true),true) StructField(id,LongType,true)
解决方案
直接在数据读取阶段指定自定义Schema即可,无需修改原有数据,也不需要调用UDF处理全量数据,实现代码如下:
from pyspark.sql.types import StructField, StructType, LongType, ArrayType jsonstr=[ """{"id":1,"a":[{"date":2001,"val":1,"num":1},{"date":2002,"val":2},{"date":2003,"val":3}],"b":[{"date":2001,"val":4},{"date":2002,"val":5},{"date":2003,"val":6}],"d":[{"date":2001,"val":21},{"date":2002,"val":22},{"date":2003,"val":23}],"c":[{"date":2000,"val":30},{"date":2001,"val":31},{"date":2002,"val":32},{"date":2003,"val":33}]}""", """{"id":2,"a":[{"date":2001,"val":1},{"date":2002,"val":2},{"date":2003,"val":3}],"b":[{"date":2001,"val":4},{"date":2002,"val":5},{"date":2003,"val":6}],"d":[{"date":2001,"val":21},{"date":2002,"val":22},{"date":2003,"val":23}],"c":[{"date":1990,"val":39},{"date":2000,"val":30},{"date":2001,"val":31},{"date":2002,"val":32},{"date":2003,"val":33},{"date":2004,"val":34}]}}""" ] # 定义统一的数组元素结构 array_schema = ArrayType(StructType([ StructField('date' ,LongType(),True), StructField('num' ,LongType(),True), StructField('val' ,LongType(),True)]), True) # 构造完整的DataFrame Schema keys = ['a', 'b', 'c', 'd'] fields = [StructField(k, array_schema, True) for k in keys] fields.append(StructField('id',LongType(),True)) df_schema = StructType(fields) # 读取数据时指定Schema,自动补全缺失的num字段为null dff = spark.read.json(sc.parallelize(jsonstr),df_schema)
内容的提问来源于stack exchange,提问作者Dan
相关产品推荐
相关产品推荐

