如何在PySpark中对关联行分组并匹配客户ID与追踪ID?
如何在PySpark中对关联行分组并匹配客户ID与追踪ID?
你遇到的场景很典型——这种按特定起始标记分组、组内长度不固定的文本数据处理,确实需要用窗口函数来实现分组关联。先把你的数据示例再明确一下:
df = spark.createDataFrame([("new entry", 1, 123), ("acct", 2, None), ("cust ID", 3, None), ("new entry", 4, 456), ("acct", 5, None), ("more text", 6, None), ("cust ID", 7, None)], ("value", "line num", "tracking ID"))
你的需求是把每个从new entry到cust ID的行划分为一个组,然后让组内的cust ID行关联上该组起始行new entry里的tracking ID值。
实现步骤
我们可以通过标记分组起始点 + 累加生成组ID + 组内填充追踪ID这三步来完成:
标记分组起始行
先给每个new entry行打上标记,用来区分分组的开始:from pyspark.sql import functions as F from pyspark.sql.window import Window df_with_start = df.withColumn( "is_group_start", F.when(F.col("value") == "new entry", 1).otherwise(0) )生成分组ID
使用累加窗口函数,基于line num的顺序,把从第一个new entry到下一个new entry之前的所有行归为同一个组:window_group = Window.orderBy("line num").rowsBetween(Window.unboundedPreceding, Window.currentRow) df_with_group = df_with_start.withColumn( "group_id", F.sum("is_group_start").over(window_group) )这里的
sum("is_group_start")会从第一行开始累加,每遇到一个new entry就加1,这样每个组就有了唯一的group_id。组内填充追踪ID
在每个组内,用last函数(忽略null值)把new entry行的tracking ID填充到组内所有行:window_fill = Window.partitionBy("group_id").orderBy("line num").rowsBetween(Window.unboundedPreceding, Window.currentRow) df_filled = df_with_group.withColumn( "matching_tracking_id", F.last(F.col("tracking ID"), ignorenulls=True).over(window_fill) )提取目标结果
如果只需要cust ID行和对应的追踪ID,可以筛选出来:result_df = df_filled.filter(F.col("value") == "cust ID").select( "value", "line num", "matching_tracking_id" )
最终结果示例
运行完上述代码后,result_df的输出会是:
| value | line num | matching_tracking_id |
|---|---|---|
| cust ID | 3 | 123 |
| cust ID | 7 | 456 |
如果你需要保留所有行并带上对应的追踪ID,直接使用df_filled即可。
备注:内容来源于stack exchange,提问作者Chuck
相关产品推荐
相关产品推荐

