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

如何使用IDEA開發(fā)Spark SQL程序(一文搞懂)

 更新時間:2021年08月10日 11:12:11   作者:小哪吒的BD  
Spark SQL 是一個用來處理結構化數(shù)據(jù)的spark組件。它提供了一個叫做DataFrames的可編程抽象數(shù)據(jù)模型,并且可被視為一個分布式的SQL查詢引擎。這篇文章主要介紹了如何使用IDEA開發(fā)Spark SQL程序(一文搞懂),需要的朋友可以參考下

前言

大家好,我是DJ丶小哪吒,我又來跟你們分享知識了。對軟件開發(fā)有著濃厚的興趣。喜歡與人分享知識。做博客的目的就是為了能與 他 人知識共享。由于水平有限。博客中難免會有一些錯誤。如有 紕漏 之處,歡迎大家在留言區(qū)指正。小編也會及時改正。

DJ丶小哪吒又來與各位分享知識了。今天我們不飆車,今天就靜靜的坐下來,我們來聊一聊關于sparkSQL。準備好茶水,聽老朽與你娓娓道來。

Spark SQL是什么

Spark SQL 是一個用來處理結構化數(shù)據(jù)的spark組件。它提供了一個叫做DataFrames的可編程抽象數(shù)據(jù)模型,并且可被視為一個分布式的SQL查詢引擎。

1、使用IDEA開發(fā)Spark SQL

Spark會根據(jù)文件信息嘗試著去推斷DataFrame/DataSet的Schema,當然我們也可以手動指定,手動指定的方式有以下幾種:

  • 第1種:指定列名添加Schema
  • 第2種:通過StructType指定Schema
  • 第3種:編寫樣例類,利用反射機制推斷Schema

 1.1、指定列名添加Schema

package cn.itcast.sql

import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.{DataFrame, SparkSession}


object CreateDFDS {
  def main(args: Array[String]): Unit = {
    //1.創(chuàng)建SparkSession
    val spark: SparkSession = SparkSession.builder().master("local[*]").appName("SparkSQL").getOrCreate()
    val sc: SparkContext = spark.sparkContext
    sc.setLogLevel("WARN")
    //2.讀取文件
    val fileRDD: RDD[String] = sc.textFile("D:\\data\\person.txt")
    val linesRDD: RDD[Array[String]] = fileRDD.map(_.split(" "))
    val rowRDD: RDD[(Int, String, Int)] = linesRDD.map(line =>(line(0).toInt,line(1),line(2).toInt))
    //3.將RDD轉成DF
    //注意:RDD中原本沒有toDF方法,新版本中要給它增加一個方法,可以使用隱式轉換
    import spark.implicits._
    val personDF: DataFrame = rowRDD.toDF("id","name","age")
    personDF.show(10)
    personDF.printSchema()
    sc.stop()
    spark.stop()
  }
}

1.2、通過StructType指定Schema

package cn.itcast.sql

import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.types._
import org.apache.spark.sql.{DataFrame, Row, SparkSession}


object CreateDFDS2 {
  def main(args: Array[String]): Unit = {
    //1.創(chuàng)建SparkSession
    val spark: SparkSession = SparkSession.builder().master("local[*]").appName("SparkSQL").getOrCreate()
    val sc: SparkContext = spark.sparkContext
    sc.setLogLevel("WARN")
    //2.讀取文件
    val fileRDD: RDD[String] = sc.textFile("D:\\data\\person.txt")
    val linesRDD: RDD[Array[String]] = fileRDD.map(_.split(" "))
    val rowRDD: RDD[Row] = linesRDD.map(line =>Row(line(0).toInt,line(1),line(2).toInt))
    //3.將RDD轉成DF
    //注意:RDD中原本沒有toDF方法,新版本中要給它增加一個方法,可以使用隱式轉換
    //import spark.implicits._
    val schema: StructType = StructType(Seq(
      StructField("id", IntegerType, true),//允許為空
      StructField("name", StringType, true),
      StructField("age", IntegerType, true))
    )
    val personDF: DataFrame = spark.createDataFrame(rowRDD,schema)
    personDF.show(10)
    personDF.printSchema()
    sc.stop()
    spark.stop()
  }
}

