【问题标题】:Spark JDBC to Read and Write from and to HiveSpark JDBC 读写 Hive
【发布时间】: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


    【解决方案1】:

    有很多上下文,我不确定它是否相关。但是,错误很简单:

    Spark 无法在 Hive 中使用 DataType“Text”创建表。

    Hive 中确实没有称为 Text 的数据类型,也许您正在寻找以下之一:

    • 字符串
    • VARCHAR
    • 字符

    我希望这会有所帮助,否则请考虑将您的问题减少到最低限度(同时保持可重复的示例)以避免分心。

    【讨论】:

    • 这里的问题是,Spark 是根据 DF 的模式动态创建表的。是否有一种通用方法可以用 Hive 支持的数据类型替换 create table 数据类型?
    • 这听起来更像是处理多个数据库的挑战,而不是与 spark 相关的东西。我不知道有一种解决方案会自动将数据类型从任何数据库转换为另一个,但也不能说它不存在! -- 鉴于数据类型的数量有限,以合理的方式应用涵盖大多数情况的转换表应该不难。
    • 是的,看起来扩展现有的 HiveDialect 可以解决问题。这个链接讨论了相同link.medium.com/jZsQwXsMy1的实现,我会试一试并发布最终解决方案。
    • 更新了问题。现在收到“java.sql.SQLException:不支持的方法”错误。如果您知道任何相同的解决方案,请提出建议。
    • 你从哪里得到这个,如果知道你想出了什么解决方案,我真的很感激
    猜你喜欢
    • 1970-01-01
    • 2023-03-31
    • 2018-02-02
    • 2017-07-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-09-03
    相关资源
    最近更新 更多