【发布时间】:2020-03-08 05:02:31
【问题描述】:
我正在尝试提出一个通用实现,以使用 Spark JDBC 支持从/向各种 JDBC 兼容数据库(如 PostgreSQL、MySQL、Hive 等)读取/写入数据。
我的代码如下所示。
val conf = new SparkConf().setAppName("Spark Hive JDBC").setMaster("local[*]")
val sc = new SparkContext(conf)
val spark = SparkSession
.builder()
.appName("Spark Hive JDBC Example")
.getOrCreate()
val jdbcDF = spark.read
.format("jdbc")
.option("url", "jdbc:hive2://host1:10000/default")
.option("dbtable", "student1")
.option("user", "hive")
.option("password", "hive")
.option("driver", "org.apache.hadoop.hive.jdbc.HiveDriver")
.load()
jdbcDF.printSchema
jdbcDF.write
.format("jdbc")
.option("url", "jdbc:hive2://127.0.0.1:10000/default")
.option("dbtable", "student2")
.option("user", "hive")
.option("password", "hive")
.option("driver", "org.apache.hadoop.hive.jdbc.HiveDriver")
.mode(SaveMode.Overwrite)
输出:
root
|-- name: string (nullable = true)
|-- id: integer (nullable = true)
|-- dept: string (nullable = true)
上面的代码可以无缝地用于 PostgreSQL、MySQL 数据库,但是一旦我使用 Hive 相关的 JDBC 配置,它就会开始导致问题。
首先,我的读取无法读取任何数据并返回空结果。经过一番搜索,我可以通过添加自定义 HiveDialect 来进行读取,但我仍然面临将数据写入 Hive 的问题。
case object HiveDialect extends JdbcDialect {
override def canHandle(url: String): Boolean = url.startsWith("jdbc:hive2")
override def quoteIdentifier(colName: String): String = s"`$colName`"
override def getJDBCType(dt: DataType): Option[JdbcType] = dt match {
case StringType => Option(JdbcType("STRING", Types.VARCHAR))
case _ => None
}
}
JdbcDialects.registerDialect(HiveDialect)
写入错误:
19/11/13 10:30:14 ERROR Executor: Exception in task 0.0 in stage 1.0 (TID 1)
java.sql.SQLException: Method not supported
at org.apache.hive.jdbc.HivePreparedStatement.addBatch(HivePreparedStatement.java:75)
at org.apache.spark.sql.execution.datasources.jdbc.JdbcUtils$.savePartition(JdbcUtils.scala:664)
at org.apache.spark.sql.execution.datasources.jdbc.JdbcUtils$$anonfun$saveTable$1.apply(JdbcUtils.scala:834)
at org.apache.spark.sql.execution.datasources.jdbc.JdbcUtils$$anonfun$saveTable$1.apply(JdbcUtils.scala:834)
at org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1$$anonfun$apply$28.apply(RDD.scala:935)
at org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1$$anonfun$apply$28.apply(RDD.scala:935)
at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:2101)
at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:2101)
如何使用 Spark JDBC 执行从 Spark 到 多个远程 Hive 服务器的 Hive 查询(读/写)?
我不能使用 Hive 元存储 URI 方法,因为在这种情况下,我将使用单个 Hive 服务器配置来限制自己。另外,正如我之前提到的,我希望该方法对所有数据库类型(PostgreSQL、MySQL、Hive)都是通用的,因此在我的情况下采用 Hive 元存储 URI 方法将不起作用。
依赖详情:
- Scala 版本:2.11
- Spark 版本:2.4.3
- Hive 版本:2.1.1。
- 使用的 Hive JDBC 驱动程序:2.0.1
【问题讨论】:
标签: apache-spark hadoop jdbc hive