【问题标题】:spark map(func).cache slowspark map(func).cache 慢
【发布时间】:2014-07-15 05:06:16
【问题描述】:

当我使用cache 存储数据时,我发现spark 运行很慢。但是,当我不使用cache 方法时,速度非常好。我的主要配置文件如下:

SPARK_JAVA_OPTS+="-Dspark.local.dir=/home/wangchao/hadoop-yarn-spark/tmp_out_info 
-Dspark.rdd.compress=true -Dspark.storage.memoryFraction=0.4 
-Dspark.shuffle.spill=false -Dspark.executor.memory=1800m -Dspark.akka.frameSize=100 
-Dspark.default.parallelism=6"

而我的测试代码是:

val file = sc.textFile("hdfs://10.168.9.240:9000/user/bailin/filename")
val count = file.flatMap(line => line.split(" ")).map(word => (word, 1)).cache()..reduceByKey(_+_)
count.collect()

非常感谢任何有关我如何解决此问题的答案或建议。

【问题讨论】:

  • 我不认为这是一个可怕的问题,不确定投票的目的是什么。 cache 函数的使用并不总是很清楚。此外,这个问题写得也不错,而且很努力

标签: performance caching apache-spark


【解决方案1】:

cache 在您使用它的上下文中是无用的。在这种情况下,cache 表示将映射的结果 .map(word => (word, 1)) 保存在内存中。而如果你不调用它,reducer 可能会被链接到地图的末尾,并且地图结果在使用后会被丢弃。 cache 更适合在创建 RDD 后调用多个转换/操作的情况下使用。例如,如果你创建一个数据集,你想加入 2 个不同的数据集,缓存它会很有帮助,因为如果你在第二次加入时不加入,整个 RDD 将被重新计算。这是 spark 网站上的一个易于理解的示例。

val file = spark.textFile("hdfs://...")
val errors = file.filter(line => line.contains("ERROR")).cache() //errors is cached to prevent recalculation when the two filters are called
// Count all the errors
errors.count()
// Count errors mentioning MySQL
errors.filter(line => line.contains("MySQL")).count()
// Fetch the MySQL errors as an array of strings
errors.filter(line => line.contains("MySQL")).collect()

缓存在内部所做的是通过将 RDD 保存在内存中/保存到磁盘(取决于存储级别)来删除 RDD 的祖先,RDD 必须保存其祖先的原因是它可以根据需要重新计算,这是RDD的恢复方法。

【讨论】:

猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-04-20
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-01-24
  • 2023-04-08
相关资源
最近更新 更多