1.3、反射推斷Schema–掌握

package cn.itcast.sql

import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.{DataFrame, SparkSession}


object CreateDFDS3 {
case class Person(id:Int,name:String,age:Int)
  def main(args: Array[String]): Unit = {
    //1.創(chuàng)建SparkSession
    val spark: SparkSession = SparkSession.builder().master("local[*]").appName("SparkSQL")
.getOrCreate()
    val sc: SparkContext = spark.sparkContext
    sc.setLogLevel("WARN")
    //2.讀取文件
    val fileRDD: RDD[String] = sc.textFile("D:\\data\\person.txt")
    val linesRDD: RDD[Array[String]] = fileRDD.map(_.split(" "))
    val rowRDD: RDD[Person] = linesRDD.map(line =>Person(line(0).toInt,line(1),line(2).toInt))
    //3.將RDD轉成DF
    //注意:RDD中原本沒有toDF方法,新版本中要給它增加一個方法,可以使用隱式轉換
    import spark.implicits._
    //注意:上面的rowRDD的泛型是Person,里面包含了Schema信息
    //所以SparkSQL可以通過反射自動獲取到并添加給DF
    val personDF: DataFrame = rowRDD.toDF
    personDF.show(10)
    personDF.printSchema()
    sc.stop()
    spark.stop()
  }
}

1.4、花式查詢

package cn.itcast.sql

import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.{DataFrame, SparkSession}


object QueryDemo {
case class Person(id:Int,name:String,age:Int)
  def main(args: Array[String]): Unit = {
    //1.創(chuàng)建SparkSession
    val spark: SparkSession = SparkSession.builder().master("local[*]").appName("SparkSQL")
.getOrCreate()
    val sc: SparkContext = spark.sparkContext
    sc.setLogLevel("WARN")
    //2.讀取文件
    val fileRDD: RDD[String] = sc.textFile("D:\\data\\person.txt")
    val linesRDD: RDD[Array[String]] = fileRDD.map(_.split(" "))
    val rowRDD: RDD[Person] = linesRDD.map(line =>Person(line(0).toInt,line(1),line(2).toInt))
    //3.將RDD轉成DF
    //注意:RDD中原本沒有toDF方法,新版本中要給它增加一個方法,可以使用隱式轉換
    import spark.implicits._
    //注意:上面的rowRDD的泛型是Person,里面包含了Schema信息
    //所以SparkSQL可以通過反射自動獲取到并添加給DF
    val personDF: DataFrame = rowRDD.toDF
    personDF.show(10)
    personDF.printSchema()
    //=======================SQL方式查詢=======================
    //0.注冊表
    personDF.createOrReplaceTempView("t_person")
    //1.查詢所有數(shù)據(jù)
    spark.sql("select * from t_person").show()
    //2.查詢age+1
    spark.sql("select age,age+1 from t_person").show()
    //3.查詢age最大的兩人
    spark.sql("select name,age from t_person order by age desc limit 2").show()
    //4.查詢各個年齡的人數(shù)
    spark.sql("select age,count(*) from t_person group by age").show()
    //5.查詢年齡大于30的
    spark.sql("select * from t_person where age > 30").show()

    //=======================DSL方式查詢=======================
    //1.查詢所有數(shù)據(jù)
    personDF.select("name","age")
    //2.查詢age+1
    personDF.select($"name",$"age" + 1)
    //3.查詢age最大的兩人
    personDF.sort($"age".desc).show(2)
    //4.查詢各個年齡的人數(shù)
    personDF.groupBy("age").count().show()
    //5.查詢年齡大于30的
    personDF.filter($"age" > 30).show()

    sc.stop()
    spark.stop()
  }
  }

1.5、 相互轉化

RDD、DF、DS之間的相互轉換有很多(6種),但是我們實際操作就只有2類:
1)使用RDD算子操作
2)使用DSL/SQL對表操作

package cn.itcast.sql

import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.{DataFrame, Dataset, Row, SparkSession}

object TransformDemo {
case class Person(id:Int,name:String,age:Int)

