【问题标题】:Executing multiple SQL queries on Spark在 Spark 上执行多个 SQL 查询
【发布时间】:2018-03-11 05:08:04
【问题描述】:

我在 test.sql 文件中有一个 Spark SQL 查询 -

CREATE GLOBAL TEMPORARY VIEW VIEW_1 AS select a,b from abc

CREATE GLOBAL TEMPORARY VIEW VIEW_2 AS select a,b from VIEW_1

select * from VIEW_2

现在,我启动我的 spark-shell 并尝试像这样执行它 -

val sql = scala.io.Source.fromFile("test.sql").mkString
spark.sql(sql).show

这会失败并出现以下错误 -

org.apache.spark.sql.catalyst.parser.ParseException:
mismatched input '<' expecting {<EOF>, 'GROUP', 'ORDER', 'HAVING', 'LIMIT', 'OR', 'AND', 'WINDOW', 'UNION', 'EXCEPT', 'MINUS', 'INTERSECT', 'SORT', 'CLUSTER', 'DISTRIBUTE'}(line 1, pos 128)

我尝试在不同的 spark.sql 语句中一一执行这些查询,并且运行良好。问题是,我有 6-7 个查询来创建临时视图,最后我需要上一个视图的输出。有没有一种方法可以让我在单个 spark.sql 语句中运行这些 SQL。我曾研究过 Postgres SQL (Redshift),它能够执行此类查询。在 spark sql 中,在这种情况下我将不得不维护很多文件。

【问题讨论】:

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


    【解决方案1】:

    问题在于mkString 将所有行连接在一个字符串中,无法将其正确解析为有效的 SQL 查询。

    脚本文件中的每一行都应作为单独的查询执行,例如:

    scala.io.Source.fromFile("test.sql").getLines()
      .filterNot(_.isEmpty)  // filter out empty lines
      .foreach(query =>
        spark.sql(query).show
      )
    

    更新

    如果查询被拆分为多行,则情况会稍微复杂一些。

    我们绝对需要一个标记查询结束的标记。让它成为分号字符,就像在标准 SQL 中一样。

    首先,我们从源文件中收集所有非空行:

    val lines = scala.io.Source.fromFile(sqlFile).getLines().filterNot(_.isEmpty)
    

    然后我们处理收集的行,将每个新行与前一行连接起来,如果它不以分号结尾:

    val queries = lines.foldLeft(List[String]()) { case(queries, line) =>
      queries match {
        case Nil => List(line) // case for the very first line
        case init :+ last =>
          if (last.endsWith(";")) {
            // if a query ended on a previous line, we simply append the new line to the list of queries
            queries :+ line.trim
          } else {
            // the query is not terminated yet, concatenate the line with the previous one
            val queryWithNextLine = last + " " + line.trim
            init :+ queryWithNextLine
          }
      }
    }
    

    【讨论】:

    • 这可以工作,但问题是我有很大的查询,基本上查询本身被分成几行。因此,如果我执行上述操作,它不会正确解析查询。有没有办法可以将 1 行结尾定义为分号。
    • 我更新了答案,将查询拆分为多行。
    • 代码中的小错误。它应该下降;也。 Spark 不支持分号。其余的答案看起来不错,谢谢。
    • @Antot 真的很好,你为什么不尝试使用定义查询条件的 yml 文件...使用它我们可以动态地构造查询?这可能吗?
    猜你喜欢
    • 2016-07-11
    • 2018-11-20
    • 1970-01-01
    • 1970-01-01
    • 2016-02-06
    • 2017-04-07
    • 1970-01-01
    • 2017-04-12
    • 1970-01-01
    相关资源
    最近更新 更多