【发布时间】:2016-01-09 02:47:53
【问题描述】:
当我尝试对 Spark-SQL 外部数据源进行一些测试时会发生此问题。
我以两种方式构建数据框,并比较收集动作的速度。而且我发现如果列数太大,从外部数据源构建的数据框会滞后。我想知道这是否是 Spark-SQL 的外部数据源的限制。 :-)
为了更清楚地提出问题,我写了一段代码:
https://github.com/sunheehnus/spark-sql-test/
在我的 External Datasource API 基准代码中,它实现了一个假的外部数据源(实际上是一个 RDD[String, Array[Int]] ),并通过
val cmpdf = sqlContext.load("com.redislabs.test.dataframeRP", Map[String, String]())
然后我构建相同的RDD并通过
val rdd = sqlContext.sparkContext.parallelize(1 to 2048, 3)
val mappedrdd = rdd.map(x =>(x.toString, (x to x + colnum).toSeq.toArray))
val df = mappedrdd.toDF()
val dataColExpr = (0 to colnum).map(_.toString).zipWithIndex.map { case (key, i) => s"_2[$i] AS `$key`" }
val allColsExpr = "_1 AS instant" +: dataColExpr
val df1 = df.selectExpr(allColsExpr: _*)
当我运行测试代码时,我可以看到结果(在我的笔记本电脑上):
9905
21427
但是当我把column less(512),我可以看到结果:
4323
2221
看起来问题是如果Schema中的列数很少,External Datasource API会受益,但是随着Schema中列数的增长,External Datasource API最终会落后......我想知道如果这是 Spark-SQL 对外部数据源 API 的限制,还是我以错误的方式使用 API? 非常感谢。 :-)
【问题讨论】:
标签: apache-spark apache-spark-sql