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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 13:35:24