【问题标题】:Spark RDD Persistence and PartitionsSpark RDD 持久化和分区
【发布时间】:2016-07-17 16:23:56
【问题描述】:

例如在 Spark 中创建某个 RDD 时:

lines = sc.textFile("README.md")

然后在这个 RDD 上调用一个转换:

pythonLines = lines.filter(lambda line: "Python" in line)

如果您在这个转换后的过滤器 RDD 上调用一个操作(例如 pythonlines.first),当他们说 an RDD will be recomputed ones again each time you run an action on them 时是什么意思?我认为在您对原始 RDD 调用 filter 转换后,您使用 textFile 方法创建的原始 RDD 不会保留。那么它会重新计算最近转换的 RDD,在这种情况下,它是我使用过滤器转换生成的 RDD?如果我的假设是正确的,我真的不明白为什么有必要这样做?

【问题讨论】:

    标签: hadoop apache-spark rdd bigdata


    【解决方案1】:

    在 spark 中,RDD 被惰性评估。这意味着如果你只是写

    lines = sc.textFile("README.md").map(xxx)
    

    您的程序将退出而不读取文件,因为您从未使用过结果。如果你写这样的东西:

    linesLength = sc.textFile("README.md").map(line => line.split(" ").length)
    sumLinesLength = linesLength.reduce(_ + _) // <-- scala way
    maxLineLength = linesLength.max()
    

    lineLength 所需的计算将进行两次,因为您在两个不同的地方重复使用它。为避免这种情况,您应该在以两种不同的方式使用之前持久化生成的 RDD

    linesLength = sc.textFile("README.md").map(line => line.split(" ").length)
    linesLength.persist()
    // ...
    

    您也可以查看https://spark.apache.org/docs/latest/programming-guide.html#rdd-persistence。希望我的解释不会太混乱!

    【讨论】:

    • 哦,所以在您的第二个示例中,当您在调用最后一行之后有这三行时:maxLineLength = linesLength.max() linesLength RDD 将消失,因为您已经完成了使用它。因此,如果您想在程序中的多个位置使用linesLength RDD,您应该将其持久化,这样您即使在使用linesLength RDD 之后也可以访问它。对吗?
    • 基本上,当您有一个中间结果(例如我的示例中的 lineLength)要多次重复使用时,您应该对其进行持久化(),否则 spark 也会对其进行多次计算。 RDD 不是数据,它是“计算列表”。所以如果你写linesLength.max(),Spark 不会理解 «你之前计算的最大值 » 但是«我可以通过读取这个文件并执行映射得到的 RDD 的最大值 »
    • 好的,现在我明白 RDD 将重新计算超时意味着什么了。感谢您的澄清!
    猜你喜欢
    • 2016-02-14
    • 2016-01-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-06-16
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多