Spark中使用Scala如何将.map操作得到的List转为RDD?
解答
你遇到的问题本质是混淆了Spark的转换算子和行动算子的特性:RDD的map()方法本身属于转换算子,执行后默认返回的就是RDD类型,你现在拿到List结果,是因为你在map()操作之后额外调用了collect()、collectAsList()这类行动算子,这类算子会把分布式存储的RDD数据拉取到当前Driver节点,生成本地集合对象。
解决方法
- 方案1(推荐,符合Spark作业的常规实现逻辑):只保留
map()转换操作,删除后续追加的行动算子调用,此时你拿到的结果天然就是RDD类型,完全满足作业要求。 - 方案2(仅适合中间需要用本地List做校验的场景):如果你中间确实需要用到List做本地验证、计算,后续要转回RDD,直接调用
SparkContext实例的parallelize()方法即可将现有List转为RDD,不同语言参考写法:// Scala 写法 val targetRDD = sc.parallelize(你的List变量名)// Java 写法 JavaRDD<你的数据类型> targetRDD = sc.parallelize(你的List变量名);# PySpark 写法 target_rdd = sc.parallelize(你的列表变量名)
注意:方案2仅适合作业这类小数据量场景,大数据量下把全量RDD数据拉取到单节点生成本地List,很容易触发Driver节点内存溢出。生产环境优先选择全程使用RDD转换算子实现业务逻辑,避免中间转本地集合的操作。
内容的提问来源于stack exchange,提问作者Tudor
相关产品推荐
相关产品推荐

