【发布时间】:2020-07-19 13:28:39
【问题描述】:
val query1 = spark.sql(s"select * from tmp_table where id IN (select id from tmp_table where date between '${fromDate}' and '$toDate' group by id having count(id) > 1)")
var finalTable = Seq[Template]()
query1.foreach { row =>
//Query on finalRtnView and get latest record for that row
//Update finalTable for that row
import spark.implicits._
finalTable = finalTable :+ finalTableTemplate
//Convert to DF so that data is updated .
val df = finalTable.toDF().createOrReplaceTempView("finalRtnView") //Throws Exception here
}
java.lang.NullPointerException 在 org.apache.spark.sql.SQLImplicits.localSeqToDatasetHolder(SQLImplicits.scala:213) 在 scala.collection.Iterator$class.foreach(Iterator.scala:893) 在 scala.collection.AbstractIterator.foreach(Iterator.scala:1336) 在 org.apache.spark.rdd.RDD$$anonfun$foreach$1$$anonfun$apply$28.apply(RDD.scala:918) 在 org.apache.spark.rdd.RDD$$anonfun$foreach$1$$anonfun$apply$28.apply(RDD.scala:918) 在 org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:2062) 在 org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:2062) 在 org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87) 在 org.apache.spark.scheduler.Task.run(Task.scala:108) 在 org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:335) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) 在 java.lang.Thread.run(Thread.java:748)
如果我在这里使用 for 循环而不是 foreach,则同样有效。在 for 循环的情况下,处理速度非常慢。
如何修改此代码以更快地运行?我每次在 for 循环中都需要更新的数据框。
【问题讨论】:
标签: scala foreach nullpointerexception apache-spark-sql scala-collections