【发布时间】:2016-01-15 14:52:10
【问题描述】:
我有两个 RDD,它们都有两列作为 (K,V)。在这些 RDD 的源中,键出现在另一个之下,并且对于每一行,为键分配了不同且不同的值。创建 RDD 的文本文件在这篇文章的底部给出。
两个 RDD 中的键完全不同,我想根据它们的值加入两个 RDD,并尝试找出每对存在多少个共同值。例如我试图达到诸如 (1-5, 10) 之类的结果,这意味着来自 RDD1 的键值“1”和来自 RDD2 的键值“5”共享 10 个值。
我在一台具有 256 GB 内存和 72 个内核的机器上工作。一个文本文件为 500 MB,而另一个为 3 MB。
这是我的代码:
val conf = new SparkConf().setAppName("app").setMaster("local[*]").set("spark.shuffle.spill", "true")
.set("spark.shuffle.memoryFraction", "0.4")
.set("spark.executor.memory","128g")
.set("spark.driver.maxResultSize", "0")
val RDD1 = sc.textFile("\\t1.txt",1000).map{line => val s = line.split("\t"); (s(0),s(1))}
val RDD2 = sc.textFile("\\t2.txt",1000).map{line => val s = line.split("\t"); (s(1),s(0))}
val emp_newBC = sc.broadcast(emp_new.groupByKey.collectAsMap)
val joined = emp.mapPartitions(iter => for {
(k, v1) <- iter
v2 <- emp_newBC.value.getOrElse(v1, Iterable())
} yield (s"$k-$v2", 1))
joined.foreach(println)
val result = joined.reduceByKey((a,b) => a+b)
我尝试通过使用从脚本中看到的广播变量来管理这个问题。如果我加入 RDD2(有 250K 行),其自身对显示在相同的分区中,因此会发生更少的随机播放,因此需要 3 分钟才能获得结果。然而,当应用 RDD1 与 RDD2 时,配对分散在分区中,导致非常昂贵的洗牌过程,它总是最终给出
ERROR TaskSchedulerImpl: Lost executor driver on localhost: Executor heartbeat timed out after 168591 ms error.
根据我的结果:
我是否应该尝试对文本文件进行分区以在较小的块中创建 RDD1 并用 RDD2 分别加入那些较小的块?
-
还有其他方法可以根据它们的 Value 字段加入两个 RDD 吗?如果我将原始值描述为键并使用连接函数将它们连接起来,则值对再次分散在分区上,这再次导致非常昂贵的 reducebykey 操作。例如
val RDD1 = sc.textFile("\\t1.txt",1000).map{line => val s = line.split("\t"); (s(1),s(0))} val RDD2 = sc.textFile("\\t2.txt",1000).map{line => val s = line.split("\t"); (s(1),s(0))}RDD1.join(RDD2).map(line => (line._2,1)).reduceByKey((a,b) => (a+b))
伪数据样本:
KEY VALUE
1 13894
1 17376
1 15688
1 22434
1 2282
1 14970
1 11549
1 26027
1 2895
1 15052
1 20815
2 9782
2 3393
2 11783
2 22737
2 12102
2 10947
2 24343
2 28620
2 2486
2 249
2 3271
2 30963
2 30532
2 2895
2 13894
2 874
2 2021
3 6720
3 3402
3 25894
3 1290
3 21395
3 21137
3 18739
...
一个小例子
RDD1
2 1
2 2
2 3
2 4
2 5
2 6
3 1
3 6
3 7
3 8
3 9
4 3
4 4
4 5
4 6
RDD2
21 1
21 2
21 5
21 11
21 12
21 10
22 7
22 8
22 13
22 9
22 11
基于此数据连接结果:
(3-22,1)
(2-21,1)
(3-22,1)
(2-21,1)
(3-22,1)
(4-21,1)
(2-21,1)
(3-21,1)
(3-22,1)
(3-22,1)
(2-21,1)
(3-22,1)
(2-21,1)
(4-21,1)
(2-21,1)
(3-21,1)
减少关键结果:
(4-21,1)
(3-21,1)
(2-21,3)
(3-22,3)
【问题讨论】:
-
你应该展示一个使用 nano-data 的例子并告诉我 spark-default.conf 的内容是什么
-
我刚刚编辑了问题。
-
我在哪里工作 我们在一个集群中有 4 台 8GB 的计算机,我们能够毫无问题地读取 4GB 以上的文件,所以我敢打赌你的 $SPARK_HOME/conf/spark-defaults 有问题.conf,再加上你没有添加我要求的示例与 IN nano-data -> OUT nano-data(为了让我们更清楚)
-
我添加了一个小例子。是的 Alberto,我认为我的来源必须足以处理这些数据。我正在通过spark.apache.org/docs/latest/tuning.html 阅读有关调整火花的信息。我在那里遇到数据局部性问题。所以我的文本文件存在于不同的驱动程序中。我的意思是数据在 S 驱动程序中,而我的火花在 D 驱动程序中。这有关系吗?谢谢。
标签: scala join apache-spark