Pyspark 2.4:无需UDF,如何从结构体数组生成新结构体数组?
不用UDF实现PySpark数据结构转换的两种方法
当然可以不用UDF搞定这个需求!在PySpark 2.4里,我们有两种更高效、更贴合Spark原生生态的方式来实现从df_1到df_2的转换,下面分别介绍:
方法一:使用DataFrame API结合Spark SQL的transform函数
这是最推荐的方式,因为Spark的内置函数会经过Catalyst优化,性能比UDF好很多,代码也简洁。核心思路是用transform遍历request数组,重新构造每个address结构体只保留street字段:
from pyspark.sql.functions import expr # 基于df_1生成df_2 df_2 = df_1.withColumn( "request", expr("transform(request, addr -> struct(addr.street as street))") )
简单解释下:
transform(request, addr -> ...)会逐个处理数组里的每个address元素(这里用addr指代)struct(addr.street as street)会重新生成一个只包含street字段的结构体,自动替换掉原来带postcode的address结构
方法二:使用RDD的map操作
如果你需要更灵活的自定义逻辑(比如后续要加更多复杂处理),可以把DataFrame转成RDD来处理,再转回DataFrame:
首先我们需要先定义好df_2的目标Schema(因为要去掉postcode字段),然后编写处理每行数据的逻辑:
from pyspark.sql.types import StructType, StructField, StringType, ArrayType from pyspark.sql import Row # 定义df_2的目标Schema target_schema = StructType([ StructField("request", ArrayType( StructType([ StructField("address", StructType([ StructField("street", StringType(), nullable=False) ]), nullable=False) ]), nullable=False )) ]) # 定义处理单行数据的函数 def clean_address(row): # 遍历request数组,只保留每个address里的street字段 cleaned_request = [ Row(address=Row(street=item["address"]["street"])) for item in row.request ] # 返回替换了request字段的新Row return row._replace(request=cleaned_request) # 转换为RDD处理后,再转回DataFrame df_2 = df_1.rdd.map(clean_address).toDF(target_schema)
两种方法的对比
- 方法一(内置函数):性能更优,代码简洁,适合大多数常规场景,推荐优先使用
- 方法二(RDD map):灵活性更高,适合复杂的自定义逻辑,但会脱离Spark的Catalyst优化,大数据量下性能不如内置函数
内容的提问来源于stack exchange,提问作者Softhinker.com
相关产品推荐
相关产品推荐

