【问题标题】:Why does RDD.foreach fail with "SparkException: This RDD lacks a SparkContext"?为什么 RDD.foreach 因“SparkException:此 RDD 缺少 SparkContext”而失败?
【发布时间】:2017-05-20 16:08:00
【问题描述】:

我有一个数据集(作为RDD),我使用不同的filter 运算符将其分成4 个RDD。

 val RSet = datasetRdd.
   flatMap(x => RSetForAttr(x, alLevel, hieDict)).
   map(x => (x, 1)).
   reduceByKey((x, y) => x + y)
 val Rp:RDD[(String, Int)] = RSet.filter(x => x._1.split(",")(0).equals("Rp"))
 val Rc:RDD[(String, Int)] = RSet.filter(x => x._1.split(",")(0).equals("Rc"))
 val RpSv:RDD[(String, Int)] = RSet.filter(x => x._1.split(",")(0).equals("RpSv"))
 val RcSv:RDD[(String, Int)] = RSet.filter(x => x._1.split(",")(0).equals("RcSv"))

我将RpRpSV发送到以下函数calculateEntropy

def calculateEntropy(Rx: RDD[(String, Int)], RxSv: RDD[(String, Int)]): Map[Int, Map[String, Double]] = {
        RxSv.foreach{item => {
               val string = item._1.split(",")
               val t = Rx.filter(x => x._1.split(",")(2).equals(string(2)))

        .
        .
    }
}

我有两个问题:

1- 当我在RxSv 上循环操作时:

RxSv.foreach{item=> { ... }}

它收集分区的所有项目,但我只想要我所在的分区。如果你说用户 map 功能但我没有更改 RDD 上的任何内容。

因此,当我在具有 4 个工作人员和一个驱动程序的集群上运行代码时,数据集分为 4 个分区,每个工作人员运行代码。但例如,我使用代码中指定的 foreach 循环。司机从工人那里收集所有数据。

2- 我在这段代码中遇到了问题

val t = Rx.filter(x => x._1.split(",")(2).equals(abc(2)))

错误:

org.apache.spark.SparkException: This RDD lacks a SparkContext.


它可能发生在以下情况:

(1) RDD transformationsactions 不是由驱动程序调用的,而是在其他转换内部;
例如,rdd1.map(x => rdd2.values.count() * x) 无效,因为值 transformationcount action 不能在 rdd1.map transformation 内部执行。有关详细信息,请参阅 SPARK-5063。

(2) 当Spark Streaming 作业从检查点恢复时,如果在DStream 操作中使用了对流作业未定义的RDD 的引用,则会触发此异常。有关详细信息,请参阅 SPARK-13​​758。

【问题讨论】:

  • 你能解释一下“但我只想在我所在的分区”?什么是“我所在的分区”
  • 我有 4 个分区,我想为每个分区运行代码块,而不从分区收集所有数据
  • 当您说“分区”时,您是指带有过滤行的各个 RDD?

标签: scala apache-spark rdd


【解决方案1】:

首先,我强烈建议使用 cache 运算符缓存第一个 RDD。

RSet.cache

这将避免每次您为其他 RDD filter 扫描和转换数据集:RpRcRpSvRcSv

引用cache的scaladoc:

cache() 以默认存储级别 (MEMORY_ONLY) 持久化此 RDD。

性能应该会提高。


其次,我会非常小心地使用术语“分区”来指代过滤的 RDD,因为该术语在 Spark 中具有特殊含义。

分区表示 Spark 为一个操作执行了多少任务。它们是 Spark 的提示,因此作为 Spark 开发人员的您可以微调您的分布式管道。

管道分布在集群节点上,每个分区方案具有一个或多个 Spark 执行器。如果你决定在一个 RDD 中有一个分区,一旦你在那个 RDD 上执行了一个动作,你就会在一个 executor 上执行一个任务。

filter 转换不会更改分区的数量(换句话说,它会保留分区)。分区数,即任务数,正好是RSet的分区数。


1- 当我在RxSv 上循环操作时,它会收集分区的所有项目,但我只想要我所在的分区

你是。不用担心,因为 Spark 会在数据所在的 executor 上执行任务。 foreach 是一个收集项目的操作,但描述了在执行程序上运行的计算,其数据分布在集群中(作为分区)。

如果您想在每个分区一次处理所有项目,请使用 foreachPartition:

foreachPartition 将函数 f 应用于此 RDD 的每个分区。


2- 我在这段代码中遇到了问题

在以下代码行中:

    RxSv.foreach{item => {
           val string = item._1.split(",")
           val t = Rx.filter(x => x._1.split(",")(2).equals(string(2)))

您正在执行foreach 操作,而该操作又使用Rx,即RDD[(String, Int)]。这是不允许的(如果可能的话不应该编译)。

出现这种行为的原因是,RDD 是一种数据结构,它仅描述执行操作时数据集发生的情况,并且存在于驱动程序(编排器)上。驱动程序使用数据结构来跟踪数据源、转换和分区数。

当驱动程序在执行程序上产生任务时,作为一个实体的 RDD 就消失了(=消失)。

当任务运行时,没有任何东西可以帮助他们知道如何运行作为他们工作一部分的 RDD。因此错误。 Spark 对此非常谨慎,会在任务执行后检查此类异常,以免它们引起问题。

【讨论】:

  • val dataset = sc.textFile("Dataset.txt", numberOfPart) 我设置了上面指定的分区数,我没有提到分区作为过滤项的数量。实际上,我的问题是我的结果会根据集群中的工作人员数量而变化。
  • 你能描述一下“我的问题是我的结果会根据集群中的worker数量而变化”你如何检查结果?你的期望是什么?你会得到什么?
猜你喜欢
  • 2018-01-04
  • 1970-01-01
  • 1970-01-01
  • 2015-09-27
  • 2019-11-09
  • 2012-04-06
  • 1970-01-01
  • 2021-01-17
  • 1970-01-01
相关资源
最近更新 更多