【问题标题】:How does lineage get passed down in RDDs in Apache Spark沿袭如何在 Apache Spark 中的 RDD 中传递
【发布时间】:2015-06-07 23:33:06
【问题描述】:

每个 RDD 是否都指向同一个谱系图

当父 RDD 将其沿袭给新 RDD 时,子 RDD 是否也复制了沿袭图,因此父和子都有不同的图。在这种情况下不是内存密集型的吗?

【问题讨论】:

    标签: apache-spark rdd


    【解决方案1】:

    每个 RDD 都维护一个指向一个或多个父级的指针以及关于它与父级之间的关系类型的元数据。例如,当我们在 RDD 上调用 val b = a.map() 时,RDD b 只保留对其父级 a 的引用(并且从不复制),这是一个沿袭

    当驱动程序提交作业时,RDD 图被序列化到工作节点,以便每个工作节点在不同的分区上应用一系列转换(如映射过滤器等)。此外,如果发生某些故障,此 RDD 沿袭将用于重新计算数据。

    为了显示一个RDD的血统,Spark提供了一个调试方法toDebugString()方法。

    考虑以下示例:

    val input = sc.textFile("log.txt")
    val splitedLines = input.map(line => line.split(" "))
                        .map(words => (words(0), 1))
                        .reduceByKey{(a,b) => a + b}
    

    splitedLines RDD上执行toDebugString(),会输出如下,

    (2) ShuffledRDD[6] at reduceByKey at <console>:25 []
        +-(2) MapPartitionsRDD[5] at map at <console>:24 []
        |  MapPartitionsRDD[4] at map at <console>:23 []
        |  log.txt MapPartitionsRDD[1] at textFile at <console>:21 []
        |  log.txt HadoopRDD[0] at textFile at <console>:21 []
    

    有关 Spark 内部工作原理的更多信息,请阅读my another post

    【讨论】:

      【解决方案2】:

      当调用转换(映射或过滤器等)时,Spark 不会立即执行它,而是为每个转换创建一个沿袭。 沿袭将跟踪必须在该 RDD 上应用的所有转换, 包括它必须从中读取数据的位置。

      例如,考虑下面的例子

      val myRdd = sc.textFile("spam.txt")
      val filteredRdd = myRdd.filter(line => line.contains("wonder"))
      filteredRdd.count()
      

      sc.textFile() 和 myRdd.filter() 不会立即执行, 它只会在 RDD 上调用 Action 时执行——这里是 filtersRdd.count()。

      Action 用于将结果保存到某个位置或显示它。 RDD沿袭信息也可以通过filteredRdd.toDebugString(filteredRdd是这里的RDD)命令打印出来。 此外,DAG 可视化以非常直观的方式显示完整的图形,如下所示:

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-02-22
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多