【问题标题】:Spark and isolating time taken for tasks任务所花费的火花和隔离时间
【发布时间】:2020-06-30 16:57:52
【问题描述】:

我最近开始使用 Spark 来处理大量数据(~1TB)。并且也能够完成工作。但是我仍在尝试了解它的工作原理。考虑以下场景:

  1. 设置参考时间(比如tref

  2. 执行以下两项任务中的任何一项:

    一个。使用 SciSpark 将数万个文件中的大量数据 (~1TB) 读取到 RDD (OR) 中

    b.如上所述读取数据并进行额外的预处理工作并将结果存储在 DataFrame 中

  3. 打印适用的 RDD 或 DataFrame 的大小以及与 tref 的时间差(即 t0a/t0b子>)

  4. 做一些计算
  5. 保存结果

换句话说,1b 在处理与 1a 完全相同的 RDD 之后创建一个 DataFrame。

我的查询如下:

推断t0b – t0a = 预处理所需时间是否正确?我在哪里可以找到相同的可靠参考?

编辑:为问题的来源添加了解释......

我的怀疑源于 Spark 的惰性计算方法及其执行异步作业的能力。它可以/是否启动可以在读取数千个输入文件时计算的后续(预处理)任务?怀疑的根源在于令人难以置信的表现(结果验证正常)我认为这看起来太棒了,难以置信。

感谢您的回复。

【问题讨论】:

    标签: apache-spark performance-testing


    【解决方案1】:

    我相信这样的事情可以帮助你(使用 Scala):

    def timeIt[T](op: => T): Float = {
      val start = System.currentTimeMillis
      val res = op
      val end = System.currentTimeMillis
      (end - start) / 1000f
    }
    
    def XYZ = {
     val r00 = sc.parallelize(0 to 999999)
     val r01 = r00.map(x => (x,(x,x,x,x,x,x,x)))
     r01.join(r01).count()
    }
    
    val time1 = timeIt(XYZ)
    // or like this on next line
    //val timeN = timeIt(r01.join(r01).count())
    
    println(s"bla bla $time1 seconds.")
    

    您需要发挥创造力并逐步使用导致实际执行的操作。因此这是有局限性的。惰性求值等。

    另一方面,Spark Web UI 记录每个动作,并记录动作的阶段持续时间。

    一般来说:在共享环境中衡量性能是很困难的。在嘈杂的集群中,Spark 中的动态分配意味着您在 Stage 期间保留获取的资源,但在相同或下一个 Stage 的连续运行时,您可能会获得更少的资源。但这至少是指示性的,您可以在不那么繁忙的时段运行。

    【讨论】:

    • 谢谢@thebluephantom。我在编辑中添加了怀疑的原因。简而言之,虽然我们可以测量一个任务所需的时间,但这并不排除其他一些任务也同时运行的可能性。
    • 是的,这就是为什么我用斜体辅助。
    • 我的错!只是想明确表达我的担忧。
    • 共享环境中的性能测量很困难。在嘈杂的集群中动态分配火花意味着在工作人员期间保持资源,但可能会获得更少的资源。但这是指示性的......
    猜你喜欢
    • 2018-01-10
    • 1970-01-01
    • 2020-05-02
    • 1970-01-01
    • 2015-05-27
    • 1970-01-01
    • 2013-03-11
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多