PySpark中如何实现可动态传入列名的withColumn自定义函数?
解决PySpark动态列名函数的NameError问题
你遇到的NameError根源很明确:在修改后的函数里,你把column_to_add用单引号括起来了,这时候它被当成了一个固定的字符串字面量,而不是你传入的参数变量。PySpark会尝试去查找名为column_to_add的列,但这个列根本不存在,所以报错。
正确的通用函数写法
只需要去掉withColumn第一个参数的单引号,直接传入参数变量即可——因为withColumn的第一个参数本身就接受字符串类型的列名,直接用你传入的参数值作为新列名就对了。另外建议把参数名list改成check_list,避免和Python内置的list类型冲突:
from pyspark.sql.functions import when def new_column(df, check_list, column_to_add, column_to_check): df1 = df.withColumn(column_to_add, when(df[column_to_check].isin(check_list), "Y").otherwise('N')) return df1
调用示例
现在你可以轻松生成不同的新列了:
# 生成city_visited列,检查city是否在指定列表中 city_list = ["New York", "London"] new_df = new_column(df, city_list, "city_visited", "city") # 生成bucket_list列,检查country是否在指定列表中 country_list = ["Japan", "Australia"] new_df = new_column(new_df, country_list, "bucket_list", "country")
额外优化建议
如果你习惯使用col()函数来引用列,也可以把df[column_to_check]换成col(column_to_check),代码会更简洁易读:
from pyspark.sql.functions import col, when def new_column(df, check_list, column_to_add, column_to_check): df1 = df.withColumn(column_to_add, when(col(column_to_check).isin(check_list), "Y").otherwise('N')) return df1
内容的提问来源于stack exchange,提问作者User12345
相关产品推荐
相关产品推荐

