【问题标题】:Connect to Hive with jdbc driver in Spark在 Spark 中使用 jdbc 驱动程序连接到 Hive
【发布时间】:2021-12-06 14:09:24
【问题描述】:

我需要使用 Spark 将数据从远程 Hive 移动到本地 Hive。我尝试使用 JDBC 驱动程序连接到远程配置单元:'org.apache.hive.jdbc.HiveDriver'。我现在正在尝试从 Hive 中读取,结果是列值中的列标题而不是实际数据:

df = self.spark_session.read.format('JDBC') \
         .option('url', "jdbc:hive2://{self.host}:{self.port}/{self.database}") \
         .option('driver', 'org.apache.hive.jdbc.HiveDriver') \
         .option("user", self.username) \
         .option("password", self.password)
         .option('dbtable', 'test_table') \
         .load()
df.show()

结果:

+----------+
|str_column|
+----------+
|str_column|
|str_column|
|str_column|
|str_column|
|str_column|
+----------+

我知道 Hive JDBC 不是 Apache Spark 中的官方支持。但我已经找到了从其他不受支持的来源(例如 IMB Informix)下载的解决方案。也许有人已经解决了这个问题。

【问题讨论】:

    标签: apache-spark jdbc pyspark hive


    【解决方案1】:

    debug&trace代码后我们会发现问题出在JdbcDialect。没有HiveDialect所以spark会使用默认的JdbcDialect.quoteIdentifier。 所以你应该实现一个 HiveDialect 来解决这个问题:

    import org.apache.spark.sql.jdbc.JdbcDialect
    
    class HiveDialect extends JdbcDialect{
      override def canHandle(url: String): Boolean = 
        url.startsWith("jdbc:hive2")
      
    
      override def quoteIdentifier(colName: String): String = {
        if(colName.contains(".")){
          var colName1 = colName.substring(colName.indexOf(".") + 1)
          return s"`$colName1`"
        }
        s"`$colName`"
      }
    }
    

    然后通过以下方式注册方言:

    JdbcDialects.registerDialect(new HiveDialect)
    

    最后,像这样将选项 hive.resultset.use.unique.column.names=false 添加到 url 中

    option("url", "jdbc:hive2://bigdata01:10000?hive.resultset.use.unique.column.names=false")
    

    参考csdn blog

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-02-02
      • 2016-12-22
      • 1970-01-01
      • 2020-05-22
      • 1970-01-01
      • 2016-09-07
      • 1970-01-01
      • 2015-06-15
      相关资源
      最近更新 更多