spark将mysql数据导入hive,pom等配置可看上一篇

import org.apache.spark.sql.{DataFrame, SaveMode, SparkSession}

object spark_from_mysql_to_hive {
  def main(args: Array[String]): Unit = {
    val spark: SparkSession = SparkSession.builder().master("local[*]").enableHiveSupport()
      .config("spark.debug.maxToStringFields", "100")
      .config("spark.sql.debug.maxToStringFields", "100")
      .config("hive.metastore.uris", "thrift://192.168.0.51:9083")
      //由于 Hive 和 SparkSQL 在 Decimal 类型上使用了不同的转换方式写入 Parquet,
      // 导致 Hive 无法正确读取 SparkSQL 所导入的数据。对于已有的使用 SparkSQL 导入的数据,
      // 如果有被 Hive/Impala 使用的需求,建议加上 spark.sql.parquet.writeLegacyFormat=true,重新导入数据。
      .config("spark.sql.parquet.writeLegacyFormat", true)
      .appName("mysql_to_hive").getOrCreate();
    import spark.implicits._
    /**
     * 使用临时表,建hive表,用insert into语句插入
     * 直接用saveAsTable生成hive表,之前要执行spark.table
     */
    val jdbcDF: DataFrame = spark.read.format("jdbc")
      .option("url","jdbc:mysql://192.168.0.54:3306/finance? characterEncoding=utf-8&serverTimezone=UTC&useSSL=false")
      .option("driver","com.mysql.jdbc.Driver")
      .option("dbtable","base_store")
      .option("user","root")
      .option("password","d0sD6Ffs7sGkHDF8mPnHJ2cl")
      .load();
//    spark.sql("create table finance.base_store( id int, student String) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t';");
    jdbcDF.show();
    jdbcDF.createOrReplaceTempView("temp");
    println("+"*500);
    spark.sql("select id,store_brand from temp").show();
    spark.sql("use finance");
//    spark.sql("set hive.stats.autogather=false");
    spark.sql("drop table if exists base_store");
    spark.sql("CREATE TABLE if not exists finance.base_store(id int,store_brand String)ROW FORMAT DELIMITED FIELDS TERMINATED BY ','");
    println("开始导入");
//    spark.sql("insert into finance.base_store select id,store_brand from temp");
    // 在saveAsTable之前要执行spark.table
    val df = spark.table("temp");
//    spark.conf.set("spark.sql.parquet.writeLegacyFormat", true);
    df.write.mode(SaveMode.Overwrite).saveAsTable("hive_records");
    spark.sql("select * from hive_records").show();
//    jdbcDF.write.format("hive").saveAsTable("base_store");
//    jdbcDF.write.mode(SaveMode.Overwrite).saveAsTable("base_store");
//    jdbcDF.write.mode("Overwrite").saveAsTable("base_store");

    println("导入完成");
    spark.stop();
  }

}

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