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

如何使用PySpark基于现有列值匹配提取字段并创建新列?

PySpark提取嵌套列表匹配值的实现方案

问题场景

现有如下结构的PySpark DataFrame:

orderidsubfilter-list
1367[['123','supply'],['367','price']]
2389[['389','supply'],['906','supply']]
3804[['173','supply'],['804','price']]

需求:从filter-list列的嵌套子列表中,筛选出第一个元素与sub列数值匹配的子列表,提取该子列表的第二个元素,生成新列filter-name,最终得到如下DataFrame:

orderidsubfilter-listfilter-name
1367[['123','supply'],['367','price']]price
2389[['389','supply'],['906','supply']]supply
3804[['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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 07:27:30