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

PySpark实现反向匹配模式并添加方案类型列

问题:PySpark DataFrame按方案类型标记行

原始DataFrame

_c0_c1_c2_c3
开放方案nullnullnull
Adnullnullnull
119551INF209KA12Z1INF209KA13Z9Ad
119552INF209MK12Z1INF209MKA13Z9Ad
关闭方案nullnullnull
Adnullnullnull
119553INF209KA12aZ1INF209KAd13Z9Ad
119554INF209MK12Zs1INF209MKAd13Z9Ad

需求

  • 为“开放方案”和“关闭方案”之间的所有行添加新列方案类型,值为“开放方案”
  • 为“关闭方案”之后的所有行添加该列,值为“关闭方案”
  • “开放方案”“关闭方案”及紧随其后的Ad行,方案类型列值为null

预期输出

_c0_c1_c2_c3方案类型
开放方案nullnullnullnull
Adnullnullnullnull
119551INF209KA12Z1INF209KA13Z9Ad开放方案
119552INF209MK12Z1INF209MKA13Z9Ad开放方案
关闭方案nullnullnullnull
Adnullnullnullnull
119553INF209KA12aZ1INF209KAd13Z9Ad关闭方案
119554INF209MK12Zs1INF209MKAd13Z9Ad关闭方案

尝试思路

原本考虑使用带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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 00:55:20