最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SparkSQl簡介及運行原理

 更新時間:2021年08月10日 11:48:06   作者:山上有風景  
Spark SQL就是將SQL轉(zhuǎn)換成一個任務,提交到集群上運行,類似于Hive的執(zhí)行方式。今天通過本文給大家分享SparkSQl簡介及運行原理,感興趣的朋友跟隨小編一起看看吧

一:什么是SparkSQL?

(一)SparkSQL簡介

Spark SQL是Spark的一個模塊,用于處理結(jié)構(gòu)化的數(shù)據(jù),它提供了一個數(shù)據(jù)抽象DataFrame(最核心的編程抽象就是DataFrame),并且SparkSQL作為分布式SQL查詢引擎。
Spark SQL就是將SQL轉(zhuǎn)換成一個任務,提交到集群上運行,類似于Hive的執(zhí)行方式。

(二)SparkSQL運行原理

將Spark SQL轉(zhuǎn)化為RDD,然后提交到集群執(zhí)行。

(三)SparkSQL特點

(1)容易整合,Spark SQL已經(jīng)集成在Spark中

(2)提供了統(tǒng)一的數(shù)據(jù)訪問方式:JSON、CSV、JDBC、Parquet等都是使用統(tǒng)一的方式進行訪問

(3)兼容 Hive

(4)標準的數(shù)據(jù)連接:JDBC、ODBC

二:DataFrame

(一)什么是DataFrame?

在Spark中,DataFrame是一種以RDD為基礎的分布式數(shù)據(jù)集,類似于傳統(tǒng)數(shù)據(jù)庫中的二維表格。

DataFrame是組織成命名列的數(shù)據(jù)集。

它在概念上等同于關系數(shù)據(jù)庫中的表,但在底層具有更豐富的優(yōu)化。

關系型數(shù)據(jù)庫中的表由表結(jié)構(gòu)和數(shù)據(jù)組成,而DataFrame也類似,由schema(結(jié)構(gòu))和數(shù)據(jù)組成,其數(shù)據(jù)集是RDD。

DataFrame可以根據(jù)很多源進行構(gòu)建,包括:結(jié)構(gòu)化的數(shù)據(jù)文件,hive中的表,外部的關系型數(shù)據(jù)庫,以及RDD

補充:Spark中的RDD、DataFrame和DataSet講解

(一)Spark中的模塊

上圖展示了Spark的模塊及各模塊之間的關系:

底層是Spark-core核心模塊,Spark每個模塊都有一個核心抽象,Spark-core的核心抽象是RDD,

Spark SQL等都基于RDD封裝了自己的抽象,在Spark SQL中是DataFrame/DataSet。

相對來說RDD是更偏底層的抽象,DataFrame/DataSet是在其上做了一層封裝,做了優(yōu)化,使用起來更加方便。

從功能上來說,DataFrame/DataSet能做的事情RDD都能做,RDD能做的事情DataFrame/DataSet不一定能做。

(二)RDD和DataFrame的區(qū)別

DataFrame與RDD的主要區(qū)別在于:

DataFrame

DataFrame帶有schema元信息,即DataFrame所表示的二維表數(shù)據(jù)集的每一列都帶有名稱和類型。

使得Spark SQL得以洞察更多的結(jié)構(gòu)信息,從而對藏于DataFrame背后的數(shù)據(jù)源以及作用于DataFrame之上的變換進行了針對性的優(yōu)化,最終達到大幅提升運行時效率的目標。

RDD

RDD,由于無從得知所存數(shù)據(jù)元素的具體內(nèi)部結(jié)構(gòu),Spark Core只能在stage層面進行簡單、通用的流水線優(yōu)化。

DataFrame和RDD聯(lián)系:

DataFrame底層是以RDD為基礎的分布式數(shù)據(jù)集,和RDD的主要區(qū)別的是:RDD中沒有schema信息,而DataFrame中數(shù)據(jù)每一行都包含schema

DataFrame = RDD[Row] + shcema

三:SparkSession

(一)SparkSession簡介

SparkSession是Spark 2.0引如的新概念。SparkSession為用戶提供了統(tǒng)一的切入點,來讓用戶學習spark的各項功能。

