PySpark实现同Unique_name分组内行匹配与自定义列生成的技术方案咨询
Hey there! 针对你遇到的这个业务场景,我们可以利用PySpark的窗口函数和分组操作来高效处理,彻底摆脱循环遍历带来的性能瓶颈。下面是贴合你需求的实现思路和完整代码:
核心实现思路
我们的核心目标是通过分区窗口对每个Unique_name下的行进行标记、配对和合并,主要分为几个关键步骤:
- 给每个分组内的
new和cleared行分别按时间递增编号,建立配对基础; - 用向前填充的方式把连续的
new行归为同一组,方便后续合并重复项; - 将每组
new行与后续对应的cleared行配对,更新状态字段; - 合并连续的
new行,更新repeatCount和updatetime; - 过滤掉无对应
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
相关产品推荐
相关产品推荐

