如何使用PySpark基于现有列值匹配提取字段并创建新列?
PySpark提取嵌套列表匹配值的实现方案
问题场景
现有如下结构的PySpark DataFrame:
| orderid | sub | filter-list |
|---|---|---|
| 1 | 367 | [['123','supply'],['367','price']] |
| 2 | 389 | [['389','supply'],['906','supply']] |
| 3 | 804 | [['173','supply'],['804','price']] |
需求:从filter-list列的嵌套子列表中,筛选出第一个元素与sub列数值匹配的子列表,提取该子列表的第二个元素,生成新列filter-name,最终得到如下DataFrame:
| orderid | sub | filter-list | filter-name |
|---|---|---|---|
| 1 | 367 | [['123','supply'],['367','price']] | price |
| 2 | 389 | [['389','supply'],['906','supply']] | supply |
| 3 | 804 | [['173','supply'],['804','price']] | price |
解决方案
方法一:使用Spark内置函数(推荐,性能更优)
利用Spark原生函数实现,无需自定义UDF,避免Python与JVM之间的性能开销:
from pyspark.sql import functions as F # 假设原始DataFrame变量名为df result_df = df.withColumn( "filter-name", # 筛选匹配的子列表,取第一个匹配项的第二个元素 F.element_at( F.filter( F.col("filter-list"), # 将sub转为字符串,与子列表第一个元素匹配 lambda item: item[0] == F.col("sub").cast("string") ), 1 # Spark的element_at为1-based索引,取第一个匹配的子列表 )[1] # 取子列表的第二个元素(数组索引从0开始) ) # 查看结果 result_df.show(truncate=False)
方法二:使用自定义UDF(适合复杂逻辑场景)
如果对内置函数逻辑不熟悉,可通过自定义UDF实现:
from pyspark.sql import functions as F from pyspark.sql.types import StringType def get_filter_name(sub_val, filter_list): # 遍历嵌套列表,匹配sub的字符串形式 for item in filter_list: if item[0] == str(sub_val): return item[1] # 无匹配项时返回None return None # 注册UDF filter_name_udf = F.udf(get_filter_name, StringType()) # 生成新列 result_df = df.withColumn("filter-name", filter_name_udf(F.col("sub"), F.col("filter-list"))) # 查看结果 result_df.show(truncate=False)
内容的提问来源于stack exchange,提问作者krishna kaushik
相关产品推荐
相关产品推荐

