【发布时间】:2019-01-04 21:13:50
【问题描述】:
我正在使用 2.1.1 版编写 Spark 应用程序。以下代码在调用带有 LocalDate 参数的方法时出错?
线程“主”java.lang.UnsupportedOperationException 中的异常:找不到 java.time.LocalDate 的编码器 - 字段(类:“java.time.LocalDate”,名称:“_2”) - 根类:“scala.Tuple2” 在 org.apache.spark.sql.catalyst.ScalaReflection$.org$apache$spark$sql$catalyst$ScalaReflection$$serializerFor(ScalaReflection.scala:602) 在 org.apache.spark.sql.catalyst.ScalaReflection$$anonfun$9.apply(ScalaReflection.scala:596) 在 org.apache.spark.sql.catalyst.ScalaReflection$$anonfun$9.apply(ScalaReflection.scala:587) 在 scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241) 在 scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241) 在 scala.collection.immutable.List.foreach(List.scala:381) 在 scala.collection.TraversableLike$class.flatMap(TraversableLike.scala:241) 在 scala.collection.immutable.List.flatMap(List.scala:344) 在 org.apache.spark.sql.catalyst.ScalaReflection$.org$apache$spark$sql$catalyst$ScalaReflection$$serializerFor(ScalaReflection.scala:587) ……val date : LocalDate = ....
val conf = new SparkConf()
val sc = new SparkContext(conf.setAppName("Test").setMaster("local[*]"))
val sqlContext = new org.apache.spark.sql.SQLContext(sc)
val itemListJob = new ItemList(sqlContext, jdbcSqlConn)
import sqlContext.implicits._
val processed = itemListJob.run(rc, priority).select("id").map(d => {
runJob.run(d, date)
})
class ItemList(sqlContext: org.apache.spark.sql.SQLContext, jdbcSqlConn: String) {
def run(date: LocalDate) = {
import sqlContext.implicits._
sqlContext.read.format("jdbc").options(Map(
"driver" -> "com.microsoft.sqlserver.jdbc.SQLServerDriver",
"url" -> jdbcSqlConn,
"dbtable" -> s"dbo.GetList('$date')"
)).load()
.select("id")
.as[Int]
}
}
更新:
我将runJob.run()的返回类型更改为元组(int, java.sql.Date)并将.map(...)的lambda中的代码更改为
val processed = itemListJob.run(rc, priority).select("id").map(d => {
val (a,b) = runJob.run(d, date)
$"$a, $b"
})
现在错误变成了
[错误] C:\....\scala\main.scala:40: 找不到存储在数据集中的类型的编码器。通过导入 spark.implicits 支持原始类型(Int、String 等)和产品类型(案例类)。未来版本中将添加对序列化其他类型的支持。 [错误] val 处理 = itemListJob.run(rc, priority).map(d => { [错误] ^ [错误] 发现一个错误 [错误] (compile:compileIncremental) 编译失败【问题讨论】:
-
请添加 spark 版本和使用的序列化(如果您更改默认值)。
-
我的本地开发 PC 上的 spark 版本是 2.1.1。我没有更改任何关于序列化的内容(默认设置)。
-
更改
runJob.run(d, date)以返回一些Spark SQL 可以理解的类,例如java.util.Date。 -
@zsxwing 谢谢,我按照您的建议更改了代码。但是,它现在得到了新的错误。我尝试在传递给
map()函数的 lambda 中添加import sqlContext.implicits._,但没有帮助。 -
import语句不应添加到 lambda 中,因为它将被map使用。只需将其添加到此行上方val processed = itemListJob.run(rc, priority).map(d => {。
标签: scala apache-spark apache-spark-dataset apache-spark-encoders