PySpark如何单次select链式调用explode并选取struct字段
问题描述
现有如下结构的PySpark DataFrame,col_name字段为存储结构体的数组类型:
from pyspark.sql import functions as F df = spark.createDataFrame([([(1, 2), (3, 4)],)], 'col_name array<struct<c1:int,c2:int>>') df.show() # +----------------+ # | col_name| # +----------------+ # |[{1, 2}, {3, 4}]| # +----------------+ df.printSchema() # root # |-- col_name: array (nullable = true) # | |-- element: struct (containsNull = true) # | | |-- c1: integer (nullable = true) # | | |-- c2: integer (nullable = true)
常规实现需要两次select操作:先调用explode展开数组得到struct类型列,再选取struct内的字段,代码如下:
df = df.select( F.explode('col_name') ).select( [f'col.{c}' for c in ('c1', 'c2')] )
执行结果符合预期:
df.show() # +---+---+ # | c1| c2| # +---+---+ # | 1| 2| # | 3| 4| # +---+---+ df.printSchema() # root # |-- c1: integer (nullable = true) # |-- c2: integer (nullable = true)
尝试直接在单次select中对explode结果按下标取struct字段时抛出异常:
df = df.select( [F.explode('col_name')[c] for c in ('c1', 'c2')] )
报错信息:
AnalysisException: No such struct field c1 in col
实现方案
可以通过内置函数inline实现单次select完成需求,这也是性能最优的写法。inline函数的作用就是直接展开array<struct>类型的列:数组中每个元素生成一行,struct内的每个字段自动拆分为独立列,不需要二次select操作,代码如下:
df = df.select(F.inline("col_name"))
执行效果和两次select的写法完全一致,字段名、类型、数据都和预期输出完全匹配。
之前写法报错的原因是:直接对F.explode()返回的生成列调用下标取字段时,Spark还未将该列纳入查询计划的列元数据中,无法识别到返回struct的内部字段,因此抛出字段不存在的异常。
如果一定要使用explode实现单次select,可以通过selectExpr在SQL表达式中完成字段提取,但这种写法会重复执行explode逻辑,性能远低于inline方案,不推荐使用:
df = df.selectExpr( "explode(col_name).c1 as c1", "explode(col_name).c2 as c2" )
内容的提问来源于stack exchange,提问作者ZygD
相关产品推荐
相关产品推荐