在spark的早期版本中,SparkContext是spark的主要切入點,由于RDD是主要的API,我們通過sparkcontext來創(chuàng)建和操作RDD。

對于每個其他的API,我們需要使用不同的context。

例如,對于Streming,我們需要使用StreamingContext;對于sql,使用sqlContext;對于Hive,使用hiveContext。

但是隨著DataSet和DataFrame的API逐漸成為標準的API,就需要為他們建立接入點。所以在spark2.0中,引入SparkSession作為DataSet和DataFrame API的切入點。

SparkSession封裝了SparkConf、SparkContext和SQLContext。為了向后兼容,SQLContext和HiveContext也被保存下來。

(二)SparkSession實質(zhì)

SparkSession實質(zhì)上是SQLContext和HiveContext的組合(未來可能還會加上StreamingContext),所以在SQLContext和HiveContext上可用的API在SparkSession上同樣是可以使用的。

SparkSession內(nèi)部封裝了sparkContext,所以計算實際上是由sparkContext完成的。

(三)SparkSession特點

   ----為用戶提供一個統(tǒng)一的切入點使用Spark各項功能

----允許用戶通過它調(diào)用DataFrame和Dataset相關 API來編寫程序

----減少了用戶需要了解的一些概念,可以很容易的與Spark進行交互

----與Spark交互之時不需要顯示的創(chuàng)建SparkConf, SparkContext以及 SQlContext,這些對象已經(jīng)封閉在SparkSession中

四:通過RDD創(chuàng)建DataFrame

(一)通過樣本類創(chuàng)建(反射)

case class People(val name:String,val age:Int)  //可以聲明數(shù)據(jù)類型

object WordCount {
  def main(args:Array[String]):Unit={
    val conf = new SparkConf()
    //設置運行模式為本地運行,不然默認是集群模式
    //conf.setMaster("local")  //默認是集群模式
    //設置任務名
    conf.setAppName("WordCount").setMaster("local")
    conf.set("spark.default.parallelism","5")
    //設置SparkContext,是SparkCore的程序入口
    val sc = new SparkContext(conf)
    val Sqlsc = new SQLContext(sc)  //根據(jù)SparkContext生成SQLContext
    
    val array = Array("mark,14","kitty,23","dasi,45")
    val peopleRDD = sc.parallelize(array).map(line=>{    //生成RDD
      People(line.split(",")(0),line.split(",")(1).trim().toInt)
    })
    
    import Sqlsc.implicits._  //引入全部方法
    //將RDD轉(zhuǎn)換成DataFrame
    val df = peopleRDD.toDF()  
    //將DataFrame轉(zhuǎn)換成一個臨時的視圖
    df.createOrReplaceTempView("people")
    //使用SQL語句進行查詢
    Sqlsc.sql("select * from people").show()
  }
}

(二)通過SparkSession創(chuàng)建DataFrame

object WordCount {
  def main(args:Array[String]):Unit={
    val conf = new SparkConf()
    //設置運行模式為本地運行,不然默認是集群模式
    //conf.setMaster("local")  //默認是集群模式
    //設置任務名
    conf.setAppName("WordCount").setMaster("local")
    conf.set("spark.default.parallelism","5")
    //設置SparkContext,是SparkCore的程序入口
    val sc = new SparkContext(conf)
    val Sqlsc = new SQLContext(sc)  //根據(jù)SparkContext生成SQLContext
    
    val array = Array("mark,14","kitty,23","dasi,45")
    //1.需要將RDD數(shù)據(jù)映射成Row,需要引入import org.apache.spark.sql.Row
    val peopleRDD = sc.parallelize(array).map(line=>{    //生成RDD
      val fields = line.split(",")
      Row(fields(0),fields(1).trim().toInt)
    })
    
    //2.創(chuàng)建StructType定義結(jié)構(gòu)
    val st:StructType = StructType(
      //字段名,字段類型,是否可以為空
      List(  //傳參是列表類型,或者使用StructField("name", StringType, true) :: StructField("age", IntegerType, true) :: Nil來構(gòu)建列表
          StructField("name",StringType,true),
          StructField("age",IntegerType,true)
          )
    )
    
    //3.使用SparkSession建立DataFrame
    val df = Sqlsc.createDataFrame(peopleRDD,st)
    //將DataFrame轉(zhuǎn)換成一個臨時的視圖
    df.createOrReplaceTempView("people")
    //使用SQL語句進行查詢
    Sqlsc.sql("select * from people").show()
  }
}

