如何从StructType类型的Array中删除元素并提取生成新列?
实现方案
以下基于Spark 3.x版本实现,Scala/PySpark/Spark SQL逻辑通用。
1. 实现删除数组中key为"electronic"的元素
直接使用Spark内置数组高阶函数filter过滤即可:
PySpark代码
from pyspark.sql import functions as F df = df.withColumn("item", F.expr("filter(item, x -> x.key != 'electronic')"))
Spark SQL代码
select *, filter(item, x -> x.key != 'electronic') as item from 你的表名
2. 实现每个key生成独立列,type为"one"时取one字段值
分两种场景选择实现方式:
场景A:已知所有key的枚举值,性能更优
如果key的取值是固定可枚举的,直接按key提取:
PySpark代码
# 替换为你实际的所有key值列表 all_keys = ["book", "clothes", "digital", "food"] for key in all_keys: df = df.withColumn( key, # 匹配key和type,取第一个符合条件的元素的one值,若不存在则为null F.expr(f"element_at(filter(item, x -> x.key = '{key}' and x.type = 'one'), 1).one") )
Spark SQL代码
select *, element_at(filter(item, x -> x.key = 'book' and x.type = 'one'), 1).one as book, element_at(filter(item, x -> x.key = 'clothes' and x.type = 'one'), 1).one as clothes -- 其他key按相同格式补全即可 from 你的表名
场景B:key取值不固定,动态生成列
如果key会动态变化,没有固定枚举范围,用炸开+透视的方式实现:
PySpark代码
from pyspark.sql import functions as F # 先炸开数组,提取需要的字段 df_explode = df.select( *[F.col(col) for col in df.columns if col != "item"], F.explode("item").alias("ele") ).select( "*", F.col("ele.key").alias("key"), F.when(F.col("ele.type") == "one", F.col("ele.one")).alias("col_val") ) # 透视key为列,保留原有所有维度 df_final = df_explode.groupBy([col for col in df.columns if col != "item"]).pivot("key").agg(F.first("col_val"))
内容的提问来源于stack exchange,提问作者boom_clap
相关产品推荐
相关产品推荐

