Scala代码中collect用法合理性及优化方案咨询
关于Spark Scala代码中collect和map的疑问解答
一、collect和map是不是同类函数?
完全不是,二者属于Spark中两类完全不同的操作:
map是转换操作(Transformation):属于懒加载逻辑,调用时不会触发实际计算,只是给原数据集添加一个转换规则,生成新的分布式数据集(DataFrame/Dataset/RDD)。collect是动作操作(Action):会触发整个计算链的执行,把集群上分散存储的所有数据拉取到Driver节点的内存里,返回一个本地的数组(Array)。
二、原代码中collect的用法合理吗?
从语法和功能实现上看,这段代码是能正常运行的,最终可以得到Driver端的Seq[String]类型学生名单,但合理性要看数据规模:
- 如果去重后的学生数据量很小(比如几万条以内),这个写法没问题;
- 如果数据量很大,
collect会把大量数据集中到Driver节点,极易引发内存溢出(OOM),这种场景下就不合理。
三、效率不足该怎么改?是不是只需移除collect?
不能直接只移除collect,得根据后续需求调整:
- 若后续仍需在Spark集群做分布式计算:
应该去掉collect和toSeq,保留分布式的Dataset[String],让计算全程在集群分布式执行,效率最优。修改后代码:val students: Dataset[String] = students_df .select(col("name")) .distinct .map(_.getString(0)) - 若确实需要在Driver端用本地的
Seq[String]:collect是必须的,但可以做小优化:- 先确认数据量足够小,避免内存溢出;
- 可以用
collectAsList()直接返回List,再转成Scala Seq,省去toSeq步骤:val students = students_df .select(col("name")) .distinct .map(_.getString(0)) .collectAsList() .asScala
内容的提问来源于stack exchange,提问作者fqv572
相关产品推荐
相关产品推荐

