You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 17:27:47