PySpark中如何更新DataFrame结构体数组内的字段值
PySpark修改数组内结构体字段的解决方案
方法1:无需explode(推荐,性能最优)
直接使用PySpark内置的transform函数操作数组内的结构体元素,不需要拆分数组,适配数组内任意数量元素的场景,不会产生额外的shuffle开销:
首先导入依赖函数:
from pyspark.sql.functions import transform, lit, struct
直接对foo字段做修改:
dfUpdated = df.withColumn("foo", transform("foo", lambda elem: struct( lit(10).alias("value"), elem["value2"].alias("value2") )) )
该方法会遍历foo数组中的每个结构体元素,保留原有value2值的同时将value修改为10,修改后的数据schema与原schema完全一致。
方法2:explode后重新组装
如果已经通过explode拆分数组,可按以下步骤组装回原格式:
- 拆分后修改结构体字段
- 用
struct函数重新生成结构体 - 用
collect_list聚合恢复数组格式(多数据场景下需要按原行的唯一标识分组,避免不同行的数组合并)
示例代码:
from pyspark.sql.functions import explode, collect_list, struct, lit, col # 拆分数组 df_exploded = df.select(explode("foo").alias("fooColumn")) # 修改value字段,同时保留原有value2 df_modified = df_exploded.withColumn("fooColumn", struct( lit(10).alias("value"), col("fooColumn.value2").alias("value2") ) ) # 聚合为数组,单条数据场景无需指定分组键,多数据场景需加上原表其他非数组字段作为分组键 dfUpdated = df_modified.groupBy().agg(collect_list("fooColumn").alias("foo"))
结果验证
两种方案执行后均可得到预期结果:
>>> dfUpdated.select(explode('foo').alias("fooColumn")).select('fooColumn.value', 'fooColumn.value2').show() +-----+------+ |value|value2| +-----+------+ | 10| null| +-----+------+
内容的提问来源于stack exchange,提问作者doc
相关产品推荐
相关产品推荐

