【问题标题】:Upgrading to Spark 2.0 dataframe.map升级到 Spark 2.0 dataframe.map
【发布时间】:2017-03-18 10:45:53
【问题描述】:

我正在将一些 Spark 1.6 代码更新到 2.0.1,并且在使用 map 时遇到了一些问题。

我看到关于 SO 问题的其他问题,例如 encoder-error-while-trying-to-map-dataframe-row-to-updated-row,但我无法让这些技术发挥作用,而且对于下面的这种情况,它们似乎很荒谬。

val df = spark.sqlContext.read.parquet(inputFile)
df: org.apache.spark.sql.DataFrame = [device_id: string, hour: string ... 9 more fields]

val deviceAggDF = df.select("device_id").distinct
deviceAggDF: org.apache.spark.sql.Dataset[org.apache.spark.sql.Row] = [device_id: string]

deviceAggDF.map( x =>
  (
    Map("ID" -> x.getAs[String](0)),
    Map()
  )
)
scala.MatchError: Nothing (of class scala.reflect.internal.Types$ClassNoArgsTypeRef)
  at org.apache.spark.sql.catalyst.ScalaReflection$.schemaFor(ScalaReflection.scala:667)
  at org.apache.spark.sql.catalyst.ScalaReflection$.toCatalystArray$1(ScalaReflection.scala:448)
  at org.apache.spark.sql.catalyst.ScalaReflection$.org$apache$spark$sql$catalyst$ScalaReflection$$serializerFor(ScalaReflection.scala:482)
  at org.apache.spark.sql.catalyst.ScalaReflection$$anonfun$9.apply(ScalaReflection.scala:592)
  at org.apache.spark.sql.catalyst.ScalaReflection$$anonfun$9.apply(ScalaReflection.scala:583)
  at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
  at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
  at scala.collection.immutable.List.foreach(List.scala:381)
  at scala.collection.TraversableLike$class.flatMap(TraversableLike.scala:241)
  at scala.collection.immutable.List.flatMap(List.scala:344)
  at org.apache.spark.sql.catalyst.ScalaReflection$.org$apache$spark$sql$catalyst$ScalaReflection$$serializerFor(ScalaReflection.scala:583)
  at org.apache.spark.sql.catalyst.ScalaReflection$.serializerFor(ScalaReflection.scala:425)
  at org.apache.spark.sql.catalyst.encoders.ExpressionEncoder$.apply(ExpressionEncoder.scala:61)
  at org.apache.spark.sql.Encoders$.product(Encoders.scala:274)
  at org.apache.spark.sql.SQLImplicits.newProductEncoder(SQLImplicits.scala:47)

【问题讨论】:

标签: apache-spark elasticsearch-hadoop


【解决方案1】:

要返回空的Map,你应该指定可以被ecnoded的类型,例如:

deviceAggDF.map( x =>
  (
    Map("ID" -> x.getAs[String](0)),
    Map[String, String]()
  )
)

Map()Map[Nothing,Nothing],不能在Dataset 中使用。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-01-28
    • 1970-01-01
    • 1970-01-01
    • 2011-12-26
    • 2012-03-12
    • 2017-07-23
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多