Spark Scala 移除结构体数组指定日期元素的最优实现方法
Spark Scala 过滤移除结构体数组列中指定条件元素
需求描述
- 处理目标:Spark DataFrame 中名为
history的数组列,数组元素为结构体类型 - 过滤规则:移除数组内所有结构体中
infodate字段值等于2019-01-26的完整元素 - 要求:仅通过
infodate字段做判断,无需硬编码结构体其他字段的取值
测试数据集构造
import org.apache.spark.sql.types._ val df=Seq( ("123","Jack",Seq(("2020-04-26","200","72","ABC"),("2020-05-26","300","71","ABC"),("2019-01-26","200","70","DEF"),("2019-01-26","200","70","DEF"),("2019-01-26","200","70","DEF"))), ("124","jones",Seq(("2020-04-26","200","72","ABC"),("2020-05-26","300","71","ABC"),("2020-06-26","200","70","ABC"),("2020-08-26","300","69","ABC"),("2020-08-26","300","69","ABC"))), ("125","daniel",Seq(("2019-01-26","200","70","DEF"),("2019-01-26","200","70","DEF"),("2019-01-26","200","70","DEF"),("2019-01-26","200","70","DEF"),("2019-01-26","200","70","DEF"))) ).toDF("id","name","history") .withColumn("history",$"history".cast("array<struct<infodate:Date,amount1:Integer,amount2:Integer,detail:string>>")) // 打印schema验证结构 df.printSchema
对应schema输出:
root |-- id: string (nullable = true) |-- name: string (nullable = true) |-- history: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- infodate: date (nullable = true) | | |-- amount1: integer (nullable = true) | | |-- amount2: integer (nullable = true) | | |-- detail: string (nullable = true)
原有硬编码方案问题
原有实现使用array_except做全结构体值匹配,需要手动写死待删除元素的所有字段值,只要待删除元素的其他字段和硬编码值不完全一致就无法命中删除逻辑,灵活性极差,代码如下:
val dfnew=df .withColumn( "history" , array_except( col("history"), array( struct( lit("2019-01-26").cast(DataTypes.DateType).alias("infodate"), lit("200").cast(DataTypes.IntegerType).alias("amount1"), lit("70").cast(DataTypes.IntegerType).alias("amount2"), lit("DEF").alias("detail") ) ) ) )
推荐实现方案
使用Spark 2.4版本开始内置的数组高阶函数filter,直接遍历数组内每个元素做条件判断,仅保留符合要求的元素,不需要匹配结构体全字段,无UDF性能损耗,代码简洁易维护:
import org.apache.spark.sql.functions._ // 定义待删除的目标日期 val deleteDate = lit("2019-01-26").cast(DateType) val dfRes = df.withColumn( "history", // 遍历history数组的每个元素elem,仅保留infodate不等于目标日期的元素 filter(col("history"), elem => elem.getField("infodate") =!= deleteDate) )
处理效果说明
- id=123的行:原数组中3条
infodate=2019-01-26的记录全部移除,剩余2条2020年的有效记录 - id=124的行:数组中无符合删除条件的元素,保留全部5条原始记录
- id=125的行:数组所有元素均符合删除条件,处理后
history列为空数组
该方案扩展成本极低,如果后续需要叠加过滤条件,比如同时过滤掉
detail为DEF的元素,只需要在判断逻辑中加&& elem.getField("detail") =!= "DEF"即可,不需要修改整体代码结构。
内容的提问来源于stack exchange,提问作者gaurav mathur
相关产品推荐
相关产品推荐

