如何使用其架构从Spark数据框创建hive表?

lse*_*ohn 9 hive scala apache-spark

我想使用Spark数据帧的架构创建一个hive表.我怎样才能做到这一点?

对于固定列,我可以使用:

val CreateTable_query = "Create Table my table(a string, b string, c double)"
sparksession.sql(CreateTable_query) 
Run Code Online (Sandbox Code Playgroud)

但是我的数据框中有很多列,所以有没有办法自动生成这样的查询?

som*_*rti 17

假设您正在使用Spark 2.1.0或更高版本,而my_DF是您的数据帧,

//get the schema split as string with comma-separated field-datatype pairs
StructType my_schema = my_DF.schema();
String columns = Arrays.stream(my_schema.fields())
                       .map(field -> field.name()+" "+field.dataType().typeName())
                       .collect(Collectors.joining(","));

//drop the table if already created
spark.sql("drop table if exists my_table");
//create the table using the dataframe schema
spark.sql("create table my_table(" + columns + ") 
    row format delimited fields terminated by '|' location '/my/hdfs/location'");
    //write the dataframe data to the hdfs location for the created Hive table
    my_DF.write()
    .format("com.databricks.spark.csv")
    .option("delimiter","|")
    .mode("overwrite")
    .save("/my/hdfs/location");
Run Code Online (Sandbox Code Playgroud)

另一种使用临时表的方法

my_DF.createOrReplaceTempView("my_temp_table");
spark.sql("drop table if exists my_table");
spark.sql("create table my_table as select * from my_temp_table");
Run Code Online (Sandbox Code Playgroud)

  • 为什么我们需要创建临时表?my_DF.write.saveAsTable(...)有什么好处吗? (3认同)
  • /sf/ask/2146480591/ TL;DR saveastable 不会创建与配置单元兼容的表。问题是专门问蜂巢表所以... (3认同)
  • 我会将 field.dataType().typeName() 更改为 field.dataType().sql() 它可以更好地处理复杂/数组类型 (2认同)
  • **Scala** 翻译 `val tableColumns = df.schema.filter(_.name != partCol).map(field => field.name + " " + field.dataType.typeName).mkString(",")` (2认同)

小智 8

根据您的问题,您似乎希望使用数据框架构在hive中创建表.但正如您所说,在该数据框中有许多列,因此有两个选项

  • 第一是通过数据框架创建直接蜂巢表.
  • 第二个是采用这个数据帧的模式并在hive中创建表.

考虑以下代码:

package hive.example

import org.apache.spark.SparkConf
import org.apache.spark.SparkContext
import org.apache.spark.sql.SQLContext
import org.apache.spark.sql.Row
import org.apache.spark.sql.SparkSession

object checkDFSchema extends App {
  val cc = new SparkConf;
  val sc = new SparkContext(cc)
  val sparkSession = SparkSession.builder().enableHiveSupport().getOrCreate()
  //First option for creating hive table through dataframe 
  val DF = sparkSession.sql("select * from salary")
  DF.createOrReplaceTempView("tempTable")
  sparkSession.sql("Create table yourtable as select * form tempTable")
  //Second option for creating hive table from schema
  val oldDFF = sparkSession.sql("select * from salary")
  //Generate the schema out of dataframe  
  val schema = oldDFF.schema
  //Generate RDD of you data 
  val rowRDD = sc.parallelize(Seq(Row(100, "a", 123)))
  //Creating new DF from data and schema 
  val newDFwithSchema = sparkSession.createDataFrame(rowRDD, schema)
  newDFwithSchema.createOrReplaceTempView("tempTable")
  sparkSession.sql("create table FinalTable AS select * from tempTable")
}
Run Code Online (Sandbox Code Playgroud)


Val*_*ack 6

另一种方法是使用 StructType 上可用的方法.. sql , simpleString, TreeString 等...

您可以从 Dataframe 的架构创建 DDL,可以从您的 DDL 创建 Dataframe 的架构 ..

这是一个例子 - (直到 Spark 2.3)

    // Setup Sample Test Table to create Dataframe from
    spark.sql(""" drop database hive_test cascade""")
    spark.sql(""" create database hive_test""")
    spark.sql("use hive_test")
    spark.sql("""CREATE TABLE hive_test.department(
    department_id int ,
    department_name string
    )    
    """)
    spark.sql("""
    INSERT INTO hive_test.department values ("101","Oncology")    
    """)

    spark.sql("SELECT * FROM hive_test.department").show()

/***************************************************************/
Run Code Online (Sandbox Code Playgroud)

现在我有 Dataframe 可以玩了。在实际情况下,您会使用 Dataframe Readers 从文件/数据库创建数据帧。让我们使用它的模式来创建 DDL

  // Create DDL from Spark Dataframe Schema using simpleString function

 // Regex to remove unwanted characters    
    val sqlrgx = """(struct<)|(>)|(:)""".r
 // Create DDL sql string and remove unwanted characters

    val sqlString = sqlrgx.replaceAllIn(spark.table("hive_test.department").schema.simpleString, " ")

// Create Table with sqlString
   spark.sql(s"create table hive_test.department2( $sqlString )")
Run Code Online (Sandbox Code Playgroud)

从 Spark 2.4 开始,您可以在 StructType 上使用 fromDDL 和 toDDL 方法 -

val fddl = """
      department_id int ,
      department_name string,
      business_unit string
      """


    // Easily create StructType from DDL String using fromDDL
    val schema3: StructType = org.apache.spark.sql.types.StructType.fromDDL(fddl)


    // Create DDL String from StructType using toDDL
    val tddl = schema3.toDDL

    spark.sql(s"drop table if exists hive_test.department2 purge")

   // Create Table using string tddl
    spark.sql(s"""create table hive_test.department2 ( $tddl )""")

    // Test by inserting sample rows and selecting
    spark.sql("""
    INSERT INTO hive_test.department2 values ("101","Oncology","MDACC Texas")    
    """)
    spark.table("hive_test.department2").show()
    spark.sql(s"drop table hive_test.department2")

Run Code Online (Sandbox Code Playgroud)


小智 5

从 Spark 2.4 开始,您可以使用该函数来获取列名称和类型(即使对于嵌套结构)

val df = spark.read....

df.schema.toDDL
Run Code Online (Sandbox Code Playgroud)

  • 我在 pyspark 中找不到这个 - 只有 Scala 吗? (3认同)