求助:基于另一列前4位修改PySpark DataFrame指定列的实现方法
PySpark DataFrame 分组逻辑修正实现
原始数据与问题
现有如下结构的PySpark DataFrame:
code 1 code 2 Fruit_Group temp_code (string) (string) (string) (boolean) 12E5-11 12E5-11 Apple True 12E5-11 ERE5-11,12E5-11 Apple True 12E5-11 MMMM-11 Apple True # 需修改为Banana 12E5-11 XXXX-11 Apple False # 需修改为Orange 12E5-11 12E5-11 Apple True 12E5-11 12E5-11, ERE5-11 Apple True
数据初始化代码:
x = [ ("12E5-11", "12E5-11", "Apple", True), ("12E5-11", "ERE5-11,12E5-11", "Apple", True), ("12E5-11", "MMMM-11", "Apple", True ), ("12E5-11", "XXXX-11" ,"Apple", False), ("12E5-11", "12E5-11", "Apple", True), ("12E5-11", "ERE5-11,12E5-11", "Apple", True) ] Fruits_df = spark.createDataFrame(x, schema=["code 1", "code 2","Fruit_Group","temp_code"])
修正规则
需按以下逻辑更新Fruit_Group字段:
- 若
code 1的前4位,与code 2按逗号拆分后的任意一项的前4位匹配,则保持Fruit_Group为Apple - 若未匹配到:
- 当
temp_code为True时,设置为Banana - 当
temp_code为False时,设置为Orange
- 当
实现代码
from pyspark.sql import functions as F # 处理code2:拆分、去除空格、提取前4位,生成匹配数组 processed_df = Fruits_df.withColumn( "code2_prefixes", F.transform( F.split(F.trim(F.col("code 2")), ",\\s*"), # 按逗号+可选空格拆分 lambda code: F.substring(code, 1, 4) # 提取每个code的前4位 ) ) # 根据规则更新Fruit_Group final_df = processed_df.withColumn( "Fruit_Group", F.when( F.array_contains(F.col("code2_prefixes"), F.substring(F.col("code 1"), 1, 4)), F.col("Fruit_Group") # 匹配成功,保持原Apple ).when( F.col("temp_code") == True, "Banana" ).otherwise( "Orange" ) ).drop("code2_prefixes") # 移除中间辅助列 # 查看结果 final_df.show(truncate=False)
代码说明
- 拆分与提取前缀:用
split拆分code 2,同时处理逗号后的空格;用transform对每个拆分后的code提取前4位,生成前缀数组 - 匹配判断:用
array_contains检查code 1的前4位是否存在于前缀数组中 - 条件更新:用
when/otherwise实现分支逻辑,最后移除辅助列
执行结果
最终DataFrame的第3、4行Fruit_Group会被修正为Banana和Orange,其余行保持Apple。
内容的提问来源于stack exchange,提问作者Bella_18
相关产品推荐
相关产品推荐

