PySpark实现反向匹配模式并添加方案类型列
问题:PySpark DataFrame按方案类型标记行
原始DataFrame
| _c0 | _c1 | _c2 | _c3 |
|---|---|---|---|
| 开放方案 | null | null | null |
| Ad | null | null | null |
| 119551 | INF209KA12Z1 | INF209KA13Z9 | Ad |
| 119552 | INF209MK12Z1 | INF209MKA13Z9 | Ad |
| 关闭方案 | null | null | null |
| Ad | null | null | null |
| 119553 | INF209KA12aZ1 | INF209KAd13Z9 | Ad |
| 119554 | INF209MK12Zs1 | INF209MKAd13Z9 | Ad |
需求
- 为“开放方案”和“关闭方案”之间的所有行添加新列方案类型,值为“开放方案”
- 为“关闭方案”之后的所有行添加该列,值为“关闭方案”
- “开放方案”“关闭方案”及紧随其后的Ad行,方案类型列值为null
预期输出
| _c0 | _c1 | _c2 | _c3 | 方案类型 |
|---|---|---|---|---|
| 开放方案 | null | null | null | null |
| Ad | null | null | null | null |
| 119551 | INF209KA12Z1 | INF209KA13Z9 | Ad | 开放方案 |
| 119552 | INF209MK12Z1 | INF209MKA13Z9 | Ad | 开放方案 |
| 关闭方案 | null | null | null | null |
| Ad | null | null | null | null |
| 119553 | INF209KA12aZ1 | INF209KAd13Z9 | Ad | 关闭方案 |
| 119554 | INF209MK12Zs1 | INF209MKAd13Z9 | Ad | 关闭方案 |
尝试思路
原本考虑使用带unbounded preceding参数的lag函数,反向遍历直到找到_c0列中包含方案关键词的值,但不清楚如何给lag函数添加匹配条件。
解决方案
可以通过标记方案行→填充方案类型→过滤不需要标记的行三步实现:
步骤1:标记方案起始行并生成分组标识
先创建辅助列标记方案起始行,再用累计求和生成分组ID,让同一方案下的行归为一组:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 标记方案起始行(开放/关闭方案) df = df.withColumn("is_scheme", F.when(F.col("_c0").isin("开放方案", "关闭方案"), 1).otherwise(0)) # 生成分组ID,维持行顺序 window_order = Window.orderBy(F.monotonically_increasing_id()) df = df.withColumn("scheme_group", F.sum("is_scheme").over(window_order))
步骤2:填充每个分组的方案类型
提取分组内的方案名称,填充到同组所有行:
# 提取分组对应的方案类型 window_group = Window.partitionBy("scheme_group") df = df.withColumn("方案类型", F.first(F.when(F.col("is_scheme") == 1, F.col("_c0"))).over(window_group))
步骤3:过滤不需要标记的行
将方案行和紧随其后的Ad行的方案类型设为null:
# 标记方案行后的Ad行 df = df.withColumn( "is_scheme_ad", F.lag(F.col("is_scheme"), 1).over(window_order) == 1 ) # 最终调整方案类型列 df = df.withColumn( "方案类型", F.when( (F.col("is_scheme") == 1) | F.col("is_scheme_ad"), None ).otherwise(F.col("方案类型")) ) # 清理辅助列 df = df.drop("is_scheme", "scheme_group", "is_scheme_ad")
完整代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 假设df为原始DataFrame # 步骤1:标记方案行并生成分组 df = df.withColumn("is_scheme", F.when(F.col("_c0").isin("开放方案", "关闭方案"), 1).otherwise(0)) window_order = Window.orderBy(F.monotonically_increasing_id()) df = df.withColumn("scheme_group", F.sum("is_scheme").over(window_order)) # 步骤2:填充分组内的方案类型 window_group = Window.partitionBy("scheme_group") df = df.withColumn("方案类型", F.first(F.when(F.col("is_scheme") == 1, F.col("_c0"))).over(window_group)) # 步骤3:过滤不需要标记的行 df = df.withColumn("is_scheme_ad", F.lag(F.col("is_scheme"), 1).over(window_order) == 1) df = df.withColumn( "方案类型", F.when( (F.col("is_scheme") == 1) | F.col("is_scheme_ad"), None ).otherwise(F.col("方案类型")) ) # 清理辅助列 df = df.drop("is_scheme", "scheme_group", "is_scheme_ad") # 查看结果 df.show()
说明
monotonically_increasing_id()用于维持原始行顺序,若数据有天然排序键可替换为对应列scheme_group将“开放方案”到“关闭方案”的行归为一组,“关闭方案”之后的行归为另一组first()窗口函数提取分组内的方案名称,自动填充到同组所有行- 最后通过
lag()识别方案行后的Ad行,将这类行和方案行的方案类型设为null
内容的提问来源于stack exchange,提问作者swap k
相关产品推荐
相关产品推荐