  def main(args: Array[String]): Unit = {
    //1.創(chuàng)建SparkSession
    val spark: SparkSession = SparkSession.builder().master("local[*]").appName("SparkSQL").getOrCreate()
    val sc: SparkContext = spark.sparkContext
    sc.setLogLevel("WARN")
    //2.讀取文件
    val fileRDD: RDD[String] = sc.textFile("D:\\data\\person.txt")
    val linesRDD: RDD[Array[String]] = fileRDD.map(_.split(" "))
    val personRDD: RDD[Person] = linesRDD.map(line =>Person(line(0).toInt,line(1),line(2).toInt))
    //3.將RDD轉成DF
    //注意:RDD中原本沒有toDF方法,新版本中要給它增加一個方法,可以使用隱式轉換
    import spark.implicits._
    //注意:上面的rowRDD的泛型是Person,里面包含了Schema信息
    //所以SparkSQL可以通過反射自動獲取到并添加給DF
    //=========================相互轉換======================
    //1.RDD-->DF
    val personDF: DataFrame = personRDD.toDF
    //2.DF-->RDD
    val rdd: RDD[Row] = personDF.rdd
    //3.RDD-->DS
    val DS: Dataset[Person] = personRDD.toDS()
    //4.DS-->RDD
    val rdd2: RDD[Person] = DS.rdd
    //5.DF-->DS
    val DS2: Dataset[Person] = personDF.as[Person]
    //6.DS-->DF
    val DF: DataFrame = DS2.toDF()

    sc.stop()
    spark.stop()
  }
  }

1.6、Spark SQL完成WordCount(案例)

1.6.1、SQL風格

package cn.itcast.sql

import org.apache.spark.SparkContext
import org.apache.spark.sql.{DataFrame, Dataset, SparkSession}


object WordCount {
  def main(args: Array[String]): Unit = {
    //1.創(chuàng)建SparkSession
    val spark: SparkSession = SparkSession.builder().master("local[*]").appName("SparkSQL").getOrCreate()
    val sc: SparkContext = spark.sparkContext
    sc.setLogLevel("WARN")
    //2.讀取文件
    val fileDF: DataFrame = spark.read.text("D:\\data\\words.txt")
    val fileDS: Dataset[String] = spark.read.textFile("D:\\data\\words.txt")
    //fileDF.show()
    //fileDS.show()
    //3.對每一行按照空格進行切分并壓平
    //fileDF.flatMap(_.split(" ")) //注意:錯誤,因為DF沒有泛型,不知道_是String
    import spark.implicits._
    val wordDS: Dataset[String] = fileDS.flatMap(_.split(" "))//注意:正確,因為DS有泛型,知道_是String
    //wordDS.show()
    /*
    +-----+
    |value|
    +-----+
    |hello|
    |   me|
    |hello|
    |  you|
      ...
     */
    //4.對上面的數(shù)據(jù)進行WordCount
    wordDS.createOrReplaceTempView("t_word")
    val sql =
      """
        |select value ,count(value) as count
        |from t_word
        |group by value
        |order by count desc
      """.stripMargin
    spark.sql(sql).show()

    sc.stop()
    spark.stop()
  }
}

1.6.2、DQL風格

package cn.itcast.sql

import org.apache.spark.SparkContext
import org.apache.spark.sql.{DataFrame, Dataset, SparkSession}


object WordCount2 {
  def main(args: Array[String]): Unit = {
    //1.創(chuàng)建SparkSession
    val spark: SparkSession = SparkSession.builder().master("local[*]").appName("SparkSQL").getOrCreate()
    val sc: SparkContext = spark.sparkContext
    sc.setLogLevel("WARN")
    //2.讀取文件
    val fileDF: DataFrame = spark.read.text("D:\\data\\words.txt")
    val fileDS: Dataset[String] = spark.read.textFile("D:\\data\\words.txt")
    //fileDF.show()
    //fileDS.show()
    //3.對每一行按照空格進行切分并壓平
    //fileDF.flatMap(_.split(" ")) //注意:錯誤,因為DF沒有泛型,不知道_是String
    import spark.implicits._
    val wordDS: Dataset[String] = fileDS.flatMap(_.split(" "))//注意:正確,因為DS有泛型,知道_是String
    //wordDS.show()
    /*
    +-----+
    |value|
    +-----+
    |hello|
    |   me|
    |hello|
    |  you|
      ...
     */
    //4.對上面的數(shù)據(jù)進行WordCount
    wordDS.groupBy("value").count().orderBy($"count".desc).show()

    sc.stop()
    spark.stop()
  }
}

