【问题标题】:Spark/Scala iterator unable to assign variables defined outside of foreach loopSpark/Scala 迭代器无法分配在 foreach 循环之外定义的变量
【发布时间】:2017-07-31 19:39:58
【问题描述】:

请注意:虽然这个问题提到了 Spark (2.1),但我认为这实际上是一个 Scala (2.11) 问题,任何精通 Scala 的开发人员都能够回答它!


我有以下代码创建一个 Spark Dataset(基本上是一个 2D 表)并逐行迭代它。如果特定行的username 列的值为“fizzbuzz”,那么我想设置一个在迭代器外部定义的变量,并在行迭代完成后使用该变量:

val myDataset = sqlContext
     .read
     .format("org.apache.spark.sql.cassandra")
     .options(Map("table" -> "mytable", "keyspace" -> "mykeyspace"))
     .load()

var foobar : String
myDataset.collect().foreach(rec =>
  if(rec.getAs("username") == "fizzbuzz") {
    foobar = rec.getAs("foobarval")
  }
)

if(foobar == null) {
  throw new Exception("The fizzbuzz user was not found.")
}

当我运行它时,我得到以下异常:

error: class $iw needs to be abstract, since:
it has 2 unimplemented members.
/** As seen from class $iw, the missing signatures are as follows.
 *  For convenience, these are usable as stub implementations.
 */
  def foobar=(x$1: String): Unit = ???

class $iw extends Serializable {
      ^

我得到这个有什么特别的原因吗?

【问题讨论】:

    标签: scala apache-spark iterator


    【解决方案1】:

    在方法或非抽象类中,必须为每个变量定义一个值;在这里,您将 foobar 保留为未定义。如果您将其定义为具有null 的初始值,事情将按预期工作:

    var foobar: String = null
    

    但是:请注意,您的代码既非惯用代码(未遵循 Scala 和 Spark 的最佳实践),也可能存在风险/缓慢:

    • 您应该避免使用可变值,例如 foobar - 不可变代码更容易推理,并且可以真正让您利用 Scala 的强大功能
    • 您应该避免在 DataFrame 上调用 collect,除非您确定它非常小,因为 collect 会将来自工作节点(其中可能有很多)的所有数据收集到单个驱动程序节点中,这会很慢,可能会导致OutOfMemoryError
    • 不鼓励使用null(因为它通常会导致意外的NullPointerExceptions)

    此代码的更惯用版本将使用 DataFrame.filter 过滤相关记录,并可能使用 Option 正确表示可能为空的值,例如:

    import spark.implicits._
    
    val foobar: Option[String] = myDataset
      .filter($"username" === "fizzbuzz") // filter only relevant records
      .take(1) // get first 1 record (if it exists) as an Array[Row]
      .headOption // get the first item in the array, or None
      .map(r => r.getAs[String]("foobarval")) // get the value of the column "foobarval", or None
    
    if (foobar.isEmpty) {
      throw new Exception("The fizzbuzz user was not found.")
    }
    

    【讨论】:

      【解决方案2】:

      foobar 变量应该被初始化:

      var foobar: String = null
      

      这看起来也不对:

      foobar = rec.getAs("foobarval")
      

      应该是:

      foobar = rec.getAs[String]("foobarval")
      

      总的来说,这不是要走的路。它根本没有从 Spark 执行模型中受益。我会过滤并取而代之:

      myDataset.filter($"username" === "fizzbuzz").select("foobarval").take(1)
      

      【讨论】:

        【解决方案3】:

        您可能应该在数据框上使用过滤器和选择:

        import spark.sqlContext.implicits._
        
        val data = spark.sparkContext.parallelize(List(
          """{ "username": "none", "foobarval":"none" }""",
          """{ "username": "fizzbuzz", "foobarval":"expectedval" }"""))
        
        val df = spark.read.json(data)
        val foobar = df.filter($"username" === "fizzbuzz").select($"foobarval").collect.head.getString(0)
        

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 2013-07-03
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2015-10-09
          • 2019-08-14
          • 1970-01-01
          • 1970-01-01
          相关资源
          最近更新 更多