基于PySpark窗口函数为多列代码填充首次出现时间戳
PySpark实现跨多列按分组填充代码首次出现时间
核心思路
要解决这个问题,关键是先将多列code转换为长格式,统一计算每个id下每个代码的首次出现时间,再关联回原表,分别映射到对应的start_time列。这种方式可以避免窗口函数只能处理单列的限制,同时不受原表时间排序的影响。
具体实现步骤
1. 导入依赖并模拟数据
from pyspark.sql import SparkSession from pyspark.sql import functions as F # 初始化SparkSession spark = SparkSession.builder.appName("CodeFirstTime").getOrCreate() # 模拟你的数据集(替换为实际表即可) data = [ (1, "2023-08-13", "A", "M", "8"), (1, "2023-05-06", "M", "8", "6"), (2, "2023-06-07", "D", None, "X"), (2, "2023-09-01", "X", "D", None) ] df = spark.createDataFrame(data, ["id", "timestamp", "code1", "code2", "code3"]) df = df.withColumn("timestamp", F.to_timestamp("timestamp"))
2. 将多列Code转成窄表
把code1/code2/code3三列转换为长格式,统一收集所有id、代码值和对应的时间:
unpivot_df = df.select( "id", "timestamp", # 使用stack函数实现宽表转长表 F.expr("stack(3, 'code1', code1, 'code2', code2, 'code3', code3) as (code_col, code)") ).filter(F.col("code").isNotNull()) # 过滤空代码值
3. 计算每个代码的首次出现时间
按id和代码分组,取最小时间作为首次出现时间:
code_first_time = unpivot_df.groupBy("id", "code").agg( F.min("timestamp").alias("first_time") )
4. 关联回原表生成目标列
这里提供两种高效关联方式:
方式一:单Join+Map映射(推荐,性能更优)
将每个id的代码-时间映射转为字典,直接通过键值对查找:
# 生成每个id对应的code到首次时间的映射字典 id_code_map = code_first_time.groupBy("id").agg( F.map_from_entries(F.collect_list(F.struct("code", "first_time"))).alias("code_time_map") ) # 关联后查找对应code的首次时间 result_df = df.join(id_code_map, on="id", how="left") \ .withColumn("start_time_1", F.col("code_time_map")[F.col("code1")]) \ .withColumn("start_time_2", F.col("code_time_map")[F.col("code2")]) \ .withColumn("start_time_3", F.col("code_time_map")[F.col("code3")]) \ .drop("code_time_map")
方式二:多次Join(逻辑更直观)
分别将每个code列与首次时间表关联:
result_df = df \ .join(code_first_time, (df.id == code_first_time.id) & (df.code1 == code_first_time.code), "left") \ .withColumnRenamed("first_time", "start_time_1") \ .join(code_first_time, (df.id == code_first_time.id) & (df.code2 == code_first_time.code), "left") \ .withColumnRenamed("first_time", "start_time_2") \ .join(code_first_time, (df.id == code_first_time.id) & (df.code3 == code_first_time.code), "left") \ .withColumnRenamed("first_time", "start_time_3") \ .select(df["*"], "start_time_1", "start_time_2", "start_time_3")
5. 查看结果
result_df.show(truncate=False)
输出结果将符合你的需求:同一id下,同一个代码无论出现在哪一列,对应的start_time列都会填充它首次出现的时间。
内容的提问来源于stack exchange,提问作者sam1234567890
相关产品推荐
相关产品推荐

