【问题标题】:reduce RDD having key as (String,String)减少键为 (String,String) 的 RDD
【发布时间】:2019-10-02 04:23:55
【问题描述】:

我想使用 ((String,String),BigDecimal) RDD 作为 PairRDD,这样我就可以使用 reduceByKey 函数。 Spark 不会将 RDD 识别为 PairRDD。有没有办法用RDD实现reduce功能。

scala> jrdd2
jrdd2: org.apache.spark.rdd.RDD[((String, String), java.math.BigDecimal)] = MapPartitionsRDD[33] at map at <console>:30

scala> val jrdd3 = jrdd2.reduceBykey((a,b)=>(a.add(b),1))
<console>:28: error: value reduceBykey is not a member of org.apache.spark.rdd.RDD[((String, String), java.math.BigDecimal)]
       val jrdd3 = jrdd2.reduceBykey((a,b)=>(a.add(b),1))

【问题讨论】:

  • 它的 .reduceByKey() 不是 .reduceBykey()
  • 感谢您发现错字。花了30分钟试图解决它。 :)

标签: scala apache-spark bigdata


【解决方案1】:

您的 reduceByKey 必须返回 BigDecimal - 而不是元组。试试这个:

val rdd = sc.parallelize(Seq((("a", "b"), new java.math.BigDecimal(2)), 
                             (("c", "d"), new java.math.BigDecimal(1)), 
                             (("a", "b"), new java.math.BigDecimal(2))))

rdd.reduceByKey(_.add(_))

返回

((c,d),1)
((a,b),4)

【讨论】:

  • 谢谢格伦尼。这只是方法名称中的一个愚蠢的错字。
  • 不,您的问题 - 正如此处发布的那样 - 不是只是一个错字。您返回的是 Tuple 而不是 BigDecimal,这将导致 value reduceBykey is not a member of org.apache.spark.rdd.RDD[((String, String), java.math.BigDecimal)] 错误。
  • 再次感谢。是的,我的代码有错字以及您指出的问题。我最终使用了 combineByKey 方法。
猜你喜欢
  • 2021-09-28
  • 1970-01-01
  • 2021-10-16
  • 2017-02-17
  • 2015-12-11
  • 2018-07-06
  • 2020-08-17
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多