如何在PySpark DataFrame中将单词列表合并为句子?
PySpark:将数组列合并为字符串句子(无需RDD方案)
问题场景
我有一个PySpark DataFrame,其中words列是字符串数组类型,需要将数组内的单词合并成单个字符串句子。尝试用RDD方法处理时未得到预期结果,希望找到无需转换RDD的解决方案。
原始DataFrame示例
temp = spark.createDataFrame([ (0, ['Julia', 'is', 'awesome']), (2, ['Data-science', 'is','cool']), (3, ['Machine','learning']) ], ["id", "words"]) # 输出预览 # +---+------------------------+ # |id |words | # +---+------------------------+ # |0 |[Julia, is, awesome] | # |2 |[Data-science, is, cool]| # |3 |[Machine, learning] | # +---+------------------------+ temp.printSchema() # root # |-- id: long (nullable = true) # |-- words: array (nullable = true) # | |-- element: string (containsNull = true)
失败的RDD尝试
rdd_df = temp.rdd.map(lambda x: [x['id'], ' '.join(x['words'])]) spark.createDataFrame(rdd_df, temp.schema).show(10, False) # 错误输出 # +---+---------------------------------------------------------+ # |id |words | # +---+---------------------------------------------------------+ # |0 |[ ' J u l i a ' , ' i s ' , ' a w e s o m e ' ] | # |2 |[ ' D a t a - s c i e n c e ' , ' i s ' , ' c o o l ' ]| # |3 |[ ' M a c h i n e ' , ' l e a r n i n g ' ] | # +---+---------------------------------------------------------+
错误原因:创建新DataFrame时强行沿用了原schema(words为数组类型),但' '.join(x['words'])生成的是字符串,Spark会自动将字符串拆分为字符数组,导致异常输出。
预期输出
+---+------------------------+ |id |words | +---+------------------------+ |0 |Julia is awesome | |2 |Data-science is cool | |3 |Machine learning | +---+------------------------+
无需RDD的解决方案
方法1:使用concat_ws函数(推荐)
PySpark内置的concat_ws函数可指定分隔符,直接将数组元素拼接为字符串,完全基于DataFrame API操作:
from pyspark.sql.functions import concat_ws result = temp.withColumn("words", concat_ws(" ", "words")) result.show(10, False)
输出与预期完全一致。
方法2:使用array_join函数(Spark 2.4+支持)
array_join是专门用于数组拼接的函数,功能与concat_ws类似:
from pyspark.sql.functions import array_join result = temp.withColumn("words", array_join("words", " ")) result.show(10, False)
输出与方法1相同。
(可选)修正后的RDD方法
如果一定要用RDD,需要修改schema,将words列改为字符串类型:
from pyspark.sql.types import StructType, StructField, LongType, StringType new_schema = StructType([ StructField("id", LongType(), True), StructField("words", StringType(), True) ]) rdd_df = temp.rdd.map(lambda x: (x['id'], ' '.join(x['words']))) spark.createDataFrame(rdd_df, new_schema).show(10, False)
但显然DataFrame API的方法更简洁高效,无需额外处理schema。
内容的提问来源于stack exchange,提问作者merkle
相关产品推荐
相关产品推荐