好了,以上內容就到這里了。你學到了嗎。

到此這篇關于如何使用IDEA開發(fā)Spark SQL程序(一文搞懂)的文章就介紹到這了,更多相關IDEA開發(fā)Spark SQL內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • java多態(tài)注意項小結

    java多態(tài)注意項小結

    面向對象的三大特性:封裝、繼承、多態(tài)。從一定角度來看,封裝和繼承幾乎都是為多態(tài)而準備的。今天通過本文給大家介紹java多態(tài)注意項總結,感興趣的朋友一起看看吧
    2021-10-10
  • Java基礎之自動裝箱,注解操作示例

    Java基礎之自動裝箱,注解操作示例

    這篇文章主要介紹了Java基礎之自動裝箱,注解操作,結合實例形式分析了java拆箱、裝箱、靜態(tài)導入、注釋等相關使用技巧,需要的朋友可以參考下
    2019-08-08
  • SpringBoot如何使用Undertow做服務器

    SpringBoot如何使用Undertow做服務器

    這篇文章主要介紹了SpringBoot如何使用Undertow做服務器,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-07-07
  • java 數(shù)據(jù)類型有哪些取值范圍多少

    java 數(shù)據(jù)類型有哪些取值范圍多少

    這篇文章主要介紹了java 數(shù)據(jù)類型有哪些取值范圍多少的相關資料,網上關于java 數(shù)據(jù)類型的資料有很多,不夠全面,這里就整理下,需要的朋友可以參考下
    2017-01-01
  • 詳解使用spring boot admin監(jiān)控spring cloud應用程序

    詳解使用spring boot admin監(jiān)控spring cloud應用程序

    這篇文章主要介紹了詳解使用spring boot admin監(jiān)控spring cloud應用程序,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-05-05
  • Spring Cloud Netflix架構淺析(小結)

    Spring Cloud Netflix架構淺析(小結)

    這篇文章主要介紹了Spring Cloud Netflix架構淺析(小結),詳解的介紹了Spring Cloud Netflix的概念和組件等,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2018-01-01
  • java LeetCode題解KMP算法示例

    java LeetCode題解KMP算法示例

    這篇文章主要為大家介紹了java LeetCode題解KMP算法示例,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-10-10
  • java使用xstream實現(xiàn)xml文件和對象之間的相互轉換

    java使用xstream實現(xiàn)xml文件和對象之間的相互轉換

    xml是一個用途比較廣泛的文件類型,在java里也自帶解析xml的包,但是本文使用的是xstream來實現(xiàn)xml和對象之間的相互轉換,xstream是一個第三方開源框架,使用起來比較方便,對java?xml和對象轉換相關知識感興趣的朋友一起看看吧
    2023-09-09
  • 解決MyBatis中Enum字段參數(shù)解析問題

    解決MyBatis中Enum字段參數(shù)解析問題

    本文主要介紹了解決MyBatis中Enum字段參數(shù)解析問題,文中通過示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-08-08
  • Java中ArrayList類的使用方法

    Java中ArrayList類的使用方法

    ArrayList就是傳說中的動態(tài)數(shù)組,用MSDN中的說法,就是Array的復雜版本,下面來簡單介紹下
    2013-12-12

最新評論

固始县| 株洲县| 六盘水市| 施秉县| 张家界市| 曲周县| 普定县| 元阳县| 江安县| 库尔勒市| 汪清县| 卢氏县| 延川县| 冀州市| 陈巴尔虎旗| 金昌市| 毕节市| 松滋市| 嫩江县| 瓦房店市| 泰宁县| 井冈山市| 达日县| 三穗县| 昆明市| 盘锦市| 柘荣县| 济源市| 霸州市| 宝丰县| 屏南县| 犍为县| 浦东新区| 红河县| 察哈| 平乐县| 岚皋县| 无极县| 朝阳区| 江西省| 阿城市|