如何用Lambda表达式从字典动态创建PySpark DataFrame列?
解决PySpark动态生成列的Lambda表达式报错问题
问题场景
需要基于给定字典动态创建PySpark DataFrame列,规则如下:
- 当字典中
Technical Column值为No时,从原DataFrame提取对应Column Mapping列的值 - 当值为
Yes时,用lit函数生成Column Mapping对应值的常量列
给定字典示例:
cols = [ {'Technical Label':'UserID', 'Column Mapping':'UID','Technical Column':'No'}, {'Technical Label':'Type', 'Column Mapping':'X_TYPE','Technical Column':'No'}, {'Technical Label':'Name.1', 'Column Mapping':'NaturalName','Technical Column':'No'}, {'Technical Label':'Address.1', 'Column Mapping':'HomeAddress','Technical Column':'Yes'}, {'Technical Label':'Identifier.1', 'Column Mapping':'Human','Technical Column':'Yes'}, {'Technical Label':'Identifier.1', 'Column Mapping':'EX_CO','Technical Column':'Yes'}, {'Technical Label':'Identifier.IdentifierValue.1', 'Column Mapping':'X_CODE','Technical Column':'No'} ]
期望实现的效果类似:
refdf = df.withColumn('UserID', df.UID).withColumn('Type', df.X_TYPE)...
用户编写的报错代码:
from functools import reduce from pyspark.sql.functions import col, lit df = spark.read.format('csv').option('header','true').option('delimiter','|').load('/my/path/file.txt') x = reduce(lambda df,colsm : df.withColumn(colsm['Technical Label'], (lambda x: x['Column Mapping'] if x['Technical Column'] == 'No' else lit(x['Column Mapping']) , cols)), cols, df) display(x)
报错信息:
PySparkTypeError: [NOT_COLUMN] Argument
colshould be a Column, got tuple.
错误原因
withColumn的第二个参数传入了一个元组(lambda x: ... , cols),而不是PySpark要求的Column对象,这是报错的核心原因。- 内部嵌套的Lambda表达式完全多余,直接针对当前遍历的
colsm字典做条件判断即可,不需要额外定义Lambda。
修正后的代码
from functools import reduce from pyspark.sql.functions import col, lit # 读取原始CSV数据 df = spark.read.format('csv').option('header','true').option('delimiter','|').load('/my/path/file.txt') # 用reduce动态批量生成目标列 refdf = reduce( lambda current_df, col_config: current_df.withColumn( col_config['Technical Label'], # 根据配置判断是取原列还是生成常量列 col(col_config['Column Mapping']) if col_config['Technical Column'] == 'No' else lit(col_config['Column Mapping']) ), cols, df # 初始值为原始DataFrame ) display(refdf)
代码说明
- 遍历
cols中的每一项配置,针对每个配置项:- 若
Technical Column为No,调用col()函数获取原DataFrame中对应列 - 若为
Yes,调用lit()函数生成对应值的常量列
- 若
reduce函数会迭代处理每一项配置,将每次withColumn生成的新DataFrame作为下一次迭代的输入,最终得到包含所有目标列的结果DataFrame
内容的提问来源于stack exchange,提问作者Alex Raj Kaliamoorthy
相关产品推荐
相关产品推荐

