在PySpark中修改嵌套Struct、Array内字段值(Spark 2.4.8)
问题:修改Spark DataFrame嵌套数组结构体字段值
数据Schema
root |-- id: string (nullable = true) |-- elements: struct (nullable = true) | |-- created: string (nullable = true) | |-- id: string (nullable = true) | |-- items: array (nullable = true) | | |-- element: struct (containsNull = true) | | | |-- field: string (nullable = true) | | | |-- fieldId: string (nullable = true) | | | |-- fieldtype: string (nullable = true) | | | |-- from: string (nullable = true) | | | |-- fromString: string (nullable = true) | | | |-- tmpFromAccountId: string (nullable = true) | | | |-- tmpToAccountId: string (nullable = true) | | | |-- to: string (nullable = true) | | | |-- toString: string (nullable = true)
需求说明
需要将elements.items数组内每个element结构体中的指定字段值统一改为"Issue",无论原字段是否为空。
修改前数据示例
+--------+--------------------------------------------------------------------------------+ | id | elements | +--------+--------------------------------------------------------------------------------+ |ABCD-123|[2023-01-16T20:25:30.875+0700, 5388402, [[field, , status,,,,, 23456, Yes]]] | +--------+--------------------------------------------------------------------------------+
修改后数据示例
+--------+----------------------------------------------------------------------------------------------------------+ | id | elements | +--------+----------------------------------------------------------------------------------------------------------+ |ABCD-123|[2023-01-16T20:25:30.875+0700, 5388402, [[Issue, Issue, Issue, Issue, Issue, Issue, Issue, Issue, Issue]]]| +-------------------------------------------------------------------------------------------------------------------+
尝试过的无效代码
replace_list = ['field', 'fieldtype', 'fieldId', 'from', 'fromString', 'to', 'toString', 'tmpFromAccountId', 'tmpToAccountId'] # Didn't work 1 for col_name in replace_list: df = df.withColumn(f"items.element.{col_name}", lit("Issue")) # Didn't work 2 for col_name in replace_list: df = df.withColumn("elements.items.element", struct(col(f"elements.items.element.*"), lit("Issue").alias(f"{col_name}")))
解决方案(Spark 2.4.8 适用,无需explode)
Spark 2.4.8支持高阶函数transform,可以直接遍历数组并修改每个元素的结构体字段,无需拆分数组。核心思路是重新构建嵌套结构:保留elements中的created和id字段不变,对items数组使用transform遍历每个element,重新构造所有指定字段为"Issue"的结构体。
代码实现
from pyspark.sql import functions as F replace_list = ['field', 'fieldtype', 'fieldId', 'from', 'fromString', 'to', 'toString', 'tmpFromAccountId', 'tmpToAccountId'] # 构造transform表达式:遍历items数组,将每个element的指定字段替换为"Issue" transform_expr = f"transform(elements.items, x -> struct({', '.join([f'\'Issue\' as {col}' for col in replace_list])}))" # 重新构建整个elements结构体 df = df.withColumn( "elements", F.struct( F.col("elements.created"), F.col("elements.id"), F.expr(transform_expr).alias("items") ) ) # 查看结果 df.show(truncate=False)
为什么之前的代码无效
- 第一种方法:Spark不支持直接通过点路径(
items.element.{col_name})修改数组内的元素字段,必须遍历数组处理每个元素。 - 第二种方法:尝试直接修改
elements.items.element,但items是数组类型,无法直接整体修改所有元素的结构体,必须通过遍历操作实现。
内容的提问来源于stack exchange,提问作者user20990375
相关产品推荐
相关产品推荐

