【问题标题】:Filtering RDD with CustomObject, Type Mismatch使用 CustomObject 过滤 RDD,类型不匹配
【发布时间】:2017-06-25 04:26:39
【问题描述】:

我有这个自定义的 Scala 对象(基本上是一个 Java POJO):

object CustomObject {

  implicit object Mapper extends JavaBeanColumnMapper[CustomObject]

}


class CustomObject extends Serializable {


  @BeanProperty
  var amount: Option[java.lang.Double] = _

  ...
}

在我的主类中,我加载了一个包含这些 CustomObjects 的 RDD。 我正在尝试过滤它们并创建一个新的 RDD,其中仅包含数量 > 5000 的对象。

val customObjectRDD = sc.objectFile[CustomObject]( "objectFiles" )
val filteredRdd = customObjectRDD.filter( x => x.amount > 5000 )
println( filteredRdd.count() )

但是,我的编辑说

类型不匹配:预期:(CustomObject) => 布尔值,实际: (CustomObject) => 任意

我必须做些什么才能让它工作?

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    > 运算符未在 Option[Double] 上定义,您的过滤谓词将需要处理 Option

    scala> case class A(amount: Option[Double])
    defined class A
    
    scala> val myRDD = sc.parallelize(Seq(A(Some(10000d)), A(None),  A(Some(5001d)), A(Some(5000d))))
    myRDD: org.apache.spark.rdd.RDD[A] = ParallelCollectionRDD[12] at parallelize at <console>:29
    
    scala> myRDD.filter(_.amount.exists(_ > 5000)).foreach{println}
    A(Some(10000.0))
    A(Some(5001.0))
    

    这假定任何带有amount = None 的对象都应该使过滤谓词失败。有关Option.exists 的定义,请参见the docs

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-06-17
      • 1970-01-01
      • 1970-01-01
      • 2019-04-09
      • 2019-01-09
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多