原创经验分享(15)spark sql limit实现原理

Posted Thinking in BigData

tags:

篇首语:本文由小常识网(cha138.com)小编为大家整理,主要介绍了原创经验分享(15)spark sql limit实现原理相关的知识,希望对你有一定的参考价值。

之前讨论过hive中limit的实现,详见 https://www.cnblogs.com/barneywill/p/10109217.html
下面看spark sql中limit的实现,首先看执行计划:

spark-sql> explain select * from test1 limit 10;
== Physical Plan ==
CollectLimit 10
+- HiveTableScan [id#35], MetastoreRelation temp, test1
Time taken: 0.201 seconds, Fetched 1 row(s)

limit对应的CollectLimit,对应的实现类是

org.apache.spark.sql.execution.CollectLimitExec

case class CollectLimitExec(limit: Int, child: SparkPlan) extends UnaryExecNode {
...
  protected override def doExecute(): RDD[InternalRow] = {
    val locallyLimited = child.execute().mapPartitionsInternal(_.take(limit))
    val shuffled = new ShuffledRowRDD(
      ShuffleExchange.prepareShuffleDependency(
        locallyLimited, child.output, SinglePartition, serializer))
    shuffled.mapPartitionsInternal(_.take(limit))
  }

可见实现非常简单,首先调用SparkPlan.execute得到结果的RDD,然后从每个partition中取前limit个row得到一个新的RDD,然后再将这个新的RDD变成一个分区,然后再取前limit个,这样就得到最终的结果。

 

以上是关于原创经验分享(15)spark sql limit实现原理的主要内容,如果未能解决你的问题,请参考以下文章

原创大叔经验分享(39)spark cache unpersist级联操作

原创大叔经验分享(55)hue导出行数限制

原创经验分享(17)编程实践对比Java vs Scala

原创大叔经验分享(23)hive metastore的几种部署方式

原创问题定位分享(18)beeline连接spark thrift有时会卡住

原创大叔经验分享(76)confluence配置