PySpark:将含词元列表的RDD转换为每行一个词元的RDD
解决Spark RDD拆分词元列表的问题
嗨,这个需求很容易实现,你需要用到Spark的flatMap操作,但要注意处理每个子列表的方式——因为你需要把每个词元都单独包装成一个单元素列表,而不是直接展开成单个元素。
具体实现代码
首先,你的原始RDD定义没问题:
mylist = [['hello'], ['cat'], ['dog'], ['hey'], ['dog'], ['I', 'need', 'coffee'], ['dance'], ['dream', 'job']] myRDD = sc.parallelize(mylist)
接下来用flatMap来转换:
# 对每个子列表,把每个词元包装成单元素列表,再扁平化为RDD的每行 result_rdd = myRDD.flatMap(lambda sub_list: [[word] for word in sub_list])
验证结果
你可以用collect()方法查看转换后的结果:
print(result_rdd.collect())
输出完全符合你的预期:
[['hello'], ['cat'], ['dog'], ['hey'], ['dog'], ['I'], ['need'], ['coffee'], ['dance'], ['dream'], ['job']]
为什么这么做?
- 如果直接用
flatMap(lambda x: x),会得到一个包含单个词元的RDD(比如['hello', 'cat', 'dog', ...]),这不符合你要每行都是单元素列表的要求。 - 而我们用列表推导式
[[word] for word in sub_list],会把每个子列表里的每个词元都包装成独立的单元素列表,再通过flatMap把这些小列表全部扁平到RDD的每一行中。
内容的提问来源于stack exchange,提问作者strv7
相关产品推荐
相关产品推荐

