【问题标题】:Spark SQL execution in scalascala中的Spark SQL执行
【发布时间】:2020-04-02 07:10:04
【问题描述】:

我有一个包含 SQL 查询和视图名称的以下数据(alldata)。

Select_Query|viewname
select v1,v2 from conditions|cond
select w1,w2 from locations|loca

我已拆分并将其正确分配给 temptable(alldata)

val Select_Querydf = spark.sql("select Select_Query,ViewName from alldata")

当我尝试执行查询并从中注册一个临时视图或表时,它显示空指针错误。但是当我注释掉 spark.sql stmt 时,PRINTLN 会显示表中的所有值。

 Select_Querydf.foreach{row => 
          val Selectstmt = row(0).toString()
          val viewname = row(1).toString()
          println(Selectstmt+"-->"+viewname)
      spark.sql(Selectstmt).registerTempTable(viewname)//.createOrReplaceTempView(viewname)
       }
output:
select v1,v2 from conditions-->cond
select w1,w2 from locations-->loca

但是当我用 spark.sql 执行它时,它显示以下错误,请帮助我哪里出错了。

19/12/09 02:43:12 错误执行程序:阶段 4.0 中任务 0.0 中的异常 (TID 4) java.lang.NullPointerException 在 org.apache.spark.sql.SparkSession.sessionState$lzycompute(SparkSession.scala:128) 在 org.apache.spark.sql.SparkSession.sessionState(SparkSession.scala:126) 在 org.apache.spark.sql.SparkSession.sql(SparkSession.scala:623) 在 sparkscalacode1.SQLQueryexecutewithheader$$anonfun$main$1.apply(SQLQueryexecutewithheader.scala:36) 在 sparkscalacode1.SQLQueryexecutewithheader$$anonfun$main$1.apply(SQLQueryexecutewithheader.scala:32) 在 scala.collection.Iterator$class.foreach(Ite​​rator.scala:891) 在 scala.collection.AbstractIterator.foreach(Ite​​rator.scala:1334) 在 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(未知来源) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(未知来源) 在 java.lang.Thread.run(Unknown Source) 19/12/09 02:43:12 错误 TaskSetManager:4.0阶段任务0失败1次;中止工作 线程“主”org.apache.spark.SparkException 中的异常:作业 由于阶段失败而中止:阶段 4.0 中的任务 0 失败 1 次,大多数 最近失败:在阶段 4.0 中丢失任务 0.0(TID 4,本地主机,执行程序 驱动程序):java.lang.NullPointerException 在 org.apache.spark.sql.SparkSession.sessionState$lzycompute(SparkSession.scala:128) 在 org.apache.spark.sql.SparkSession.sessionState(SparkSession.scala:126) 在 org.apache.spark.sql.SparkSession.sql(SparkSession.scala:623) 在 sparkscalacode1.SQLQueryexecutewithheader$$anonfun$main$1.apply(SQLQueryexecutewithheader.scala:36) 在 sparkscalacode1.SQLQueryexecutewithheader$$anonfun$main$1.apply(SQLQueryexecutewithheader.scala:32) 在 scala.collection.Iterator$class.foreach(Ite​​rator.scala:891) 在 scala.collection.AbstractIterator.foreach(Ite​​rator.scala:1334) 在 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(未知来源) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(未知来源) 在 java.lang.Thread.run(Unknown Source)

【问题讨论】:

  • 您根本无法这样做,因为您无法在执行程序中执行驱动程序的代码。执行者对火花会话和/或上下文一无所知,因此您会遇到异常。唯一可以使用spark.sql(Selectstmt).registerTempTable(viewname)//.createOrReplaceTempView(viewname) 之类的代码的地方是驱动程序
  • 谢谢亚历山德罗斯。但是你能建议我一种方法吗?我有一个包含 5 行 5 个不同 sql 查询的数据框。如果我必须一个接一个地执行它,该怎么做?但我希望它在驱动程序代码中。让我粘贴完整的代码。

标签: scala apache-spark


【解决方案1】:

这里的spark.sqlSparkSession 不能在Dataframe 的foreach 中使用。 Sparksession 在 Driver 中创建,foreach 在 worker 中执行而不是序列化。

我希望你有一个Select_Querydf 的小列表,如果有的话,你可以收集为一个列表并如下使用它。

Select_Querydf.collect().foreach { row =>
  val Selectstmt = row.getString(0)
  val viewname = row.getString(1)
  println(Selectstmt + "-->" + viewname)
  spark.sql(Selectstmt).createOrReplaceTempView(viewname)
}

希望这会有所帮助!

【讨论】:

  • 完美,它工作.. 新的火花,试图做一些自我项目。非常感谢..
猜你喜欢
  • 2021-11-09
  • 2015-12-03
  • 1970-01-01
  • 2020-05-22
  • 1970-01-01
  • 2018-11-20
  • 1970-01-01
  • 2017-06-12
  • 1970-01-01
相关资源
最近更新 更多