Pyspark 如何根据df1的questions列指定列名从df2取值新增列
问题分析与解决方法
报错原因
- 语法参数错误:
df1.select('questions').collect()返回的是df1全量行的Row对象列表,无法直接作为列名参数传给df2.select(),只有取单行列值的collect()[0][0]才能得到合法的列名字符串列表。 - 运算逻辑错误:Spark的
withColumn是行级并行运算,你没有建立df1和df2的行关联关系,直接在df1的列运算中引用无关联的df2列,Spark无法解析对应取值逻辑。
实现代码
从你给出的样例来看df2为单行的全局映射表,优先采用将小表广播为字典+UDF的方案,性能最优:
import pyspark.sql.functions as F from pyspark.sql.types import ArrayType, IntegerType, StringType # 先将单行df2转为{Q列名: 对应值}的字典,全量拉到Driver内存,开销极低 q_value_map = df2.first().asDict() # 定义UDF:输入问题列表,返回对应值的数组 @F.udf(returnType=ArrayType(IntegerType())) def map_question_to_value(question_list): return [q_value_map.get(q) for q in question_list] # 直接给df1加values列 df1 = df1.withColumn('values', map_question_to_value('questions'))
如果你需要的是逗号分隔的字符串格式而非数组,修改UDF即可:
@F.udf(returnType=StringType()) def map_question_to_value(question_list): return ', '.join([str(q_value_map.get(q)) for q in question_list])
如果你的df2是多行结构,需要先给两个DataFrame添加相同的关联主键(比如行ID),执行join关联后再做行内的数组取值处理即可。
内容的提问来源于stack exchange,提问作者Turvy
相关产品推荐
相关产品推荐

