【发布时间】: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"))
我将Rp和RpSV发送到以下函数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 transformations 和 actions 不是由驱动程序调用的,而是在其他转换内部;
例如,rdd1.map(x => rdd2.values.count() * x) 无效,因为值 transformation 和 count action 不能在 rdd1.map transformation 内部执行。有关详细信息,请参阅 SPARK-5063。
(2) 当Spark Streaming 作业从检查点恢复时,如果在DStream 操作中使用了对流作业未定义的RDD 的引用,则会触发此异常。有关详细信息,请参阅 SPARK-13758。
【问题讨论】:
-
你能解释一下“但我只想在我所在的分区”?什么是“我所在的分区”?
-
我有 4 个分区,我想为每个分区运行代码块,而不从分区收集所有数据
-
当您说“分区”时,您是指带有过滤行的各个 RDD?
标签: scala apache-spark rdd