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

PySpark生成lookup_l/lookup_r列:按指定列名列表提取对应值

PySpark实现动态提取指定列值生成列表列

原始DataFrame

first_name_lfirst_name_rlast_name_llast_name_rdob_ldob_rcity_lcity_raverage_scorematched_columnslookup_l_listlookup_r_list
robertrobertnullallen1971-06-241971-05-24nullnull49.57[first_name_score, dob_score][first_name_l, dob_l][first_name_r, dob_r]
nullrobertallenalen1971-06-241971-06-24londonlonon69.95[dob_score, city_score, last_name_score][dob_l, city_l, last_name_l][dob_r, city_r, last_name_r]

需求

  • 新增lookup_l列:值为当前行lookup_l_list中指定列名对应的行值组成的列表,顺序与lookup_l_list一致(如第一行应为["robert","1971-06-24"])
  • 新增lookup_r列:值为当前行lookup_r_list中指定列名对应的行值组成的列表,顺序与lookup_r_list一致(如第一行应为["robert","1971-05-24"])

解决方案

由于需要根据每行动态指定的列名提取值,可通过PySpark UDF(用户自定义函数)实现,代码如下:

1. 导入依赖模块

from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
from pyspark.sql.types import ArrayType, StringType

2. 定义提取值的UDF函数

def get_values_from_columns(row, col_list):
    # 按列名列表顺序提取对应行值
    return [row[col] for col in col_list]

# 注册UDF,指定返回类型为字符串数组
extract_values_udf = udf(lambda row, cols: get_values_from_columns(row, cols), ArrayType(StringType()))

3. 应用UDF生成新列

# 假设原始DataFrame变量名为df
df = df.withColumn("lookup_l", extract_values_udf(df, df["lookup_l_list"]))
df = df.withColumn("lookup_r", extract_values_udf(df, df["lookup_r_list"]))

4. 查看结果

df.show(truncate=False)

执行后即可得到包含原有所有列及新增lookup_l、lookup_r列的目标DataFrame。

内容的提问来源于stack exchange,提问作者mahak tirole

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 22:57:55