【问题标题】:How to know which count query is the fastest?如何知道哪个计数查询最快?
【发布时间】:2017-10-06 05:01:45
【问题描述】:

我一直在探索最近版本的 Spark SQL 2.3.0-SNAPSHOT 中的查询优化,并注意到语义相同查询的不同物理计划。

假设我必须计算以下数据集中的行数:

val q = spark.range(1)

我可以按如下方式计算行数:

  1. q.count
  2. q.collect.size
  3. q.rdd.count
  4. q.queryExecution.toRdd.count

我最初的想法是,它几乎是一个恒定的操作(肯定是由于本地数据集),以某种方式已被 Spark SQL 优化并立即给出结果,尤其是。第一个 Spark SQL 完全控制查询执行。

看过查询的物理计划后,我相信最有效的查询将是最后一个:

q.queryExecution.toRdd.count

原因是:

  1. 它避免了从 InternalRow 二进制格式反序列化行
  2. 查询是代码生成的
  3. 只有一个工作有一个阶段

物理计划就是这么简单。

我的推理正确吗?如果是这样,如果我从外部数据源(例如文件、JDBC、Kafka)读取数据集,答案会有所不同吗?

主要问题是要考虑哪些因素来判断查询是否比其他查询更有效(根据此示例)?


完整性的其他执行计划。

q.count

q.collect.size

q.rdd.count

【问题讨论】:

  • 您是如何获得这些带有附加参数(行数等)的 DAG 图形的,从未见过。这是 spark 2.3 中的新功能吗?
  • @JacekLaskowski 基准测试? :) 2 节点集群,10000000 行,10 次迭代 - 应该回答你的问题。您还可以启动 Java Mission Control 来测量 GC,在线代码编译
  • 只是预感,但q.count 似乎是这里唯一合理的选择。这是唯一可以应用源特定优化的方法(Daniel Darabos 的相关问题:stackoverflow.com/q/40629435/1560062)。 q.queryExecution.toRdd.count 可能很快(毕竟它只是一个带有可变累加器的天真 while,所以 JVM 应该喜欢它)但它完全不知道上下文。例如,如果你通过 JDBC 运行它,它只会获取所有行,而不是一堆。
  • @zero323 尽管如此。不在主界面上运行应用程序有什么好处吗?我认为我们不应该直接使用查询,而是通过数据集
  • 你们都想要基准测试的任何其他“计数”方法吗?列出的 4 个不同数据大小(1M、1B、1T)、方法、数据源(范围、镶木地板、文本文件)的编译时间

标签: performance apache-spark query-optimization apache-spark-sql


【解决方案1】:

我对@9​​87654322@做了一些测试:

  1. q.count:~50 毫秒
  2. q.collect.size:大约一分钟后我停止了查询……
  3. q.rdd.count: ~1100 毫秒
  4. q.queryExecution.toRdd.count: ~600 毫秒

一些解释:

选项 1 是迄今为止最快的,因为它同时使用部分聚合和整个阶段代码生成。整个阶段的代码生成让 JVM 变得非常聪明并进行一些剧烈的优化(参见:https://databricks.com/blog/2017/02/16/processing-trillion-rows-per-second-single-machine-can-nested-loop-joins-fast.html)。

选项 2。速度很慢,将驱动程序上的所有内容都具体化,这通常是个坏主意。

选项 3。类似于选项 4,但是这首先将内部行转换为常规行,这非常昂贵。

选项 4. 与不生成整个阶段代码的速度差不多。

【讨论】:

  • 顺便说一句,为什么选项 4. 没有 codegen? QueryExecution.preparations 物理计划准备规则在 toRdd(如 executedPlan)之前执行,因此 codegen 应该是选项 4 的一部分。直到今天我才探索该区域,所以我可能会混淆。我会很感激你的回答。谢谢。
  • 好的,让我改写一下。选项 4 是部分代码生成的。范围运算符是代码生成的,聚合不是。这有两个缺点:我们需要具体化一行,而且(更重要的是)JVM 很难优化它。
  • 答案的基本原理很好,我刚刚完成了针对不同文件类型和大小的一组基准的编译,您的基本原理似乎(大部分)遵循我在数据中看到的内容,我会当我可以使用我的个人电脑时发布我的发现
【解决方案2】:

格式化的东西很糟糕,哦,好吧

/*note: I'm using spark 1.5.2 so some of these might be different than what you might find in a newer version
 *note: These were all done using a 25 node , 40 core/node and started with --num-executors 64 --executor-cores 1 --executor-memory 4g
 *note: the displayed values are the mean from 10 runs
 *note: the spark-shell was restarted every time I noticed any spikes intra-run
 *
 *million/billion = sc.parallelize(1  to 1000000).toDF("col1")
 *
 *val s0 = sc.parallelize(1  to 1000000000)
 *//had to use this to get around maxInt constraints for Seq
 *billion10 = sc.union(s0,s1,s2,s3,s4,s5,s6,s7,s8,s9).toDF("col1")
 *
 *for parquet files
 *compression=uncompressed
 *written with:    million/billion/billion10.write.parquet
 *read with:    sqlContext.read.parquet
 *
 *for text files
 *written with:    million/billion/billion10.map(x=> x.mkString(",")).saveAsTextFile
 *read with:    sc.textFile.toDF("col1")
 *
 *excluded the collect() because that would have murdered my machine
 *made them all dataframes for consistency
/*


size       type     query         
billion10  text     count              81.594582
                    queryExecution     81.949047
                    rdd.count         119.710021
           Seq      count              18.768544
                    queryExecution     14.257751
                    rdd.count          36.404834
           parquet  count              12.016753
                    queryExecution     24.305452
                    rdd.count          41.932466
billion    text     count              14.120583
                    queryExecution     14.346528
                    rdd.count          22.240026
           Seq      count               2.191781
                    queryExecution      1.655651
                    rdd.count           2.831840
           parquet  count               2.004464
                    queryExecution      5.010546
                    rdd.count           7.815010
million    text     count               0.975095
                    queryExecution      0.113718
                    rdd.count           0.184904
           Seq      count               0.192044
                    queryExecution      0.029069
                    rdd.count           0.036061
           parquet  count               0.963874
                    queryExecution      0.217661
                    rdd.count           0.262279

观察:

  • 对于百万条记录,Seq 是最快的,但只有十分之一 秒数很难衡量整个集群的真实速度差异
  • 一般而言,将其存储为 TEXT 很慢,这可能是我阅读它的方式,我几乎只使用 Parquet,所以错过一些我很容易错过的东西
  • count 和 queryExecution 比 rdd.count 快 每个案例(正如赫尔曼在他的回答中解释的那样)
  • count 和 queryExecution 轮流更快,并且在数据类型之间有所不同:
    • 镶木地板的计数更快
    • 对于 Seq,queryExecution 更快
    • 对于文本,它们几乎相同
  • 随着尺寸的增加,速度是非线性的

如果有人想要不同的存储类型、计算不同的 dtype、压缩、更多列,请继续评论或给我发消息,我会看看我能做什么

【讨论】:

  • 谢谢詹姆斯。看起来这个问题需要检查所有不同的文件格式、数据源和大小。大量工作,当然取决于 Spark 的版本(因为每个版本都会发生变化)。
猜你喜欢
  • 2023-04-04
  • 1970-01-01
  • 1970-01-01
  • 2017-09-09
  • 1970-01-01
  • 1970-01-01
  • 2013-05-16
  • 1970-01-01
  • 2018-11-10
相关资源
最近更新 更多