(三)通過 json 文件創(chuàng)建DataFrames

[{"name":"dafa","age":12},{"name":"safaw","age":17},{"name":"ge","age":34}]
def main(args:Array[String]):Unit={
    val conf = new SparkConf()
    //設置運行模式為本地運行,不然默認是集群模式
    //conf.setMaster("local")  //默認是集群模式
    //設置任務名
    conf.setAppName("WordCount").setMaster("local")
    //設置SparkContext,是SparkCore的程序入口
    val sc = new SparkContext(conf)
    val Sqlsc = new SQLContext(sc)  //根據(jù)SparkContext生成SQLContext
    
    //通過json數(shù)據(jù)直接創(chuàng)建DataFrame
    val df = Sqlsc.read.json("E:\\1.json")
    
    //將DataFrame轉(zhuǎn)換成一個臨時的視圖
    df.createOrReplaceTempView("people1")
    //使用SQL語句進行查詢
    Sqlsc.sql("select * from people1").show()
  }

五:臨時視圖

(一)什么是視圖

視圖是一個虛表,跟Mysql里的概念是一樣的,視圖基于實際的表而存在,其實質(zhì)是一系列的查詢語句

(二)類型

局部視圖(Temoporary View):只在當前會話中有效,如果創(chuàng)建它的會話終止,則視圖也會消失。

全局視圖(Global Temporary View): 在全局范圍內(nèi)有效,不同的Session中都可以訪問,生命周期是Spark的Application運行周期,全局視圖會綁定到系統(tǒng)保留的數(shù)據(jù)庫global_temp中,因此使用它的時候必須加上相應前綴。

(三)創(chuàng)建視圖

創(chuàng)建局部視圖:df.createOrReplaceTempView("emp")
創(chuàng)建全局視圖:df.createOrReplaceGlobalTempView("empG")

(四)視圖查詢

spark.sql("select * from emp").show
spark.sql("select * from global_temp.empG").show  //查詢?nèi)忠晥D,需要添加前綴

(五)會話周期

spark.newSession.sql("select * from emp").show -----> 報錯,Table or View Not Found
spark.newSession.sql("select * from global_temp.empG").show ---->可以正常查詢

六:DataFrame的read和save和savemode

(一)數(shù)據(jù)讀取

val conf = new SparkConf().setAppName("TestDataFrame2").setMaster("local")
    val sc = new SparkContext(conf)
    val sqlContext = new SQLContext(sc)
    //方式一
    val df1 = sqlContext.read.json("E:\\666\\people.json")
    val df2 = sqlContext.read.parquet("E:\\666\\users.parquet")
    //方式二
    val df3 = sqlContext.read.format("json").load("E:\\666\\people.json")
    val df4 = sqlContext.read.format("parquet").load("E:\\666\\users.parquet")
    //方式三,默認是parquet格式
    val df5 = sqlContext.load("E:\\666\\users.parquet")
   //方式四,使用MySQL進行數(shù)據(jù)源讀取
    val url = "jdbc:mysql://192.168.123.102:3306/hivedb"
    val table = "dbs"
    val properties = new Properties()
    properties.setProperty("user","root")
    properties.setProperty("password","root")
    //需要傳入Mysql的URL、表明、properties(連接數(shù)據(jù)庫的用戶名密碼)
    val df = sqlContext.read.jdbc(url,table,properties)
    df.createOrReplaceTempView("dbs")
    sqlContext.sql("select * from dbs").show()

使用Hive作為數(shù)據(jù)源:需要在pom.xml文件中添加依賴

<!-- https://mvnrepository.com/artifact/org.apache.spark/spark-hive -->
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-hive_2.11</artifactId>
            <version>2.3.0</version>
        </dependency>

開發(fā)環(huán)境則把resource文件夾下添加hive-site.xml文件,集群環(huán)境把hive的配置文件要發(fā)到$SPARK_HOME/conf目錄下

