PySpark中map()处理单列转DF报TypeError多列正常的原因
问题场景
使用PySpark的map()做列转换时,返回单列结果调用toDF()会抛出TypeError,相同逻辑返回多列时可正常运行。
报错的单列处理代码
from pyspark.sql import SparkSession data = [('James','Smith','M',30), ('Anna','Rose','F',41), ('Robert','Williams','M',62)] rdd = spark.sparkContext.parallelize(data) columns = ["firstname","lastname","gender","salary"] df = spark.createDataFrame(data,schema = columns) rdd2=df.rdd.map(lambda x: (x["firstname"]+","+ x["lastname"]) ) df2=rdd2.toDF(["fullname"]) df2.show()
运行抛出TypeError:
可正常运行的多列处理代码
from pyspark.sql import SparkSession data = [('James','Smith','M',30), ('Anna','Rose','F',41), ('Robert','Williams','M',62)] rdd = spark.sparkContext.parallelize(data) columns = ["firstname","lastname","gender","salary"] df = spark.createDataFrame(data,schema = columns) rdd2=df.rdd.map(lambda x: (x["firstname"]+","+ x["lastname"],x["gender"]) ) df2=rdd2.toDF(["fullname","gender"])
运行结果:
问题根因
PySpark的RDD.toDF()方法转换DataFrame时,要求RDD内的每一条记录必须是长度和传入列名数匹配的结构化序列(元组、列表、pyspark.sql.Row对象均可)。
两段代码的核心差异在lambda的返回值类型:
- 多列场景下lambda返回
(x["firstname"]+","+ x["lastname"],x["gender"]),是长度为2的标准元组,和传入的列名列表["fullname","gender"]长度匹配,因此可以正常转换。 - 单列场景下写的
(x["firstname"]+","+ x["lastname"])在Python语法中不属于元组——括号仅作为优先级运算符存在,没有类型标记作用,实际返回值是普通字符串。字符串本身是逐字符可迭代对象,toDF()会尝试把单个字符串拆分为多个独立字符作为一行的多个字段,此时字段数量和传入的1个列名完全不匹配,直接触发TypeError。
修复方案
方案1:修正RDD map的返回值格式
在单元素末尾加逗号,显式标记为长度1的元组即可:
from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() data = [('James','Smith','M',30), ('Anna','Rose','F',41), ('Robert','Williams','M',62)] columns = ["firstname","lastname","gender","salary"] df = spark.createDataFrame(data,schema = columns) # 注意返回值末尾的逗号,标识为单元素元组 rdd2=df.rdd.map(lambda x: (x["firstname"]+","+ x["lastname"],) ) df2=rdd2.toDF(["fullname"]) df2.show()
运行输出:
+---------------+ | fullname| +---------------+ | James,Smith| | Anna,Rose| |Robert,Williams| +---------------+
方案2:直接使用DataFrame原生API(推荐)
转RDD再map会触发JVM和Python进程间的序列化开销,性能远低于原生API,列拼接直接用内置函数即可:
from pyspark.sql.functions import concat_ws df2 = df.select(concat_ws(",", "firstname", "lastname").alias("fullname"))
内容的提问来源于stack exchange,提问作者awesome_sangram
相关产品推荐
相关产品推荐

