【问题标题】:AWS glueContext read doesn't allow a sql queryAWS glueContext 读取不允许 sql 查询
【发布时间】:2019-01-08 14:55:34
【问题描述】:

我想使用 AWS 粘合作业从 Mysql 实例中读取过滤后的数据。由于胶水 jdbc 连接不允许我下推谓词,因此我试图在我的代码中显式创建 jdbc 连接。

我想使用如下所示的 jdbc 连接对 Mysql 数据库运行带有 where 子句的选择查询

import com.amazonaws.services.glue.GlueContext
import org.apache.spark.SparkContext
import org.apache.spark.sql.SparkSession


object TryMe {

  def main(args: Array[String]): Unit = {
    val sc: SparkContext = new SparkContext()
    val glueContext: GlueContext = new GlueContext(sc)
    val spark: SparkSession = glueContext.getSparkSession

    // Read data into a DynamicFrame using the Data Catalog metadata
    val t = glueContext.read.format("jdbc").option("url","jdbc:mysql://serverIP:port/database").option("user","username").option("password","password").option("dbtable","select * from table1 where 1=1").option("driver","com.mysql.jdbc.Driver").load()

  }
}

失败并出现错误

com.mysql.jdbc.exceptions.jdbc4.MySQLSyntaxErrorException 你有一个 SQL 语法错误;检查与您对应的手册 MySQL 服务器版本,用于在 'select * from 附近使用正确的语法 table1 where 1=1 WHERE 1=0' at line 1

这不应该吗?如何在不将整个表读入数据框的情况下使用 JDBC 连接检索过滤后的数据?

【问题讨论】:

    标签: aws-glue mssql-jdbc


    【解决方案1】:

    我认为问题的发生是因为您没有使用括号中的查询并提供别名。在我看来,它应该类似于以下示例:

     val t = glueContext.read.format("jdbc").option("url","jdbc:mysql://serverIP:port/database").option("user","username").option("password","password").option("dbtable","(select * from table1 where 1=1) as t1").option("driver","com.mysql.jdbc.Driver").load()
    

    有关 SQL 数据源中参数的更多信息:

    https://spark.apache.org/docs/latest/sql-data-sources-jdbc.html

    Glue 和 Glue 提供的框架也有“push_down_predicate”选项,但我只在基于 S3 的数据源上使用过该选项。我认为它不适用于 S3 和非分区数据以外的其他来源。

    https://docs.aws.amazon.com/glue/latest/dg/aws-glue-programming-etl-partitions.html

    【讨论】:

      【解决方案2】:

      对于仍在寻找更多答案/示例的任何人,我可以确认push_down_predicate 选项适用于 ODBC 数据源。以下是我从 SQL Server 读取数据的方式(使用 Python)。

      df = glueContext.read.format("jdbc")
          .option("url","jdbc:sqlserver://server-ip:port;databaseName=db;")
          .option("user","username")
          .option("password","password")
          .option("dbtable","(select t1.*, t2.name from dbo.table1 t1 join dbo.table2 t2 on t1.id = t2.id) as users")
          .option("driver","com.microsoft.sqlserver.jdbc.SQLServerDriver")
          .load()
      

      这也有效,但不像我预期的那样。谓词不会下推到数据源。

      df = glueContext.create_dynamic_frame.from_catalog(database = "db", table_name = "db_dbo_table1", push_down_predicate = "(id >= 2850700 AND statusCode = 'ACT')")
      

      pushDownPredicate 上的文档指出:启用或禁用谓词下推到 JDBC 数据源的选项。默认值为true,在这种情况下,Spark 会尽可能将过滤器下推到 JDBC 数据源。

      【讨论】:

        【解决方案3】:
        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多