<configuration>
        <property>
                <name>javax.jdo.option.ConnectionURL</name>
                <value>jdbc:mysql://localhost:3306/hivedb?createDatabaseIfNotExist=true</value>
                <description>JDBC connect string for a JDBC metastore</description>
                <!-- 如果 mysql 和 hive 在同一個服務器節(jié)點,那么請更改 hadoop02 為 localhost -->
        </property>
        <property>
                <name>javax.jdo.option.ConnectionDriverName</name>
                <value>com.mysql.jdbc.Driver</value>
                <description>Driver class name for a JDBC metastore</description>
        </property>
        <property>
                <name>javax.jdo.option.ConnectionUserName</name>
                <value>root</value>
                <description>username to use against metastore database</description>
        </property>
        <property>
                <name>javax.jdo.option.ConnectionPassword</name>
                <value>root</value>
        <description>password to use against metastore database</description>
        </property>
    <property>
                <name>hive.metastore.warehouse.dir</name>
                <value>/hive/warehouse</value>
                <description>hive default warehouse, if nessecory, change it</description>
        </property>  
</configuration>

hive-site.xml配置文件
val conf = new SparkConf().setMaster("local").setAppName(this.getClass.getSimpleName)
    val sc = new SparkContext(conf)
    val sqlContext = new HiveContext(sc)
    sqlContext.sql("select * from myhive.student").show() 

(二)數(shù)據(jù)保存

val conf = new SparkConf().setAppName("TestDataFrame2").setMaster("local")
    val sc = new SparkContext(conf)
    val sqlContext = new SQLContext(sc)
    val df1 = sqlContext.read.json("E:\\666\\people.json")
    //方式一
    df1.write.json("E:\\111")
    df1.write.parquet("E:\\222")
    //方式二
    df1.write.format("json").save("E:\\333")
    df1.write.format("parquet").save("E:\\444")
    //方式三
    df1.write.save("E:\\555")

(三)數(shù)據(jù)保存模式

df1.write.format("parquet").mode(SaveMode.Ignore).save("E:\\444")

七:數(shù)據(jù)集DataSet

Dataset也是一個分布式數(shù)據(jù)容器,簡單來說是類似二維表,Dataset里頭存有schema數(shù)據(jù)結(jié)構(gòu)信息和原生數(shù)據(jù),Dataset的底層封裝的是RDD,當RDD的泛型是Row類型的時候,我們也可以稱它為DataFrame。即Dataset<Row> = DataFrame。DataFrame是特殊的Dataset。

Spark整合了Dataset和DataFrame,前者是有明確類型的數(shù)據(jù)集,后者是無明確類型的數(shù)據(jù)集。根據(jù)官方的文檔:

Dataset是一種強類型集合,與領域?qū)ο笙嚓P,可以使用函數(shù)或者關系進行分布式的操作。
每個Dataset也有一個無類型的視圖,叫做DataFrame,也就是關于Row的Dataset。
簡單來說,Dataset一般都是Dataset[T]形式,這里的T是指數(shù)據(jù)的類型,如上圖中的Person,而DataFrame就是一個Dataset[Row]。

Datasets是懶加載的,即只有actions被調(diào)用的時候才會觸發(fā)計算。在內(nèi)部,Dataset代表一個邏輯計劃,用來描述產(chǎn)生數(shù)據(jù)需要的計算。當一個action被調(diào)用的時候,Spark的query優(yōu)化器會優(yōu)化這個邏輯計劃并以分布式的方式在物理上進行實際的計算操作。

(一)創(chuàng)建和使用DataSet---使用序列

(1,"Tom")  (2,"Mary")

測試數(shù)據(jù)

(1)定義case class
             case class MyData(a:Int,b:String)
(2)使用序列創(chuàng)建DataSet
             val DS = Seq(MyData(1,"Tom"),MyData(2,"Mary")).toDS

(二)創(chuàng)建和使用DataSet---通過case class作為編碼器,將DataFrame轉(zhuǎn)換成DataSet

(1)定義case class
             case class Person(name:String,age:BigInt)
(2)讀入JSON的數(shù)據(jù)
             val df = spark.read.json("/root/temp/people.json")
