• spark入门学习-3-SparkSQL数据抽象


    SparkSQL的数据抽象:DataFrame和DataSet,底层是RDD。
    在这里插入图片描述
    DataFrame = RDD - 泛型 +Schema(指定了字段名和类型)+ SQL操作 + 优化
    DataFrame 就是在RDD的基础之上做了进一步的封装,支持SQL操作!
    DataFrame 就是一个分布式表!

    DataSet = DataFrame + 泛型
    DataSet = RDD + Schema约束(指定了字段名和类型) + SQL操作 + 优化

    2 加载数据

    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.sql.{DataFrame, SparkSession}
    
    object Demo extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo").master("local[*]").getOrCreate()
        // 加载数据
        val dfText: DataFrame =
          spark.read.text("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/helloworld.txt")
    
        val dfJson: DataFrame =
          spark.read.json("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/people.json")
    
        val dfCSV =
          spark.read.csv("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/country.csv")
    
        // 输出结果
        dfText.printSchema()
    
        /**
         * root
         * |-- value: string (nullable = true)
         */
        dfText.show()
    
        /**
         * +------------------+
         * |             value|
         * +------------------+
         * |hello me you her !|
         * |      hello me you|
         * |          hello me|
         * |             hello|
         * |                 #|
         * |                 #|
         * |                 #|
         * |                  |
         * +------------------+
         */
    
        dfJson.printSchema()
    
        /**
         * root
         * |-- age: long (nullable = true)
         * |-- height: double (nullable = true)
         * |-- name: string (nullable = true)
         */
        dfJson.show()
    
        /**
         * +---+------+-------+
         * |age|height|   name|
         * +---+------+-------+
         * | 10| 168.8|Michael|
         * | 30| 168.8|   Andy|
         * | 19| 169.8| Justin|
         * | 32| 188.8| 王启峰|
         * | 10| 168.8|   John|
         * | 19| 179.8|   Domu|
         * | 13| 178.8| 郭英伟|
         * | 18| 175.8| 王  荃|
         * | 19| 190.8| 米  鼎|
         * +---+------+-------+
         */
    
        dfCSV.printSchema()
    
        /**
         * root
         * |-- _c0: string (nullable = true)
         * |-- _c1: string (nullable = true)
         * |-- _c2: string (nullable = true)
         */
        dfCSV.show()
    
        /** 
         * +---+----------------+---+
         * |_c0|             _c1|_c2|
         * +---+----------------+---+
         * |  1|            中国|  1|
         * |  2|      阿尔巴尼亚|ALB|
         * |  3|      阿尔及利亚|DZA|
         * |  4|          阿富汗|AFG|
         * |  5|          阿根廷|ARG|
         * |  6|阿拉伯联合酋长国|ARE|
         * |  7|          阿鲁巴|ABW|
         * |  8|            阿曼|OMN|
         * |  9|        阿塞拜疆|AZE|
         * | 10|        阿森松岛|ASC|
         * | 11|            埃及|EGY|
         * | 12|      埃塞俄比亚|ETH|
         * | 13|          爱尔兰|IRL|
         * | 14|        爱沙尼亚|EST|
         * | 15|          安道尔|AND|
         * | 16|          安哥拉|AGO|
         * | 17|          安圭拉|AIA|
         * | 18|安提瓜岛和巴布达|ATG|
         * | 19|        澳大利亚|AUS|
         * | 20|          奥地利|AUT|
         * +---+----------------+---+
         *
         *
         */
      }
    
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    • 57
    • 58
    • 59
    • 60
    • 61
    • 62
    • 63
    • 64
    • 65
    • 66
    • 67
    • 68
    • 69
    • 70
    • 71
    • 72
    • 73
    • 74
    • 75
    • 76
    • 77
    • 78
    • 79
    • 80
    • 81
    • 82
    • 83
    • 84
    • 85
    • 86
    • 87
    • 88
    • 89
    • 90
    • 91
    • 92
    • 93
    • 94
    • 95
    • 96
    • 97
    • 98
    • 99
    • 100
    • 101
    • 102
    • 103
    • 104
    • 105
    • 106
    • 107
    • 108
    • 109
    • 110
    • 111

    3 RDD 转 DF

    1 样例类

    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.rdd.RDD
    import org.apache.spark.sql.{DataFrame, SparkSession}
    
    object Demo_rdd_df extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo_rdd_df").master("local[*]").getOrCreate()
        val sc = spark.sparkContext
        // 加载数据
        val lines: RDD[String] = sc.textFile("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/person.txt")
    
        val personRDD: RDD[Person] = lines.map(line => {
          val arr: Array[String] = line.split(" ")
          Person(arr(0).toInt, arr(1), arr(2).toInt)
        })
    
        // RDD -> DF
        import spark.implicits._
        val personDF: DataFrame = personRDD.toDF()
    
        personDF.printSchema()
    
        /**
         * root
         * |-- id: integer (nullable = false)
         * |-- name: string (nullable = true)
         * |-- age: integer (nullable = false)
         */
        personDF.show()
    
        /**
         * +---+--------+---+
         * | id|    name|age|
         * +---+--------+---+
         * |  1|zhangsan| 20|
         * |  2|    lisi| 29|
         * |  3|  wangwu| 25|
         * |  4| zhaoliu| 30|
         * |  5|  tianqi| 35|
         * |  5|    kobe| 40|
         * +---+--------+---+
         */
    
      }
    
    }
    
    case class Person(id: Int, name: String, age: Int)
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51

    2 指定类型+列名

    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.rdd.RDD
    import org.apache.spark.sql.{DataFrame, SparkSession}
    
    object Demo_rdd_df_no_case extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo_rdd_df_no_case").master("local[*]").getOrCreate()
        val sc = spark.sparkContext
        // 加载数据
        val lines: RDD[String] = sc.textFile("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/person.txt")
    
        val personRDD: RDD[(Int, String, Int)] = lines.map(line => {
          val arr: Array[String] = line.split(" ")
          (arr(0).toInt, arr(1), arr(2).toInt)
        })
    
        // RDD -> DF
        import spark.implicits._
        val personDF: DataFrame = personRDD.toDF("id", "name", "age")
    
        personDF.printSchema()
    
        /**
         * root
         * |-- id: integer (nullable = false)
         * |-- name: string (nullable = true)
         * |-- age: integer (nullable = false)
         */
        personDF.show()
    
        /**
         * +---+--------+---+
         * | id|    name|age|
         * +---+--------+---+
         * |  1|zhangsan| 20|
         * |  2|    lisi| 29|
         * |  3|  wangwu| 25|
         * |  4| zhaoliu| 30|
         * |  5|  tianqi| 35|
         * |  5|    kobe| 40|
         * +---+--------+---+
         */
    
      }
    
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50

    3 自定义Scheme

    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.rdd.RDD
    import org.apache.spark.sql.types.{BooleanType, IntegerType, LongType, StringType, StructField, StructType}
    import org.apache.spark.sql.{DataFrame, Row, SparkSession}
    
    object Demo_rdd_df_scheme extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo_rdd_df_no_case").master("local[*]").getOrCreate()
        val sc = spark.sparkContext
        // 加载数据
        val lines: RDD[String] = sc.textFile("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/person.txt")
    
        val personRDD: RDD[Row] = lines.map(line => {
          val arr: Array[String] = line.split(" ")
          Row(arr(0).toInt, arr(1), arr(2).toInt)
        })
    
        // RDD -> DF
        import spark.implicits._
        val scheme = StructType(
              List(StructField("id", IntegerType, true),
              StructField("name", StringType, false),
              StructField("age", IntegerType, false))
        )
        val personDF = spark.createDataFrame(personRDD, scheme)
    
        personDF.printSchema()
    
        /**
         * root
         * |-- id: integer (nullable = false)
         * |-- name: string (nullable = true)
         * |-- age: integer (nullable = false)
         */
        personDF.show()
    
        /**
         * +---+--------+---+
         * | id|    name|age|
         * +---+--------+---+
         * |  1|zhangsan| 20|
         * |  2|    lisi| 29|
         * |  3|  wangwu| 25|
         * |  4| zhaoliu| 30|
         * |  5|  tianqi| 35|
         * |  5|    kobe| 40|
         * +---+--------+---+
         */
    
      }
    
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56

    4 RDD DF DS

    在这里插入图片描述

    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.rdd.RDD
    import org.apache.spark.sql.{DataFrame, Dataset, Row, SparkSession}
    
    object Demo_rdd_df_ds extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo_rdd_df").master("local[*]").getOrCreate()
        val sc = spark.sparkContext
        // 加载数据
        val lines: RDD[String] = sc.textFile("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/person.txt")
    
        val personRDD: RDD[Person] = lines.map(line => {
          val arr: Array[String] = line.split(" ")
          Person(arr(0).toInt, arr(1), arr(2).toInt)
        })
    
        import spark.implicits._
        // RDD -> DF
        val personDF: DataFrame = personRDD.toDF()
    
        // RDD -> DS
        val personDS: Dataset[Person] = personRDD.toDS()
        personDS.show()
        /**
         * +---+--------+---+
         * | id|    name|age|
         * +---+--------+---+
         * |  1|zhangsan| 20|
         * |  2|    lisi| 29|
         * |  3|  wangwu| 25|
         * |  4| zhaoliu| 30|
         * |  5|  tianqi| 35|
         * |  5|    kobe| 40|
         * +---+--------+---+
         */
    
        // DF -> RDD DF没有泛型,所以返回Row
        val rdd: RDD[Row] = personDF.rdd
    
        // DS -> RDD
        val rdd1: RDD[Person] = personDS.rdd
    
        // DF(没有泛型) -> DS(有泛型)
        val ds: Dataset[Person] = personDF.as[Person]
    
        // DS -> DF
        val df: DataFrame = personDS.toDF()
        df.show()
    
        /**
         * +---+--------+---+
         * | id|    name|age|
         * +---+--------+---+
         * |  1|zhangsan| 20|
         * |  2|    lisi| 29|
         * |  3|  wangwu| 25|
         * |  4| zhaoliu| 30|
         * |  5|  tianqi| 35|
         * |  5|    kobe| 40|
         * +---+--------+---+
         */
      }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    • 57
    • 58
    • 59
    • 60
    • 61
    • 62
    • 63
    • 64
    • 65
    • 66

    5 SparkSQL 花式查询

    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.rdd.RDD
    import org.apache.spark.sql.{DataFrame, Dataset, SparkSession}
    
    object Demo_spark_sql_get extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo_rdd_df").master("local[*]").getOrCreate()
        val sc = spark.sparkContext
        // 加载数据
        val lines: RDD[String] = sc.textFile("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/person.txt")
    
        val personRDD: RDD[Person] = lines.map(line => {
          val arr: Array[String] = line.split(" ")
          Person(arr(0).toInt, arr(1), arr(2).toInt)
        })
    
        import spark.implicits._
        // RDD -> DF
        val personDF: DataFrame = personRDD.toDF()
    
        /**
         * sql 查询
         */
        // 注册表名
        personDF.createOrReplaceTempView("tb_person")
        spark.sql(
          s"""
            |select name from tb_person
            |""".stripMargin).show()
        spark.sql(
          """
            |select id, name, age + 100 from tb_person
            |""".stripMargin).show()
    
        personDF.select("name").show()
        personDF.select("name", "age").show()
    
        // $ 把字符串转为对象
        personDF.select(($"age" + 100).as("age")).show()
    
        personDF.filter($"age" > 25).show()
    
        println(personDF.filter($"age" > 30).count())
    
        personDF.groupBy("age").count().show()
    
        personDF.filter($"name" === "zhangsan").show()
        
    
        // RDD -> D
        val personDS: Dataset[Person] = personRDD.toDS()
    
      }
    
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    • 57
    • 58
    • 59

    6 wordCount

    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.rdd.RDD
    import org.apache.spark.sql.{Dataset, Row, SparkSession}
    
    object Demo_spark_wordcount extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo_spark_wordcount").master("local[*]").getOrCreate()
        // 加载数据
        val ds: Dataset[String] =
          spark.read.textFile("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/helloworld.txt")
        ds.show()
        /**
         * +------------------+
         * |             value|
         * +------------------+
         * |hello me you her !|
         * |      hello me you|
         * |          hello me|
         * |             hello|
         * |                 #|
         * |                 #|
         * |                 #|
         * |                  |
         * +------------------+
         */
        import spark.implicits._
        val words: Dataset[String] = ds.flatMap(_.split(" "))
        words.createOrReplaceTempView("tb_words")
        spark.sql(
          """
            |select value, count(*) as counts from tb_words group by  value order by counts desc
            |""".stripMargin).show()
    
        words.groupBy('value)
          .count()
          .orderBy('count.desc).show()
    
      }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42

    7 多数据源支持

    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.sql.{DataFrame, Dataset, SaveMode, SparkSession}
    
    import java.util.Properties
    
    object Demo_spark_save_load extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo_spark_save_load").master("local[*]").getOrCreate()
        // 加载数据
        val df: DataFrame =
          spark.read.json("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/people.json")
        df.show()
    //    df.coalesce(1).write.json("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/save/people_json")
    //    df.coalesce(1).write.csv("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/save/people_csv")
    //    df.coalesce(1).write.orc("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/save/people_orc")
    
        val url = "jdbc:mysql://localhost:3306/test"
        val table = "tb_people"
        val props = new Properties()
        props.put("user", "root");
        props.put("password", "root")
        df.write.mode(SaveMode.Overwrite).jdbc(url, table, props)
      }
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28

    8 案例

    对电影评分数据进行统计分析,分别使用DSL编程和SQL编程,获取电影平均分Top10,要求电影的评分次数大于200

    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.sql.{DataFrame, Dataset, SaveMode, SparkSession}
    
    import java.util.Properties
    
    object Demo_spark_file extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo_spark_file").master("local[*]").getOrCreate()
        // 加载数据
        val ds: Dataset[String] = spark.read.textFile("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/file.csv")
        import spark.implicits._
        val df: DataFrame = ds.map(line => {
          val arr: Array[String] = line.split(",")
          (arr(1), arr(2).toInt)
        }).toDF("movieId", "score").cache()
        df.printSchema()
        df.show()
    
        df.createOrReplaceTempView("t_movies")
    
        /**
         * +-------+-----+
         * |movieId|score|
         * +-------+-----+
         * |    302|    3|
         */
        // 统计评分次数大于200的电影的平均分Top10
        val movieDf = spark.sql(
          s"""
             |select movieId, avg(score) as avg_score, count(1) as counts
             |from t_movies
             |group by movieId
             |having counts > 10
             |order by avg_score desc
             |limit 10
             |""".stripMargin).cache()
        movieDf.show()
        
        import org.apache.spark.sql.functions._
        movieDf.groupBy("movieId")
          .agg(
            avg('score) as "avg_score",
            count('movieId) as "counts"
          ).filter('counts > 10)
          .orderBy('avg_score.desc)
          .limit(10)
      }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51

    9 UDF

    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.sql.{Dataset, SparkSession}
    
    object Demo_spark_udf extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo_spark_udf").master("local[*]").getOrCreate()
        // 加载数据
        val ds: Dataset[String] =
          spark.read.textFile("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/helloworld.txt").cache()
        ds.createOrReplaceTempView("tb_word")
        spark.udf.register("small2big", (value: String) => {
          value.toUpperCase()
        })
        val df = spark.sql(
          """
            |select value, small2big(value) as bigvalue
            |from tb_word
            |""".stripMargin)
        df.show()
    
        import org.apache.spark.sql.functions._
        import spark.implicits._
        val small3big = udf((value: String) => {
          value.toUpperCase()
        })
        val dfBig = ds.select('value, small3big('value))
        dfBig.show()
      }
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    package com.baidu.sparkcodetest.sql
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.sql.api.java.UDF1
    import org.apache.spark.sql.types.DataTypes
    import org.apache.spark.sql.{Dataset, SparkSession}
    
    object Demo_spark_udf extends LoggerTrait {
      def main(args: Array[String]): Unit = {
        // 环境创建
        val spark = SparkSession.builder().appName("Demo_spark_udf").master("local[*]").getOrCreate()
        // 加载数据
        val ds: Dataset[String] =
          spark.read.textFile("/Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/helloworld.txt").cache()
        ds.createOrReplaceTempView("tb_word")
    
        spark.udf.register("small2big", new UDF1[String, String] {
          override def call(str: String): String = {
            str.toUpperCase
          }},DataTypes.StringType)
    
        val df = spark.sql(
          """
            |select value, small2big(value) as bigvalue
            |from tb_word
            |""".stripMargin)
        df.show()
    
        import org.apache.spark.sql.functions._
        import spark.implicits._
        val small3big = udf((value: String) => {
          value.toUpperCase()
        })
        val dfBig = ds.select('value, small3big('value))
        dfBig.show()
      }
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38

    10 sparkOnHive

    https://blog.csdn.net/qq_43659234/article/details/117298401
    https://blog.csdn.net/weixin_44911081/article/details/121858327

    SparkOnHive: 仅仅使用Hive的元数据(库/表/字段/位置等信息);剩下的用SparkSQL的,如:执行引擎,语法解析,物理执行计划,SQL优化。
    在这里插入图片描述
    首先要启动: hive --service metastore

    package com.baidu.sparkcodetest.hive
    
    import com.baidu.utils.LoggerTrait
    import org.apache.spark.sql.SparkSession
    
    object Demo_hive extends LoggerTrait{
      def main(args: Array[String]): Unit = {
        // 环境创建  hive --service metastore
        val spark = SparkSession.builder().appName("Demo_flatMap_map").master("local[*]")
          .config("spark.sql.shuffle.partitions", "4")
          .config("spark.sql.warehouse.dir", "hdfs://localhost:8020/user/hive/warehouse") // 指定hdfs位置
          .config("hive.metastore.uris", "thrift://localhost:9083")
          .enableHiveSupport()
          .getOrCreate()
    
        import spark.implicits._
        spark.sql("show databases").show(false)
        spark.sql("show tables").show(false)
        spark.sql("create table person(id int, name string, age int) row format delimited fields terminated by ' '")
        spark.sql("load data local inpath 'file:///Users/zhaoshuai11/Desktop/baidu/ebiz/stu-scala/src/main/resources/person.txt'" +
          " into table person")
        spark.sql("select * from person").show(false)
      }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
  • 相关阅读:
    第十四届蓝桥杯模拟赛(第二场)题解·2022年·C/C++
    java中几种对象存储(文件存储)中间件的介绍
    java毕业设计日租房管理系统源码+lw文档+mybatis+系统+mysql数据库+调试
    机器学习(四十六):Streamlit 构建机器学习 Web
    uniapp中tabbar导航的点击事件
    营销复盘秘籍,6步法让你的活动效果翻倍
    WireShark抓包工具的安装
    最新中文版本FLStudio21水果音乐软件更新下载
    阿里云域名动态解析
    Java常见面试题21-30(集合类)
  • 原文地址:https://blog.csdn.net/zs18753479279/article/details/126099364