PySpark多条件elif逻辑实现:基于数组位置提取指定值
问题描述
现有如下结构的PySpark DataFrame:
| productNo | prodcuctMT | productPR | productList |
|---|---|---|---|
| 2389 | ['xy-5', 'yz-12','zb-56','iu-30'] | ['pr-1', 'pr-2', 'pr-3', 'pr-4'] | ['67230','7839','1339','9793'] |
| 6745 | ['xy-4', 'yz-34','zb-8','iu-9'] | ['pr-6', 'pr-1', 'pr-3', 'pr-7'] | ['1111','0987','8910','0348'] |
需求:实现多条件匹配逻辑,找到同一位置满足prodcuctMT = 'xy-5'且productPR = 'pr-1'的行,提取对应位置的productList值作为新增列new_productList,预期结果仅保留符合条件的行:
| productNo | prodcuctMT | productPR | productList | new_productList |
|---|---|---|---|---|
| 2389 | ['xy-5', 'yz-12','zb-56','iu-30'] | ['pr-1', 'pr-2', 'pr-3', 'pr-4'] | ['67230','7839','1339','9793'] | 67230 |
之前尝试用filter处理但无法适配多条件遍历的需求,需要更完善的实现方案。
解决方案
可以通过PySpark的数组关联+条件筛选组合操作实现,步骤如下:
实现代码
from pyspark.sql import SparkSession import pyspark.sql.functions as F # 初始化SparkSession spark = SparkSession.builder.appName("array_match").getOrCreate() # 创建测试DataFrame data = [ ("2389", ["xy-5", "yz-12", "zb-56", "iu-30"], ["pr-1", "pr-2", "pr-3", "pr-4"], ["67230", "7839", "1339", "9793"]), ("6745", ["xy-4", "yz-34", "zb-8", "iu-9"], ["pr-6", "pr-1", "pr-3", "pr-7"], ["1111", "0987", "8910", "0348"]) ] df = spark.createDataFrame(data, ["productNo", "prodcuctMT", "productPR", "productList"]) # 1. 关联三个数组的对应位置为结构体数组 df_zip = df.withColumn("zipped", F.arrays_zip("prodcuctMT", "productPR", "productList")) # 2. 筛选符合条件的结构体元素 df_filtered = df_zip.withColumn( "matched", F.filter( "zipped", lambda x: (x.prodcuctMT == "xy-5") & (x.productPR == "pr-1") ) ) # 3. 提取匹配到的productList值,同时过滤无匹配的行 result_df = df_filtered.filter(F.size("matched") > 0)\ .withColumn("new_productList", F.col("matched")[0].productList)\ .drop("zipped", "matched") # 展示结果 result_df.show(truncate=False)
代码说明
arrays_zip:将三个数组列的对应位置元素打包成结构体,确保同一位置的prodcuctMT、productPR、productList被关联在一起。filter:对结构体数组进行多条件筛选,只保留满足prodcuctMT='xy-5'且productPR='pr-1'的元素。size("matched") > 0:过滤掉没有匹配结果的行,和预期结果一致。F.col("matched")[0].productList:提取第一个匹配项的productList值,如果存在多个匹配可根据需求调整索引或用聚合函数处理。
扩展:处理多条件elif逻辑
如果需要实现多优先级的elif逻辑(比如先匹配条件A,不满足再匹配条件B,以此类推),可以用transform结合case when来实现优先级判断,示例如下:
df_priority = df_zip.withColumn( "matched", F.transform( "zipped", lambda x: F.when( (x.prodcuctMT == "xy-5") & (x.productPR == "pr-1"), x.productList ).when( (x.prodcuctMT == "xy-4") & (x.productPR == "pr-1"), x.productList ).otherwise(None) ) ).withColumn( "new_productList", F.first(F.col("matched"), ignorenulls=True).over() ).filter(F.col("new_productList").isNotNull())
这段代码会按顺序匹配条件,取第一个符合优先级的productList值。
内容的提问来源于stack exchange,提问作者Yas
相关产品推荐
相关产品推荐