(3)將DataFrame轉(zhuǎn)換成DataSet
             val PersonDS =df.as[Person]

(三)創(chuàng)建和使用DataSet---讀取HDFS數(shù)據(jù)文件

(1)讀取HDFS的文件,直接創(chuàng)建DataSet
             val lineDS = spark.read.text("hdfs://bigdata111:9000/input/data.txt").as[String]
(2)分詞操作,查詢長度大于3的單詞
             val words = lineDS.flatMap(_.split(" ")).filter(_.length > 3)
             words.show
             words.collect

到此這篇關于SparkSQl簡介及運行原理的文章就介紹到這了,更多相關SparkSQl使用內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • SimpleDateFormat格式化日期問題

    SimpleDateFormat格式化日期問題

    這篇文章主要介紹了SimpleDateFormat格式化日期問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-06-06
  • spring boot中使用http請求的示例代碼

    spring boot中使用http請求的示例代碼

    本篇文章主要介紹了spring boot中 使用http請求的示例代碼,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-12-12
  • idea如何查看安裝插件的位置

    idea如何查看安裝插件的位置

    這篇文章主要介紹了idea如何查看安裝插件的位置問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-06-06
  • 解決mybatis-plus 查詢耗時慢的問題

    解決mybatis-plus 查詢耗時慢的問題

    這篇文章主要介紹了解決mybatis-plus 查詢耗時慢的問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-07-07
  • Java 鏈表的定義與簡單實例

    Java 鏈表的定義與簡單實例

    這篇文章主要介紹了 Java 鏈表的定義與簡單實例的相關資料,需要的朋友可以參考下
    2017-06-06
  • Springboot讀取配置文件及自定義配置文件的方法

    Springboot讀取配置文件及自定義配置文件的方法

    這篇文章主要介紹了Springboot讀取配置文件及自定義配置文件的方法,非常不錯,具有參考借鑒價值,需要的朋友可以參考下
    2017-12-12
  • Jenkins安裝和插件管理配置入門教程

    Jenkins安裝和插件管理配置入門教程

    這篇文章主要介紹了Jenkins安裝和插件管理知識,本文通過圖文并茂的形式給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2022-02-02
  • 深入了解JAVA HASHMAP的死循環(huán)

    深入了解JAVA HASHMAP的死循環(huán)

    HASHMAP基于哈希表的 Map 接口的實現(xiàn)。此實現(xiàn)提供所有可選的映射操作,并允許使用 null 值和 null 鍵。(除了非同步和允許使用 null 之外,HashMap 類與 Hashtable 大致相同。)下面小編來帶大家詳細了解下吧
    2019-06-06
  • Java8中Lambda表達式的理解與應用

    Java8中Lambda表達式的理解與應用

    Java8最值得學習的特性就是Lambda表達式和Stream?API,如果有python或者javascript的語言基礎,對理解Lambda表達式有很大幫助,下面這篇文章主要給大家介紹了關于Java8中Lambda表達式的相關資料,需要的朋友可以參考下
    2022-02-02
  • SpringBoot整合Retry實現(xiàn)錯誤重試過程逐步介紹

    SpringBoot整合Retry實現(xiàn)錯誤重試過程逐步介紹

    重試的使用場景比較多,比如調(diào)用遠程服務時,由于網(wǎng)絡或者服務端響應慢導致調(diào)用超時,此時可以多重試幾次。用定時任務也可以實現(xiàn)重試的效果,但比較麻煩,用Spring Retry的話一個注解搞定所有,感興趣的可以了解一下
    2023-02-02

最新評論

石河子市| 桐梓县| 泸溪县| 丰宁| 朝阳区| 阜南县| 杂多县| 万载县| 山东| 清新县| 长顺县| 凤城市| 金川县| 莲花县| 全州县| 新丰县| 万宁市| 五峰| 津市市| 中超| 平泉县| 遂昌县| 清流县| 蓬莱市| 蒙城县| 周宁县| 中阳县| 乌海市| 曲水县| 广元市| 金乡县| 蒙城县| 山东| 镇雄县| 都昌县| 包头市| 永靖县| 无极县| 浦县| 义乌市| 南投县|