【问题标题】:How can I find the size of a RDD如何找到 RDD 的大小
【发布时间】:2015-10-02 13:37:51
【问题描述】:

我有RDD[Row],需要将其持久化到第三方存储库。 但是这个第三方存储库在一次调用中最多接受 5 MB。

所以我想根据 RDD 中存在的数据大小而不是基于 RDD 中存在的行数来创建分区。

如何找到RDD 的大小并根据它创建分区?

【问题讨论】:

    标签: apache-spark apache-spark-sql


    【解决方案1】:

    正如 Justin 和 Wang 所提到的,要获得 RDD 的大小并不容易。我们可以做一个估计。

    我们可以对一个 RDD 进行采样,然后使用SizeEstimator 来获取样本的大小。 正如王和贾斯汀所说, 基于离线采样的大小数据,例如,X 行使用 Y GB 离线,Z 行在运行时可能需要 Z*Y/X GB

    这是获取 RDD 大小/估计值的示例 scala 代码。

    我是 scala 和 spark 的新手。下面的示例可能会以更好的方式编写

    def getTotalSize(rdd: RDD[Row]): Long = {
      // This can be a parameter
      val NO_OF_SAMPLE_ROWS = 10l;
      val totalRows = rdd.count();
      var totalSize = 0l
      if (totalRows > NO_OF_SAMPLE_ROWS) {
        val sampleRDD = rdd.sample(true, NO_OF_SAMPLE_ROWS)
        val sampleRDDSize = getRDDSize(sampleRDD)
        totalSize = sampleRDDSize.*(totalRows)./(NO_OF_SAMPLE_ROWS)
      } else {
        // As the RDD is smaller than sample rows count, we can just calculate the total RDD size
        totalSize = getRDDSize(rdd)
      }
    
      totalSize
    }
    
    def getRDDSize(rdd: RDD[Row]) : Long = {
        var rddSize = 0l
        val rows = rdd.collect()
        for (i <- 0 until rows.length) {
           rddSize += SizeEstimator.estimate(rows.apply(i).toSeq.map { value => value.asInstanceOf[AnyRef] })
        }
    
        rddSize
    }
    

    【讨论】:

    • YARN 是如何获取 RDD 大小的?我正在运行作业并估计我的 RDD 大小(以 GB 为单位),但我无法在我的 Spark 代码中访问此信息。
    • 我发现 value.asInstanceOf[AnyRef] 比 toString 更好地估计 value.toString 如果 value 为 null 可以抛出一个空指针,并且似乎强制转换不会有这个问题,所以它也是更安全。
    • @SamuelAlexander rdd.sample(true, NO_OF_SAMPLE_ROWS) 将返回完整的 RDD,第二个参数应该是 0 到 1 之间的数字
    • @TheProletariat - 是的,您也可以使用您的代码编辑答案
    • 这个例子是错误的。作为@VictorP。指出,NO_OF_SAMPLE_ROWS 必须是小数
    【解决方案2】:

    一种直接的方法是调用following,取决于您是否要以序列化形式存储数据,然后转到spark UI“存储”页面,您应该能够计算出RDD的总大小(内存+磁盘):

    rdd.persist(StorageLevel.MEMORY_AND_DISK)
    
    or
    
    rdd.persist(StorageLevel.MEMORY_AND_DISK_SER)
    

    在运行时计算准确的内存大小并不容易。不过,您可以尝试在运行时进行估计:根据离线采样的大小数据,例如,X 行离线使用 Y GB,运行时 Z 行可能占用 Z*Y/X GB;这与 Justin 之前建议的类似。

    希望这会有所帮助。

    【讨论】:

    • 感谢您的回答。是的,这将有助于找到尺寸。但我想在我的管道/代码执行过程中检查这一点。所以手动签入 Spark UI 对我来说不是一个选项。
    • 我认为在运行时计算准确的内存大小并不容易。不过,您可以尝试在运行时进行估计:根据离线采样的大小数据,例如,X 行离线使用 Y GB,运行时 Z 行可能占用 Z*Y/X GB;这与 Justin 之前建议的类似。
    • 随机问题,当我执行 rdd.cache() 时,我在 UI 中看不到它。仅内存存储不显示?
    【解决方案3】:

    我认为 RDD.count() 会给你 RDD 中元素的数量

    【讨论】:

    • 你好@Yiying,欢迎来到 StackOverflow。发帖人要求的是 RDD 的大小,而不仅仅是行数。也许您可以扩展您的答案,以便海报不需要任何进一步的澄清。一旦你有足够的声望,如果你愿意,你就可以离开 cmets。
    • 据推测,该问题要求以信息单位(字节)为单位的大小。但count 也是衡量大小的标准——这个答案并没有真正回答问题,但确实为理想答案添加了信息。
    【解决方案4】:

    这将取决于序列化等因素,因此不切实际。但是,您可以获取一个样本集并对该样本数据进行一些实验,然后从那里进行推断。

    【讨论】:

    • 考虑我有一个包含字符串的 RDD。是否需要遍历所有 RDD 并使用 String.size() 来获取大小?
    • @sag 这是一种方法,但它会增加执行时间。如果你的 rdd 不是很大,你可以这样做。
    【解决方案5】:

    如果您实际上是在集群上处理大数据,则可以使用此版本 - 即它消除了收集。

    def calcRDDSize(rdd: RDD[Row]): Long = {
      rdd.map(_.mkString(",").getBytes("UTF-8").length.toLong)
         .reduce(_+_) //add the sizes together
    }
    
    def estimateRDDSize( rdd: RDD[Row], fraction: Double ) : Long = {
      val sampleRDD = rdd.sample(true,fraction)
      val sampleRDDsize = calcRDDSize(sampleRDD)
      println(s"sampleRDDsize is ${sampleRDDsize/(1024*1024)} MB")
    
      val sampleAvgRowSize = sampleRDDsize / sampleRDD.count()
      println(s"sampleAvgRowSize is $sampleAvgRowSize")
    
      val totalRows = rdd.count()
      println(s"totalRows is $totalRows")
    
      val estimatedTotalSize = totalRows * sampleAvgRowSize
      val formatter = java.text.NumberFormat.getIntegerInstance
      val estimateInMB = formatter.format(estimatedTotalSize/(1024*1024))
      println(s"estimatedTotalSize is ${estimateInMB} MB")
    
      return estimatedTotalSize
    }
    
    // estimate using 15% of data
    val size = estimateRDDSize(df.rdd,0.15)
    

    【讨论】:

    • 我认为可能有一个解决方案可以避免使用collect 来解决上述问题(你在这个解决方案中所做的),但仍然使用spark 的SizeEstimator.estimate,这可能比在行上运行 mkString 并查看字符串长度。大概这个答案只在它们存储为字符串时才有效,并且取决于RDD的持久化方式(序列化为字符串,序列化为Java对象等)
    猜你喜欢
    • 2015-01-09
    • 1970-01-01
    • 1970-01-01
    • 2011-05-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-30
    相关资源
    最近更新 更多