如何修改PySpark DataFrame中数组结构体类型列的long_value值
解决PySpark数组嵌套结构体字段修改问题
你的代码存在两个核心问题:
- 仅修改了
properties数组的第一个元素(col("properties")[0]),未处理数组内所有元素 - 直接用
"properties.value.long_value"作为新列名是错误的逻辑,你需要替换整个properties列,而非创建一个不存在的嵌套列
正确实现代码
要批量修改数组中所有结构体的long_value,需用transform函数遍历数组元素,重构结构体时保留其他字段并修改目标值:
from pyspark.sql import functions as F df = df.withColumn( "properties", F.transform( "properties", lambda elem: F.struct( elem["key"].alias("key"), F.struct( elem["value"]["string_value"].alias("string_value"), F.when(elem["value"]["long_value"].isNotNull(), elem["value"]["long_value"] / 10).alias("long_value") ).alias("value") ) ) )
代码说明
transform("properties", lambda elem: ...):遍历properties数组的每个元素elem- 用
F.struct重构外层结构体,保留原key字段 - 内层重构
value结构体:- 保留原
string_value字段 - 用
F.when判断long_value不为null时才执行除法,避免null值运算报错
- 保留原
- 将修改后的数组重新赋值给
properties列,覆盖原列数据
处理value为null的场景
如果value本身可能为null,需增加一层判断避免报错:
df = df.withColumn( "properties", F.transform( "properties", lambda elem: F.struct( elem["key"].alias("key"), F.when( elem["value"].isNotNull(), F.struct( elem["value"]["string_value"].alias("string_value"), F.when(elem["value"]["long_value"].isNotNull(), elem["value"]["long_value"] / 10).alias("long_value") ) ).alias("value") ) ) )
内容的提问来源于stack exchange,提问作者Ktos
相关产品推荐
相关产品推荐

