PySpark窗口函数无法实现?求时序数据前后关联转换方案
PySpark实现按分组关联下一条记录的解决方案
你遇到的这个需求完全可以用PySpark的窗口函数lead()来实现,之前你觉得需要聚合可能是对窗口函数的用法有误解——lead()属于分析函数,不需要分组聚合,直接在分区内获取下一条记录的值即可。
具体实现步骤:
- 导入必要模块并创建测试数据(已有DataFrame可跳过此步)
- 定义窗口规则:按
Id分区,按TimeStamp排序 - 使用
lead()获取下一条记录的时间和数值 - 重命名列并过滤无效行
代码示例:
# 导入依赖 from pyspark.sql import Window, functions as F # 创建测试DataFrame data = [ (1, "01/01/2023 10:15", 10), (1, "01/01/2023 10:30", 20), (1, "01/01/2023 10:45", 40), (2, "01/01/2023 10:15", 15), (2, "01/01/2023 10:30", 25), (2, "01/01/2023 10:45", 35) ] df = spark.createDataFrame(data, ["Id", "TimeStamp", "value"]) # 可选:将字符串类型的TimeStamp转为时间戳类型,确保排序准确 df = df.withColumn("TimeStamp", F.to_timestamp("TimeStamp", "MM/dd/yyyy HH:mm")) # 定义窗口:按Id分组,按TimeStamp升序排列 window_spec = Window.partitionBy("Id").orderBy("TimeStamp") # 生成目标DataFrame result_df = df.withColumn("EndTimeStamp", F.lead("TimeStamp").over(window_spec)) \ .withColumn("End Reading", F.lead("value").over(window_spec)) \ .withColumnRenamed("TimeStamp", "StartTimeStamp") \ .withColumnRenamed("value", "Starting Reading") \ .filter(F.col("EndTimeStamp").isNotNull()) \ .select("Id", "StartTimeStamp", "Starting Reading", "EndTimeStamp", "End Reading") # 查看结果 result_df.show(truncate=False)
输出结果:
+---+-------------------+---------------+-------------------+-----------+ |Id |StartTimeStamp |Starting Reading|EndTimeStamp |End Reading| +---+-------------------+---------------+-------------------+-----------+ |1 |2023-01-01 10:15:00|10 |2023-01-01 10:30:00|20 | |1 |2023-01-01 10:30:00|20 |2023-01-01 10:45:00|40 | |2 |2023-01-01 10:15:00|15 |2023-01-01 10:30:00|25 | |2 |2023-01-01 10:30:00|25 |2023-01-01 10:45:00|35 | +---+-------------------+---------------+-------------------+-----------+
关键说明:
lead()函数的作用是在指定窗口分区内,获取当前行之后第N条记录的值(默认N=1),属于分析函数范畴——它不会合并行,只是为每行补充后续行的信息,完全不需要使用groupBy聚合操作。你之前的困惑可能是混淆了分析函数与聚合函数的差异,聚合函数会将分区内的多行合并为一行,而分析函数是对每行独立计算并保留原有行结构。
内容的提问来源于stack exchange,提问作者Krishna
相关产品推荐
相关产品推荐

