SparkSQL的数据抽象:DataFrame和DataSet,底层是RDD。

DataFrame = RDD - 泛型 +Schema(指定了字段名和类型)+ SQL操作 + 优化
DataFrame 就是在RDD的基础之上做了进一步的封装,支持SQL操作!
DataFrame 就是一个分布式表!
DataSet = DataFrame + 泛型
DataSet = RDD + Schema约束(指定了字段名和类型) + SQL操作 + 优化
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|
* +---+----------------+---+
*
*
*/
}
}
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)
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|
* +---+--------+---+
*/
}
}
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|
* +---+--------+---+
*/
}
}

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|
* +---+--------+---+
*/
}
}
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()
}
}
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()
}
}
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)
}
}
对电影评分数据进行统计分析,分别使用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)
}
}
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()
}
}
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()
}
}
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)
}
}