PySpark如何修改Struct列内部嵌套属性 无需重构结构体实现过滤
问题解答
可以实现,Spark 3.1及以上版本内置了StructType的withField方法,支持直接更新结构体内部的指定属性,不需要先提取所有属性、处理后再重新构建整个结构体。
实现代码
PySpark 写法
from pyspark.sql import functions as F # df 为你的原始DataFrame df_result = df.withColumn( "arrayCol", # 直接更新arrayCol结构体的b属性 F.col("arrayCol").withField( "b", # 过滤b数组中不等于1的元素 F.filter(F.col("arrayCol.b"), lambda x: x != 1) ) )
Scala Spark 写法
import org.apache.spark.sql.functions._ val dfResult = df.withColumn( "arrayCol", col("arrayCol").withField( "b", filter(col("arrayCol.b"), _ =!= 1) ) )
低版本兼容方案
如果你使用的Spark版本低于3.1,不支持withField方法,才需要用拆分重构结构体的写法:
from pyspark.sql import functions as F df_result = df.withColumn( "arrayCol", F.struct( F.col("arrayCol.a").alias("a"), F.filter(F.col("arrayCol.b"), lambda x: x != 1).alias("b") ) )
内容的提问来源于stack exchange,提问作者gael
相关产品推荐
相关产品推荐

