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

Spark如何按Id分组提取Msg_name列start紧随end的时序对模式

问题根因

原有代码仅判断相邻行的Msg_name不一致,未限定合法配对的规则:必须是前一行为start、当前行为end,导致end→start这类无效配对也被保留,输出不符合预期。

正确实现代码
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

val window = Window.partitionBy("Id_num").orderBy("TimeStamp")

val result = latenies_df
  // 窗口取前一行的消息类型和时间戳
  .withColumn("prev_msg", lag("Msg_name", 1).over(window))
  .withColumn("prev_time", lag("TimeStamp", 1).over(window))
  // 仅保留start紧跟end的合法配对
  .filter($"Msg_name" === "end" && $"prev_msg" === "start")
  // 按要求输出字段
  .select(
    lit("pattern1").alias("pattern_name"),
    $"prev_time".alias("Timestamp_start"),
    $"TimeStamp".alias("Timestamp_end"),
    $"Id_num"
  )
  .orderBy("Id_num", "Timestamp_start")

result.show(false)
输出结果
+------------+--------------------------+--------------------------+------+
|pattern_name|Timestamp_start           |Timestamp_end             |Id_num|
+------------+--------------------------+--------------------------+------+
|pattern1    |2021-10-25 21:26:23.48667 |2021-10-25 21:26:23.48759 |304875|
|pattern1    |2021-10-25 21:26:23.509683|2021-10-25 21:26:23.509689|304875|
|pattern1    |2021-10-25 21:26:23.510156|2021-10-25 21:26:23.510159|201345|
|pattern1    |2021-10-25 21:26:23.510335|2021-10-25 21:26:23.513177|201345|
+------------+--------------------------+--------------------------+------+

完全符合预期要求:连续多个start仅保留最后一个与后续end配对,连续多个end仅保留第一个与前序start配对,自动过滤无效配对。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 21:57:02