PySpark生成lookup_l/lookup_r列:按指定列名列表提取对应值
PySpark实现动态提取指定列值生成列表列
原始DataFrame
| first_name_l | first_name_r | last_name_l | last_name_r | dob_l | dob_r | city_l | city_r | average_score | matched_columns | lookup_l_list | lookup_r_list |
|---|---|---|---|---|---|---|---|---|---|---|---|
| robert | robert | null | allen | 1971-06-24 | 1971-05-24 | null | null | 49.57 | [first_name_score, dob_score] | [first_name_l, dob_l] | [first_name_r, dob_r] |
| null | robert | allen | alen | 1971-06-24 | 1971-06-24 | london | lonon | 69.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
相关产品推荐
相关产品推荐

