You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 01:03:26