PySpark实现按Session分组生成next_timestamp并填充最后一行空值
PySpark实现分组后生成next_timestamp并填充最后一行值
要实现和Pandas中groupby.shift(-1).fillna(df['timestamp'])完全一致的逻辑,在PySpark里可以通过**窗口函数lead结合coalesce**来完成,具体步骤如下:
核心逻辑
- 按
session分组、timestamp排序,用lead函数获取分组内下一行的timestamp,分组最后一行会返回null - 用
coalesce函数将null值替换为当前行的timestamp,完成填充
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import lead, coalesce, col # 初始化Spark会话 spark = SparkSession.builder.appName("session_next_ts").getOrCreate() # 构造示例数据 sample_data = [ ("session1", "2023-10-01 10:00:00"), ("session1", "2023-10-01 10:05:00"), ("session1", "2023-10-01 10:10:00"), ("session2", "2023-10-01 11:00:00"), ("session2", "2023-10-01 11:05:00") ] # 创建DataFrame并转换timestamp类型 df = spark.createDataFrame(sample_data, ["session", "timestamp"]) df = df.withColumn("timestamp", col("timestamp").cast("timestamp")) # 定义窗口:按session分组,按timestamp升序排序 session_window = Window.partitionBy("session").orderBy("timestamp") # 生成next_timestamp列 result_df = df.withColumn( "next_timestamp", # 先取下一行的timestamp,空值则用当前行的timestamp填充 coalesce(lead(col("timestamp"), 1).over(session_window), col("timestamp")) ) # 查看结果 result_df.show(truncate=False)
输出结果
+--------+-------------------+-------------------+ |session |timestamp |next_timestamp | +--------+-------------------+-------------------+ |session1|2023-10-01 10:00:00|2023-10-01 10:05:00| |session1|2023-10-01 10:05:00|2023-10-01 10:10:00| |session1|2023-10-01 10:10:00|2023-10-01 10:10:00| |session2|2023-10-01 11:00:00|2023-10-01 11:05:00| |session2|2023-10-01 11:05:00|2023-10-01 11:05:00| +--------+-------------------+-------------------+
对应Pandas逻辑对比
你的Pandas参考代码逻辑如下:
import pandas as pd df_pd = pd.DataFrame(sample_data, columns=["session", "timestamp"]) df_pd["timestamp"] = pd.to_datetime(df_pd["timestamp"]) df_pd["next_timestamp"] = df_pd.groupby("session")["timestamp"].shift(-1) df_pd["next_timestamp"] = df_pd["next_timestamp"].fillna(df_pd["timestamp"])
PySpark的实现完全对齐了这一逻辑:lead(1)对应shift(-1),coalesce(..., col("timestamp"))对应fillna(df_pd["timestamp"])。
内容的提问来源于stack exchange,提问作者user20480394
相关产品推荐
相关产品推荐

