【问题标题】:Unable to run spark map function that reads a Tuple RDD and returns a Tuple RDD无法运行读取元组 RDD 并返回元组 RDD 的 spark map 函数
【发布时间】:2017-07-10 17:09:24
【问题描述】:

我需要从另一个配对的 RDD 生成配对的 RDD。基本上,我正在尝试编写一个执行以下操作的地图函数。

RDD[Polygon,HashSet[Point]] => RDD[Polygon,Integer]

这是我写的代码:

迭代 HashSet 并从“Point”对象中累加一个值的 Scala 函数。

def outCountPerCell( jr: Tuple2[Polygon,HashSet[Point]] ) : Tuple2[Polygon,Integer] = {
  val setIter = jr._2.iterator()
  var outageCnt: Int = 0
  while(setIter.hasNext()) {
    outageCnt += setIter.next().getCoordinate().getOrdinate(2).toInt
  }
  return Tuple2(jr._1,Integer.valueOf(outageCnt))
}

在配对的 RDD 上应用该函数,这会引发错误:

scala> val mappedJoinResult = joinResult.map((t: Tuple2[Polygon,HashSet[Point]]) => outCountPerCell(t))
<console>:82: error: type mismatch;
found   : ((com.vividsolutions.jts.geom.Polygon, java.util.HashSet[com.vividsolutions.jts.geom.Point])) => (com.vividsolutions.jts.geom.Polygon, Integer)
required: org.apache.spark.api.java.function.Function[(com.vividsolutions.jts.geom.Polygon, java.util.HashSet[com.vividsolutions.jts.geom.Point]),?]
       val mappedJoinResult = joinResult.map((t: Tuple2[Polygon,HashSet[Point]]) => outCountPerCell(t))

有人可以看看我缺少什么,或者分享任何在 map() 操作中使用自定义函数的示例代码。

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    这里的问题是 joinResult 是来自 Java API 的 JavaPairRDD。此数据结构的 map 期望 Java 类型 lambdas (Function) 不能(至少是微不足道地)与 Scala lambdas 互换。

    因此有两种解决方案:尝试将给定方法转换为 Java Function 以传递给 map,或者按照开发人员的意图简单地使用 Scala RDD:

    设置虚拟数据

    在这里,我创建了一些备用类,并制作了一个与 OP 结构相似的 Java RDD:

    scala> case class Polygon(name: String)
    defined class Polygon
    
    scala> case class Point(ordinate: Int)
    defined class Point
    
    scala> :pa
    // Entering paste mode (ctrl-D to finish)
    
    /* More idiomatic method */
    def outCountPerCell( jr: (Polygon,java.util.HashSet[Point])) : (Polygon, Integer) =
    {
        val count = jr._2.asScala.map(_.ordinate).sum
        (jr._1, count)
    }
    
    // Exiting paste mode, now interpreting.
    
    outCountPerCell: (jr: (Polygon, java.util.HashSet[Point]))(Polygon, Integer)
    
    scala> val hs = new java.util.HashSet[Point]()
    hs: java.util.HashSet[Point] = []
    
    scala> hs.add(Point(2))
    res13: Boolean = true
    
    scala> hs.add(Point(3))
    res14: Boolean = true
    
    scala> val javaRDD = new JavaPairRDD(sc.parallelize(Seq((Polygon("a"), hs))))
    javaRDD: org.apache.spark.api.java.JavaPairRDD[Polygon,java.util.HashSet[Point]] = org.apache.spark.api.java.JavaPairRDD@14fc37a
    

    使用 Scala RDD

    可以使用.rdd从Java RDD中检索底层Scala RDD:

    scala> javaRDD.rdd.map(outCountPerCell).foreach(println)
    (Polygon(a),5)
    

    更好的是,使用 mapValues 和 Scala RDD

    由于只有元组的第二部分发生了变化,这个问题可以用.mapValues 彻底解决:

    scala> javaRDD.rdd.mapValues(_.asScala.map(_.ordinate).sum).foreach(println)
    (Polygon(a),5)
    

    【讨论】:

    • 谢谢evan058。我通过删除 tuple2 尝试了上述方法。我能够保存该功能。当试图调用该函数时,它会抛出以下错误。 scala&gt; joinResult.map(outCountPerCell(_)).foreach(println) &lt;console&gt;:82: error: missing parameter type for expanded function ((x$1) =&gt; outCountPerCell(x$1)) joinResult.map(outCountPerCell(_)).foreach(println)
    • 我的输入 Paired RDD 有一个 java.util.HashSet 而不是 scala 的 HashSet。你会导致这个问题吗? scala&gt; val joinResult = JoinQuery.SpatialJoinQuery(myPointsRDD,myPolygonRDD,true,true) joinResult: org.apache.spark.api.java.JavaPairRDD[com.vividsolutions.jts.geom.Polygon,java.util.HashSet[com.vividsolutions.jts.geom.Point]] = org.apache.spark.api.java.JavaPairRDD@3dc4e185
    • @IWonderHow 它对我有用java.util.HashSetjoinResult 的类型是什么?试试println(joinResult.getClass)
    • 它是一个javaPairRDD。 scala> println(joinResult.getClass) 类 org.apache.spark.api.java.JavaPairRDD
    • 感谢埃文,它成功了!非常感谢您抽出时间提供帮助。我必须使用 import scala.collection.JavaConverters._ 才能使用 .asScala
    猜你喜欢
    • 2016-06-23
    • 2018-09-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-10-31
    • 2015-09-09
    • 2017-08-16
    相关资源
    最近更新 更多