You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Scala基于首个DataFrame校验结果创建事件DataFrame

Spark Scala 实现DataFrame双规则校验打标

实现思路

  • 先将字符串类型的date列统一解析为Spark原生日期类型,避免字符串直接比较带来的逻辑错误
  • 用Spark内置日期函数计算当前日期往前推1个月的阈值,不依赖外部时间类,适配分布式执行环境
  • 按规则筛选命中校验逻辑的记录,为不同命中场景赋值对应的events字段值
  • 提供两种命中逻辑适配:单条记录多规则命中取优先级打标、单条记录多规则命中生成多条独立事件

核心代码

首先导入需要依赖的内置函数:

import org.apache.spark.sql.functions._

版本1:单条记录多规则命中时按优先级打单个标签

优先级按需求顺序:日期过老规则 > 合同列为空规则,即同时命中两个规则时,优先取日期规则对应的标签值

// 解析字符串日期为日期类型,若日期字符串不是yyyy-MM-dd格式,需在to_date第二个参数传入对应格式,比如"yyyyMMdd"/"yyyy/MM/dd"
val baseDf = df1.withColumn("parsed_date", to_date(col("date")))
// 计算1个月前的日期阈值
val thresholdDate = add_months(current_date(), -1)

val resultDf = baseDf
  // 筛选命中任意一个校验规则的记录
  .filter(col("parsed_date") < thresholdDate or col("contracts").isNull)
  .withColumn("events",
    when(col("parsed_date") < thresholdDate, "NULL")
    .when(col("contracts").isNull, "OLD_DATE")
  )
  // 删除中间解析用的临时日期列,需要保留可注释该行
  .drop("parsed_date")

注意:当前代码严格按照需求描述赋值事件值,若实际业务逻辑为「合同列为空标记NULL、日期过老标记OLD_DATE」属于需求笔误,直接调换两个when分支内的字符串值即可

版本2:单条记录多规则命中时生成多条独立事件

如果要求只要命中规则就单独生成一条事件记录(同一条原始数据命中两个规则就输出两条),用两个分支筛选后合并即可:

val baseDf = df1.withColumn("parsed_date", to_date(col("date")))
val thresholdDate = add_months(current_date(), -1)

// 筛选日期早于1个月的记录
val oldDateEventDf = baseDf
  .filter(col("parsed_date") < thresholdDate)
  .withColumn("events", lit("NULL"))
  .drop("parsed_date")

// 筛选合同列为空的记录
val nullContractEventDf = baseDf
  .filter(col("contracts").isNull)
  .withColumn("events", lit("OLD_DATE"))
  .drop("parsed_date")

// 合并两个结果集得到最终输出
val resultDf = oldDateEventDf.unionByName(nullContractEventDf)

注意事项

如果原始date列存在格式非法的脏数据,to_date函数会返回null,这类数据默认不会命中日期校验规则。如果需要把日期解析失败的记录也纳入异常范围,可在filter条件中补充col("parsed_date").isNull判断,自定义对应事件标签即可。

内容的提问来源于stack exchange,提问作者stdinstoud

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.29 10:27:29