如何在PySpark DataFrame中更新嵌套数组内指定字段值(保留原Schema)
PySpark DataFrame嵌套字段值替换方案
需求:将PySpark DataFrame中ticket.money数组内每个struct的total字段值为-9999的条目替换为None,必须保留原Schema且不新增列。
数据Schema
root |-- date: string |-- ticket: struct | |-- money: array | | |-- element: struct | | | |-- currency: string | | | |-- total: double
解决方案代码
from pyspark.sql import functions as F # 假设原DataFrame名为df df_updated = df.withColumn( "ticket", F.struct( F.transform( F.col("ticket.money"), lambda x: F.struct( x["currency"].alias("currency"), F.when(x["total"] == -9999, None).otherwise(x["total"]).alias("total") ) ).alias("money") ) )
代码说明
- 遍历数组元素:使用
transform函数遍历ticket.money数组,对每个数组元素(struct类型)进行处理 - 替换指定值:通过
when条件判断,当total等于-9999时替换为None,否则保留原数值 - 重构嵌套结构:用
struct函数重新构造ticket字段,将处理后的数组赋值给money字段,直接覆盖原ticket列,确保Schema与原表完全一致
验证结果
处理后的DataFrame数据如下:
+-----------------------++----------------------------------------+ |date ||ticket | +-----------------------++----------------------------------------+ |2024-02-02T04:31:06 ||[{CAD, None}, {CAD, None}] | |2024-02-02T04:31:06 ||[{CAD, None}, {CAD, None}] | +-----------------------++----------------------------------------+
内容的提问来源于stack exchange,提问作者Maxime Pelletier Legault
相关产品推荐
相关产品推荐

