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
相关产品推荐
相关产品推荐

