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

如何删除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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 22:57:04