【问题标题】:Scala + Spark - How to cast from 'scala.collection.mutable.WrappedArray$ofRef' to our custom object?Scala + Spark - 如何从“scala.collection.mutable.WrappedArray$ofRef”转换为我们的自定义对象?
【发布时间】:2020-09-01 16:48:25
【问题描述】:

我们在项目中使用sparkScala

我们正在使用自定义encoders 来创建spark datasets

我们的数据集架构类型如下:

(String, util.ArrayList[MyObject])

当我们执行 df.printSchema 时,我们得到以下信息:

root
 |-- key: string (nullable = true)
 |-- listOfMyObject: array (nullable = true)
 |    |-- element: binary (containsNull = true)

MyObject 是如下的 scala 案例类:

case class MyObject(key: String, dataList: java.util.ArrayList[MyObject2])

当我们在这个数据框上应用地图函数时,

df.map((row) => {
val key  = row.get(0)
val values = row.get(1)
})

在运行时,行包含以下架构:

StructField(key,StringType,true)
StructField(myObject,ArrayType(BinaryType,true),true)

我们能够检索到String 值,但是在尝试检索util.ArrayList[MyObject] 时,我们得到了scala.collection.mutable.WrappedArray$ofRef

我们使用ref.getClass 方法得到了这个。

有什么方法可以解决这个问题吗?

谢谢

阿努杰

【问题讨论】:

  • 能否提供更多代码?
  • 添加代码供参考
  • 请像这样val values: java.util.ArrayList[MyObject] = row.get(1) 投射它并让我知道结果
  • val 值:java.util.ArrayList[MyObject] = row.getAs[java.util.ArrayList[MyObject]](1) 由于阶段失败而中止作业:阶段 81.0 中的任务 123 失败 1次,最近的失败:在阶段 81.0 中丢失任务 123.0(TID 411,本地主机,执行程序驱动程序):java.lang.ClassCastException:scala.collection.mutable.WrappedArray$ofRef 无法转换为 java.util.ArrayList
  • 您需要提供一个可重现的示例。不清楚您遵循了哪些步骤以及您的代码是什么样的

标签: scala apache-spark apache-spark-sql scala-collections


【解决方案1】:

尝试使用spark可以通过编码器管理的Scala对象,在您的情况下,您可以将java.util.ArrayList更改为scala的简单List,然后如果列的顺序和类型正确,则转换仅使用以下简单行将您的数据框转换为您的类类型的数据集:

case class MyObject(key: String, dataList: List[MyObject2])
val ds: Dataset[MyObject] = df.as[MyObject]

请记住,如果架构有数组,则只能将其转换为列表,不能转换为数组。

【讨论】:

    【解决方案2】:

    我会将java.util.ArrayList 更改为scala.collection.immutable.List

    package playground
    
    import org.apache.spark.sql.SparkSession
    // import java.util.ArrayList
    
    object Casting {
    
      val spark = SparkSession.builder()
        .appName("Casting")
        .config("spark.master", "local[*]")
        .getOrCreate()
    
      case class MyObject2(fname: String, age: Int)
    
      // case class MyObject(key: String, dataList: ArrayList[MyObject2])
      case class MyObject(key: String, dataList: List[MyObject2])
    
      // val arrLst = new util.ArrayList[MyObject2]()
      // arrLst.add(MyObject2("Marie", 25))
      // arrLst.add(MyObject2("Peter", 27))
    
      val lst = List(MyObject2("Marie", 25),MyObject2("Peter", 27))
    
      // val data: Seq[MyObject] = Seq(MyObject("key",arrLst))
      val data: Seq[MyObject] = Seq(MyObject("key",lst))
    
      def main(args: Array[String]): Unit = {
    
        import spark.implicits._
    
        try {
          val df = spark.createDataFrame(data)
    
          df.printSchema()
          df.map(row => {
            val key  = row.key
            val values = row.dataList
            (key, values)
          }).show(false)
    
          df.show()
        } finally {
          spark.stop()
          println("Spark Session has stopped.")
        }
      }
    }
    
    root
     |-- key: string (nullable = true)
     |-- dataList: array (nullable = true)
     |    |-- element: struct (containsNull = true)
     |    |    |-- fname: string (nullable = true)
     |    |    |-- age: integer (nullable = false)
    
    +-----+--------------------------+
    |_1   |_2                        |
    +-----+--------------------------+
    |Lucia|[[Marie, 25], [Peter, 27]]|
    +-----+--------------------------+
    

    ArrayList 不是Scala collection 的成员。 我会尝试与Scala objects 合作,Spark 可以通过其encoders 管理。

    【讨论】:

    • 我做到了; val array = row.get(1).asInstanceOf[mutable.WrappedArray[MyObject]].toArray 出现以下异常:作业因阶段失败而中止:阶段 81.0 中的任务 123 失败 1 次,最近一次失败:阶段中丢失任务 123.0 81.0(TID 411,本地主机,执行程序驱动程序):java.lang.ArrayStoreException:[B
    • 而不是 toArray,尝试 toList。为什么要使用 ArrayList?是强制使用吗?
    • 我还能用什么来代替arraylist?
    • ArrayList 它不是 Scala 集合的成员,您可以尝试 List 例如。你试过 val array = arr.asInstanceOf[mutable.WrappedArray[MyObject]].toList.
    • 我没有看到,确实在您的帖子中,您的问题在于属于案例类 MyObject(key: String, dataList: java.util.ArrayList[ MyObject2]) 并且您可以检索字符串键。您的 DF 是扩展为 [String, WrappedArray] 的 Row[MyObject] 并且您得到 String 但不是 WrappedArray。首先,我认为您必须通过 List 或 Array 更改 java.util.ArrayList,然后调用 asInstanceOf 将 WrappedArray 转换为 List 或 Array。
    猜你喜欢
    • 2021-05-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-03-08
    • 1970-01-01
    • 2014-02-15
    • 1970-01-01
    • 2011-01-15
    相关资源
    最近更新 更多