【问题标题】:Cache and Query a Dataset In Parallel Using Spark使用 Spark 并行缓存和查询数据集
【发布时间】:2017-12-11 22:58:11
【问题描述】:

我有一个要求,我想缓存一个数据集,然后通过在该数据集上并行触发“N”个查询来计算一些指标,所有这些查询都会计算相似的指标,只是过滤器会改变并且我想运行这些查询是并行的,因为响应时间很重要,而且我要缓存的数据集的大小总是小于 GB。

我知道如何在 Spark 中缓存数据集,然后随后对其进行查询,但是如果我必须对同一个数据集并行运行查询,我该如何实现呢?引入 alluxio 是一种方式,但在 Spark 世界中我们可以通过其他方式实现同​​样的效果吗?

例如使用Java,我可以将数据缓存在内存中,然后通过使用多线程我可以实现相同的效果,但是如何在Spark中做到这一点?

【问题讨论】:

  • 到目前为止您尝试过什么?你必须先尝试,然后才寻求帮助
  • 正如我在问题中提到的,我知道如何缓存数据集并在此之上执行查询,如果我知道方法,我需要一些关于如何在 Spark 中并行实现相同目标的指导/concept 使用,我早就这样做了
  • 默认情况下:查询可以在分布式数据集上并行运行,查询可以串行运行。现在,如果您想并行运行多个查询,那么您将不得不使用线程概念。 :)
  • 现在试试这个russellspitzer.com/2017/02/27/Concurrency-In-Spark,很快就会发布我的观察结果

标签: scala hadoop apache-spark


【解决方案1】:

使用 Scala 的并行集合在 Spark 的驱动程序代码中触发并行查询非常简单。这是一个最小的例子:

val dfSrc = Seq(("Raphael",34)).toDF("name","age").cache()


// define your queries, instead of returning a dataframe you could also write to a table etc
val query1: (DataFrame) => DataFrame = (df:DataFrame) => df.select("name")
val query2: (DataFrame) => DataFrame = (df:DataFrame) => df.select("age")

// Fire queries in parallel
import scala.collection.parallel.ParSeq
ParSeq(query1,query2).foreach(query => query(dfSrc).show())

编辑:

要在地图中收集查询 ID 和结果,您应该这样做:

val resultMap  = ParSeq(
 (1,query1), 
 (2,query2)
).map{case (queryId,query) => (queryId,query(dfSrc))}.toMap

【讨论】:

  • 太棒了,谢谢,这就是我要找的,不知道并行集合,将研究它们。
  • 如果查询总是返回单个值作为输出,不知道我们如何将所有查询的输出收集到一个 Map 中,其中 key 表示 query_id,value 表示查询输出(单个值)跨度>
  • @RajivChodisetti 然后使用地图而不是 foreach,例如ParSeq(query1,query2).map(query => (getQueryId(query),query(dfSrc))).toMap
  • 是的,想到了使用map,但是这里的查询是匿名函数,不知道如何为每个查询分配唯一标识符
  • 谢谢!! Scala 集合很棒,看来我需要更彻底地练习 n 学习 Scala 集合
猜你喜欢
  • 1970-01-01
  • 2020-04-03
  • 2019-04-01
  • 1970-01-01
  • 2022-01-14
  • 1970-01-01
  • 2012-03-02
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多