【问题标题】:Iteratively running queries on Apache Spark在 Apache Spark 上迭代运行查询
【发布时间】:2016-11-25 18:38:16
【问题描述】:

我一直试图在一个相对较大的数据集 11M 上执行 10,000 个查询。更具体地说,我正在尝试基于一些 predicate 使用 filter 转换 RDD,然后通过应用 COUNT 计算有多少记录符合该过滤器行动。

我在具有 16GB 内存和 8 核 CPU 的本地计算机上运行 Apache Spark。我已将 --driver-memory 设置为 10G,以便将 RDD 缓存在内存中。

但是,由于我必须重新执行此操作 10,000 次,因此完成此操作需要异常长的时间。我还附上了我的代码,希望它能让事情更清楚。

加载查询和我要查询的数据框。

//load normalized dimensions
val df = spark.read.parquet("/normalized.parquet").cache()
//load query ranges
val rdd = spark.sparkContext.textFile("part-00000")

并行执行查询

在这里,我的查询收集在一个列表中,并使用 par 并行执行。然后我收集查询所需的参数,以过滤数据集。 isWithin 函数调用一个函数并测试我的数据集中包含的向量是否在我的查询给定的范围内。

现在过滤我的数据集后,我执行 count 以获取过滤后的数据集中存在的记录数,然后创建一个字符串报告有多少。

val results = queries.par.map(q => {
  val volume = q(q.length-1)
  val dimensions = q.slice(0, q.length-1)
  val count = df.filter(row => {
    val v = row.getAs[DenseVector]("scaledOpen")
    isWithin(volume, v, dimensions)
  }).count
  q.mkString(",")+","+count
})

现在,我想到的是,鉴于我拥有的大型数据集并试图在单台机器上运行这样的东西,这项任务通常非常困难。我知道这在 Spark 上运行的东西或通过使用索引可能会快得多。但是,我想知道是否有办法让它更快。

【问题讨论】:

    标签: performance scala apache-spark range


    【解决方案1】:

    仅仅因为您并行访问本地集合并不意味着任何事情都是并行执行的。可以并发执行的作业数量受集群资源而非驱动程序代码的限制。

    同时,Spark 专为高延迟批处理作业而设计。如果工作数量达到数万,您就无法期望事情会很快。

    您可以尝试的一件事是将过滤器下放到一个作业中。将DataFrame 转换为RDD

    import org.apache.spark.mllib.linalg.{Vector => MLlibVector}
    import org.apache.spark.rdd.RDD
    
    val vectors: RDD[org.apache.spark.mllib.linalg.DenseVector] = df.rdd.map(
      _.getAs[MLlibVector]("scaledOpen").toDense
    )
    

    map 向量到 {0, 1} 指标:

    import breeze.linalg.DenseVector
    
    // It is not clear what is the type of queries
    type Q = ???
    val queries: Seq[Q] = ???
    
    val inds: RDD[breeze.linalg.DenseVector[Long]] = vectors.map(v => {
      //  Create {0, 1} indicator vector
      DenseVector(queries.map(q => {
        // Define as before
        val volume = ???
        val dimensions = ???
    
        // Output 0 or 1 for each q
        if (isWithin(volume, v, dimensions)) 1L else 0L
      }): _*)
    })
    

    aggregate 部分结果:

    val counts: breeze.linalg.DenseVector[Long] = inds
      .aggregate(DenseVector.zeros[Long](queries.size))(_ += _, _ += _)
    

    并准备最终输出:

    queries.zip(counts.toArray).map {
      case (q, c) => s"""${q.mkString(",")},$c"""
    }
    

    【讨论】:

      猜你喜欢
      • 2015-09-30
      • 2017-03-11
      • 2015-05-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-10-17
      相关资源
      最近更新 更多