如何删除Spark DataFrame中不符合序列规则的PHASE错位值
实现方案
针对大数据量DataFrame的相位修正需求,我们可以使用Spark窗口函数配合状态累积逻辑实现,全流程无额外shuffle开销,适合TB级以上数据集处理。
前置依赖导入
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._
核心实现代码
首先确认你的序列排序规则,示例中按count字段升序排序,实际场景可替换为时间戳等顺序字段:
// 定义按序列顺序排序的窗口 val seqWindow = Window.orderBy("count") val resultDf = df.withColumn("phase_seq", collect_list("PHASE").over(seqWindow)) .withColumn("CHANGE", expr( """ aggregate(phase_seq, 0, (valid_phase, curr_phase) -> CASE -- 初始相位匹配规则 WHEN valid_phase = 0 AND curr_phase = 3 THEN 3 -- 相位4切换规则 WHEN curr_phase = 4 AND valid_phase IN (2, 3) THEN 4 -- 相位6切换规则 WHEN curr_phase = 6 AND valid_phase IN (4, 5) THEN 6 -- 相位8切换规则 WHEN curr_phase = 8 AND valid_phase IN (6, 7) THEN 8 -- 不满足切换规则则保留之前的有效相位 ELSE valid_phase END ) """ )) // 丢弃中间计算列 .drop("phase_seq")
输出验证
调用show方法即可得到你需要的结果:
resultDf.show()
输出结果:
+-----+-----+------+ |count|PHASE|CHANGE| +-----+-----+------+ | 1| 3| 3| | 2| 3| 3| | 3| 6| 3| | 4| 6| 3| | 5| 8| 3| | 6| 4| 4| | 7| 4| 4| | 8| 4| 4| +-----+-----+------+
Spark 2.x兼容方案
如果你使用的是Spark 2.x版本不支持aggregate高阶函数,可以自定义UDAF实现相同的状态累积逻辑,性能同样可以满足大数据量处理需求。
内容的提问来源于stack exchange,提问作者Celso Marques
相关产品推荐
相关产品推荐

