Spark如何转置展开带动态嵌套数组的列并解决类型不匹配问题
问题解决方案
错误原因
stack函数要求所有待堆叠的列数据类型完全一致,你的a、b列数组内的结构体仅包含date、val两个字段,c列的结构体多了val_dynamic字段,类型不匹配触发了AnalysisException。
解决思路
统一所有待转置数组列的结构体Schema,给缺失val_dynamic字段的列手动补充该字段,值设为null即可。
完整可运行代码
import pyspark.sql.functions as f from pyspark.sql.types import StructType, StructField, LongType df = spark.read.json(sc.parallelize([ """{"id":1,"a":[{"date":1,"val":1},{"date":11,"val":11}]}""", """{"id":2,"b":[{"date":2,"val":2}]}}""", """{"id":3,"c":[{"date":3,"val":3, "val_dynamic":3}]}}""" ])) # 定义统一的结构体Schema,包含所有可能出现的字段 common_struct = StructType([ StructField("date", LongType(), nullable=True), StructField("val", LongType(), nullable=True), StructField("val_dynamic", LongType(), nullable=True) ]) cols = ['a', 'b', 'c'] # 对每个列统一Schema,缺失字段自动补null for c in cols: df = df.withColumn(c, f.col(c).cast(f"array<{common_struct.simpleString()}>")) # 构造stack表达式 expr = f"stack({len(cols)}," + \ ",".join([f"'{c}',{c}" for c in cols]) + \ ")" # 目标输出1:仅转置 transpose_df = df.selectExpr("id", expr) \ .withColumnRenamed("col0", "cols") \ .withColumnRenamed("col1", "arrays") \ .filter("arrays is not null") transpose_df.show() # 目标输出2:转置+展开 explode_df = transpose_df.selectExpr('id', 'cols', 'inline(arrays)') explode_df.show()
运行结果
目标输出1(transpose_df)
+---+----+------------------------+ | id|cols| arrays| +---+----+------------------------+ | 1| a|[{1, 1, null}, {11, 1...| | 2| b| [{2, 2, null}] | | 3| c| [{3, 3, 3}] | +---+----+------------------------+
目标输出2(explode_df)
+---+----+----+---+-----------+ | id|cols|date|val|val_dynamic| +---+----+----+---+-----------+ | 1| a| 1| 1| null| | 1| a| 11| 11| null| | 2| b| 2| 2| null| | 3| c| 3| 3| 3| +---+----+----+---+-----------+
内容的提问来源于stack exchange,提问作者Dan
相关产品推荐
相关产品推荐

