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

