【发布时间】:2018-05-08 12:43:57
【问题描述】:
我有两个不同的 RDD,并对它们都应用了 foreach,并注意到我无法解决的差异。
第一个:
val data = Array(("CORN",6), ("WHEAT",3),("CORN",4),("SOYA",4),("CORN",1),("PALM",2),("BEANS",9),("MAIZE",8),("WHEAT",2),("PALM",10))
val rdd = sc.parallelize(data,3) // NOT sorted
rdd.foreach{ x => {
println (x)
}}
rdd: org.apache.spark.rdd.RDD[(String, Int)] = ParallelCollectionRDD[103] at parallelize at command-325897530726166:8
在这个意义上工作得很好。
第二个:
rddX.foreach{ x => {
val prod = x(0)
val vol = x(1)
val prt = counter
val cnt = counter * 100
println(prt,cnt,prod,vol)
}}
rddX: org.apache.spark.rdd.RDD[org.apache.spark.sql.Row] = MapPartitionsRDD[128] at rdd at command-686855653277634:51
工作正常。
问题:为什么我不能像第一个例子的第二种情况那样做 val prod = x(0) ?我怎么能用 foreach 做到这一点?或者我们是否需要在第一种情况下总是使用地图?由于第二个示例中的 Row 内部结构?
【问题讨论】:
标签: apache-spark