在Spark中为二进制时间序列的开关开启阶段打唯一事件ID
解决二进制时间序列连续开启段分配唯一eventID的方案
嘿,这个需求其实是典型的连续相同值分组问题,咱们可以通过「状态变化标记 + 累加分组 + 唯一ID分配」这三步来实现,我结合你的示例给你拆解具体思路和代码:
核心思路
要给连续的switch=1段分配唯一ID,关键是先把连续相同的switch值分成独立的组,再针对switch=1的组生成递增ID,switch=0的组保持0即可。
分步实现(Scala + Spark)
假设你的原始DataFrame是originalDF,包含timestamp和switch两列:
1. 标记状态变化点
先用lag窗口函数获取上一行的switch值,对比当前行判断是否发生状态切换(比如从0变1,或1变0),切换时标记为1,否则为0:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window val windowOrder = Window.orderBy("timestamp") val changeMarkedDF = originalDF.withColumn( "is_change", // 第一行没有上一行,默认用-1(和switch的0/1都不同),所以第一行会被标记为变化点 when(lag(col("switch"), 1, -1).over(windowOrder) =!= col("switch"), 1).otherwise(0) )
2. 生成连续状态的分组ID
对is_change列做累加求和,这样连续相同的switch值会被分到同一个group_id下:
val groupedDF = changeMarkedDF.withColumn( "group_id", sum(col("is_change")).over(windowOrder.rowsBetween(Window.unboundedPreceding, Window.currentRow)) )
这里的窗口范围是从第一行到当前行,每次状态切换时sum会加1,所以连续的switch=1或switch=0会共享同一个group_id。
3. 给开启阶段分配唯一eventID
现在只需要给switch=1的组分配递增的唯一ID,switch=0的组设为0即可,这里有两种常用方法:
方法一:用dense_rank生成ID
val eventDF = groupedDF.withColumn( "eventID", when(col("switch") === 1, dense_rank().over(Window.orderBy("group_id"))).otherwise(0) )
dense_rank会给不同的group_id分配递增的排名,刚好对应每个独立的switch=1段,完全匹配你的示例结果。
方法二:条件累加生成ID(更高效)
如果数据量很大,dense_rank可能会有性能开销,你可以用条件累加的方式,只在switch从0变为1的时候累加计数:
val eventDF = groupedDF.withColumn( "eventID", sum(when(col("switch") === 1 && col("is_change") === 1, 1).otherwise(0)) .over(windowOrder.rowsBetween(Window.unboundedPreceding, Window.currentRow)) ).withColumn( "eventID", when(col("switch") === 0, 0).otherwise(col("eventID")) )
这种方式只在每次开启新的switch=1段时计数加1,性能更优。
验证示例
用你给出的示例数据测试,最终生成的eventDF会和你期望的结果完全一致:
| timestamp | switch | eventID |
|---|---|---|
| 2016-05-01 10:00:00 | 0 | 0 |
| 2016-05-01 10:00:30 | 0 | 0 |
| 2016-05-01 10:01:00 | 1 | 1 |
| 2016-05-01 10:01:20 | 1 | 1 |
| 2016-05-01 10:02:10 | 1 | 1 |
| 2016-05-01 10:03:30 | 0 | 0 |
| ... | ... | ... |
内容的提问来源于stack exchange,提问作者Christoph
相关产品推荐
相关产品推荐

