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

PySpark实现同Unique_name分组内行匹配与自定义列生成的技术方案咨询

Hey there! 针对你遇到的这个业务场景,我们可以利用PySpark的窗口函数和分组操作来高效处理,彻底摆脱循环遍历带来的性能瓶颈。下面是贴合你需求的实现思路和完整代码:

核心实现思路

我们的核心目标是通过分区窗口对每个Unique_name下的行进行标记、配对和合并,主要分为几个关键步骤:

  1. 给每个分组内的new和cleared行分别按时间递增编号,建立配对基础;
  2. 用向前填充的方式把连续的new行归为同一组,方便后续合并重复项;
  3. 将每组new行与后续对应的cleared行配对,更新状态字段;
  4. 合并连续的new行,更新repeatCount和updatetime;
  5. 过滤掉无对应new行的cleared行,组装最终结果。

完整代码实现

第一步:初始化环境并创建示例DataFrame

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 初始化SparkSession
spark = SparkSession.builder.appName("BusinessScenarioProcessing").getOrCreate()

# 构建输入数据
input_data = [
    ("Unique-01", "2018-09-01", "new", "2018-09-01", None, 0, 0),
    ("Unique-01", "2018-09-15", "cleared", None, "2018-09-15", 0, 0),
    ("Unique-01", "2018-09-16", "new", "2018-09-16", None, 0, 0),
    ("Unique-01", "2018-09-27", "cleared", None, "2018-09-27", 0, 0),
    ("Unique-01", "2018-09-30", "cleared", None, "2018-09-30", 0, 0),
    ("Unique-02", "2018-09-21", "new", "2018-09-21", None, 0, 0),
    ("Unique-02", "2018-09-28", "cleared", None, "2018-09-28", 0, 0),
    ("Unique-02", "2018-09-28", "new", "2018-09-28", None, 0, 0),
    ("Unique-02", "2018-10-05", "new", "2018-10-05", None, 0, 0),
    ("Unique-02", "2018-10-15", "cleared", None, "2018-10-15", 0, 0),
    ("Unique-02", "2018-10-15", "new", "2018-10-15", None, 0, 0)
]

# 定义Schema并创建DataFrame
schema = ["Unique_name", "Input_time", "class", "time_up", "time_down", "repeatCount", "updatetime"]
df = spark.createDataFrame(input_data, schema=schema)

# 将时间字段转换为日期类型,确保排序和比较的准确性
df = df.withColumn("Input_time", F.to_date("Input_time")) \
       .withColumn("time_up", F.to_date("time_up")) \
       .withColumn("time_down", F.to_date("time_down"))

第二步:为new/cleared行标记分组编号

我们给每个Unique_name下的new和cleared行分别按时间顺序编号,方便后续配对:

# 定义分组窗口:按Unique_name分区,Input_time升序排序
window_new_seq = Window.partitionBy("Unique_name").orderBy("Input_time")
window_cleared_seq = Window.partitionBy("Unique_name").orderBy("Input_time")

# 给new行添加序号,cleared行同理
df = df.withColumn("new_seq", F.when(F.col("class") == "new", F.row_number().over(window_new_seq)).otherwise(None))
df = df.withColumn("cleared_seq", F.when(F.col("class") == "cleared", F.row_number().over(window_cleared_seq)).otherwise(None))

第三步:将连续new行归为同一组,并匹配对应的cleared行

用向前填充的方式把连续的new行归为同一组,然后将每组new行与对应的cleared行配对:

# 向前填充new_seq,把连续的new行归为同一个group_id
window_fill_group = Window.partitionBy("Unique_name").orderBy("Input_time").rowsBetween(Window.unboundedPreceding, 0)
df = df.withColumn("group_id", F.last("new_seq", ignorenulls=True).over(window_fill_group))

# 提取cleared行的信息,用于后续配对
cleared_info_df = df.filter(F.col("class") == "cleared").select(
    "Unique_name", "cleared_seq", "time_down", "Input_time"
).withColumnRenamed("Input_time", "cleared_input_time")

# 匹配每组new行对应的cleared行(group_id与cleared_seq相等,且cleared行在new组之后)
df = df.join(
    cleared_info_df,
    (df["Unique_name"] == cleared_info_df["Unique_name"]) 
    & (df["group_id"] == cleared_info_df["cleared_seq"]) 
    & (df["Input_time"] <= cleared_info_df["cleared_input_time"]),
    "left"
).drop(cleared_info_df["Unique_name"])

第四步:合并连续new行,更新repeatCount和updatetime

对同一组内的连续new行,保留第一行并更新重复计数和更新时间:

# 定义组内窗口:按Unique_name和group_id分区,Input_time升序排序
window_group_repeat = Window.partitionBy("Unique_name", "group_id").orderBy("Input_time")

# 计算重复次数:组内第n个new行的repeatCount为n-1
df = df.withColumn(
    "repeatCount",
    F.when(F.col("class") == "new", F.row_number().over(window_group_repeat) - 1).otherwise(F.col("repeatCount"))
)

# 更新updatetime为组内最后一个new行的time_up
df = df.withColumn(
    "updatetime",
    F.when(F.col("class") == "new", F.last("time_up", ignorenulls=True).over(window_group_repeat)).otherwise(F.col("updatetime"))
)

# 同一组内只保留第一个new行(后续重复行的信息已经合并到第一行)
df = df.withColumn("row_in_group", F.row_number().over(window_group_repeat))
df = df.filter((F.col("class") != "new") | (F.col("row_in_group") == 1)).drop("row_in_group")

第五步:组装最终结果并过滤无效行

更新状态字段,过滤掉无对应new行的cleared行:

# 更新class和time_down:如果new行有匹配的cleared行,则标记为cleared并设置time_down
final_df = df.withColumn(
    "class",
    F.when((F.col("class") == "new") & (F.col("time_down").isNotNull()), "cleared").otherwise(F.col("class"))
).withColumn(
    "time_down",
    F.when((F.col("class") == "new") & (F.col("time_down").isNotNull()), F.col("time_down")).otherwise(F.col("time_down"))
)

# 过滤掉无对应new行的cleared行
final_df = final_df.filter(
    ~((F.col("class") == "cleared") & (F.col("group_id").isNull()))
)

# 整理输出字段,确保Input_time为对应new行的原时间
final_df = final_df.select(
    "Unique_name",
    F.when(
        (F.col("class") == "cleared") & (F.col("group_id").isNotNull()),
        F.first("Input_time").over(Window.partitionBy("Unique_name", "group_id"))
    ).otherwise(F.col("Input_time")).alias("Input_time"),
    "class",
    "time_up",
    "time_down",
    "repeatCount",
    "updatetime"
).distinct()

# 按要求排序并展示结果
final_df.orderBy("Unique_name", "Input_time").show()

逻辑匹配说明

  • 配对new与cleared行:通过group_id和cleared_seq的对应关系,精准匹配每组new行后续的cleared行,更新class为cleared并填充time_down;
  • 连续new行处理:利用窗口函数的row_number统计重复次数,last函数获取最新的time_up作为updatetime,只保留每组的第一个new行;
  • 过滤无效cleared行:通过判断group_id是否为空,直接过滤掉没有对应new行的cleared行。

这个方案完全基于PySpark的分布式计算能力,即使数据量达到千万级也能高效运行,彻底解决你之前循环遍历的性能问题。

内容的提问来源于stack exchange,提问作者Albin Chandy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 09:38:12