【问题标题】:Execute Spark sql query within withColumn clause is Spark Scala在 withColumn 子句中执行 Spark sql 查询是 Spark Scala
【发布时间】:2021-11-09 13:07:02
【问题描述】:

我有一个数据框,其中有一个名为“Query”的列,其中存在 select 语句。想要执行此查询并创建一个包含 TempView 实际结果的新列。

+--------------+-----------+-----+----------------------------------------+
|DIFFCOLUMNNAME|DATATYPE   |ISSUE|QUERY                                   |
+--------------+-----------+-----+----------------------------------------+
|Firstname     |StringType |YES  |Select Firstname from TempView  limit 1 |
|LastName      |StringType |NO   |Select LastName from TempView  limit 1  |
|Designation   |StringType |YES  |Select Designation from TempView limit 1|
|Salary        |IntegerType|YES  |Select Salary from TempView    limit 1  |
+--------------+-----------+-----+----------------------------------------+

由于类型不匹配而出现错误,找到所需的字符串列。 我需要在这里使用UDF吗?但不确定如何编写和使用。请推荐

DF.withColumn("QueryResult", spark.sql(col("QUERY")))

TempView 是我创建的具有所有必需列的临时视图。 预期的最终 Dataframe 将是这样的,添加了新列 QUERYRESULT。

+--------------+-----------+-----+----------------------------------------+------------+
|DIFFCOLUMNNAME|DATATYPE   |ISSUE|QUERY                                   | QUERY RESULT
+--------------+-----------+-----+----------------------------------------+------------+
|Firstname     |StringType |YES  |Select Firstname from TempView  limit 1 | Bunny      |
|LastName      |StringType |NO   |Select LastName from TempView  limit 1  | Gummy      |
|Designation   |StringType |YES  |Select Designation from TempView limit 1| Developer  |
|Salary        |IntegerType|YES  |Select Salary from TempView    limit 1  | 100        |
+--------------+-----------+-----+----------------------------------------+------------+

【问题讨论】:

  • 显示一些代码供其他人查看。不寻常的构造
  • 我添加了代码,一行有 withColumn 子句。由于预期的错误是字符串和获取列
  • 简短的回答是“不,你不能那样做”。您可以做的是 pasha701s 回答中说明的解决方法:收集查询,以便它们在驱动程序中可用,然后逐个执行查询。但是,当数据无论如何都需要存在于驱动程序进程中时,为什么要将查询存储在 Spark 数据帧中呢?使用案例类列表而不是 Spark 数据框来保存查询可能会更容易。

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


【解决方案1】:

如果查询数量有限,您可以收集它们,执行每个查询,然后加入原始查询数据帧(Kieran 的回答更快,但我的回答有示例):

val queriesDF = Seq(
  ("Firstname", "StringType", "YES", "Select Firstname from TempView  limit 1 "),
  ("LastName", "StringType", "NO", "Select LastName from TempView  limit 1 "),
  ("Designation", "StringType", "YES", "Select Designation from TempView limit 1"),
  ("Salary", "IntegerType", "YES", "Select Salary from TempView limit 1 ")
).toDF(
  "DIFFCOLUMNNAME", "DATATYPE", "ISSUE", "QUERY"
)
val data = Seq(
  ("Bunny", "Gummy", "Developer", 100)
)
  .toDF("Firstname", "LastName", "Designation", "Salary")

data.createOrReplaceTempView("TempView")

// get all queries and evaluate results
val queries = queriesDF.select("QUERY").distinct().as(Encoders.STRING).collect().toSeq
val queryResults = queries.map(q => (q, spark.sql(q).as(Encoders.STRING).first()))
val queryResultsDF = queryResults.toDF("QUERY", "QUERY RESULT")

// Join original queries and results
queriesDF.alias("queriesDF")
  .join(queryResultsDF, Seq("QUERY"))
  .select("queriesDF.*", "QUERY RESULT")

输出:

+----------------------------------------+--------------+-----------+-----+------------+
|QUERY                                   |DIFFCOLUMNNAME|DATATYPE   |ISSUE|QUERY RESULT|
+----------------------------------------+--------------+-----------+-----+------------+
|Select Firstname from TempView  limit 1 |Firstname     |StringType |YES  |Bunny       |
|Select LastName from TempView  limit 1  |LastName      |StringType |NO   |Gummy       |
|Select Designation from TempView limit 1|Designation   |StringType |YES  |Developer   |
|Select Salary from TempView limit 1     |Salary        |IntegerType|YES  |100         |
+----------------------------------------+--------------+-----------+-----+------------+

【讨论】:

  • 不错!感谢您的说明
【解决方案2】:

假设您没有那么多“查询行”,只需使用 df.collect() 将结果收集到驱动程序,然后使用普通 Scala 映射查询。

【讨论】:

  • 这不是这个实际问题的答案。这是另一种选择 - 可能。
  • 当然,这不是当前措辞的答案,但它回答了他们问题的意图......
  • 他们在这里往往很挑剔,我也是。
  • 然后编辑问题并更具体。鉴于这个问题,他们可能需要更多信息。
  • 他们问题的第一行是“我有一个数据框,其中有一个名为“Query”的列,其中存在 select 语句。想要执行此查询并创建一个具有实际结果的新列TempView。”
猜你喜欢
  • 2018-12-26
  • 2018-11-20
  • 1970-01-01
  • 1970-01-01
  • 2021-07-13
  • 2023-04-03
  • 2020-04-02
  • 1970-01-01
  • 2016-07-11
相关资源
最近更新 更多