SpringBoot使用Spark過程詳解
前提: 可以參考文章 SpringBoot 接入 Spark
- SpringBoot 已經(jīng)接入 Spark
- 已配置 JavaSparkContext
- 已配置 SparkSession
@Resource private SparkSession sparkSession; @Resource private JavaSparkContext javaSparkContext;
讀取 txt 文件
測(cè)試文件 word.txt

java 代碼
- textFile:獲取文件內(nèi)容,返回 JavaRDD
- flatMap:過濾數(shù)據(jù)
- mapToPair:把每個(gè)元素都轉(zhuǎn)換成一個(gè)<K,V>類型的對(duì)象,如 <123,1>,<456,1>
- reduceByKey:對(duì)相同key的數(shù)據(jù)集進(jìn)行預(yù)聚合
public void testSparkText() {
String file = "D:\\TEMP\\word.txt";
JavaRDD<String> fileRDD = javaSparkContext.textFile(file);
JavaRDD<String> wordsRDD = fileRDD.flatMap(line -> Arrays.asList(line.split(" ")).iterator());
JavaPairRDD<String, Integer> wordAndOneRDD = wordsRDD.mapToPair(word -> new Tuple2<>(word, 1));
JavaPairRDD<String, Integer> wordAndCountRDD = wordAndOneRDD.reduceByKey((a, b) -> a + b);
//輸出結(jié)果
List<Tuple2<String, Integer>> result = wordAndCountRDD.collect();
result.forEach(System.out::println);
}
結(jié)果得出,123 有 3 個(gè),456 有 2 個(gè),789 有 1 個(gè)

讀取 csv 文件
測(cè)試文件 testcsv.csv

java 代碼
public void testSparkCsv() {
String file = "D:\\TEMP\\testcsv.csv";
JavaRDD<String> fileRDD = javaSparkContext.textFile(file);
JavaRDD<String> wordsRDD = fileRDD.flatMap(line -> Arrays.asList(line.split(",")).iterator());
//輸出結(jié)果
System.out.println(wordsRDD.collect());
}
輸出結(jié)果

讀取 MySQL 數(shù)據(jù)庫(kù)表
- format:獲取數(shù)據(jù)庫(kù)建議是 jdbc
- option.url:添加 MySQL 連接 url
- option.user:MySQL 用戶名
- option.password:MySQL 用戶密碼
- option.dbtable:sql 語(yǔ)句
- option.driver:數(shù)據(jù)庫(kù) driver,MySQL 使用 com.mysql.cj.jdbc.Driver
public void testSparkMysql() throws IOException {
Dataset<Row> jdbcDF = sparkSession.read()
.format("jdbc")
.option("url", "jdbc:mysql://192.168.140.1:3306/user?useUnicode=true&characterEncoding=UTF-8&serverTimezone=Asia/Shanghai")
.option("dbtable", "(SELECT * FROM xxxtable) tmp")
.option("user", "root")
.option("password", "xxxxxxxxxx*k")
.option("driver", "com.mysql.cj.jdbc.Driver")
.load();
jdbcDF.printSchema();
jdbcDF.show();
//轉(zhuǎn)化為RDD
JavaRDD<Row> rowJavaRDD = jdbcDF.javaRDD();
System.out.println(rowJavaRDD.collect());
}
也可以把表內(nèi)容輸出到文件,添加以下代碼
List<Row> list = rowJavaRDD.collect();
BufferedWriter bw;
bw = new BufferedWriter(new FileWriter("d:/test.txt"));
for (int j = 0; j < list.size(); j++) {
bw.write(list.get(j).toString());
bw.newLine();
bw.flush();
}
bw.close();
結(jié)果輸出

讀取 Json 文件
測(cè)試文件 testjson.json,內(nèi)容如下
[{
"name": "name1",
"age": "1"
}, {
"name": "name2",
"age": "2"
}, {
"name": "name3",
"age": "3"
}, {
"name": "name4",
"age": "4"
}]
注意:testjson.json 文件的內(nèi)容不能帶格式,需要進(jìn)行壓縮

java 代碼
- createOrReplaceTempView:讀取 json 數(shù)據(jù)后,創(chuàng)建數(shù)據(jù)表 t
- sparkSession.sql:使用 sql 對(duì) t 進(jìn)行查詢,輸出 age 大于 3 的數(shù)據(jù)
public void testSparkJson() {
Dataset<Row> df = sparkSession.read().json("D:\\TEMP\\testjson.json");
df.printSchema();
df.createOrReplaceTempView("t");
Dataset<Row> row = sparkSession.sql("select age,name from t where age > 3");
JavaRDD<Row> rowJavaRDD = row.javaRDD();
System.out.println(rowJavaRDD.collect());
}
輸出結(jié)果

中文輸出亂碼
測(cè)試文件 testcsv.csv

public void testSparkCsv() {
String file = "D:\\TEMP\\testcsv.csv";
JavaRDD<String> fileRDD = javaSparkContext.textFile(file);
JavaRDD<String> wordsRDD = fileRDD.flatMap(line -> Arrays.asList(line.split(",")).iterator());
//輸出結(jié)果
System.out.println(wordsRDD.collect());
}
輸出結(jié)果,發(fā)現(xiàn)中文亂碼,可惡

原因:textFile 讀取文件沒有解決亂碼問題,但 sparkSession.read() 卻不會(huì)亂碼
解決辦法:獲取文件方式由 textFile 改成 hadoopFile,由 hadoopFile 指定具體編碼
public void testSparkCsv() {
String file = "D:\\TEMP\\testcsv.csv";
String code = "gbk";
JavaRDD<String> gbkRDD = javaSparkContext.hadoopFile(file, TextInputFormat.class, LongWritable.class, Text.class).map(p -> new String(p._2.getBytes(), 0, p._2.getLength(), code));
JavaRDD<String> gbkWordsRDD = gbkRDD.flatMap(line -> Arrays.asList(line.split(",")).iterator());
//輸出結(jié)果
System.out.println(gbkWordsRDD.collect());
}
輸出結(jié)果

到此這篇關(guān)于SpringBoot使用Spark過程詳解的文章就介紹到這了,更多相關(guān)SpringBoot Spark內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
使用JPA進(jìn)行CriteriaQuery進(jìn)行查詢的注意事項(xiàng)
這篇文章主要介紹了使用JPA進(jìn)行CriteriaQuery進(jìn)行查詢的注意事項(xiàng),具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-12-12
Java設(shè)置JSON字符串參數(shù)編碼的示例詳解
在Java中創(chuàng)建JSON字符串,我們可以使用多個(gè)庫(kù),其中最流行的是Jackson、Gson和org.json,,下面給大家分享Java設(shè)置JSON字符串參數(shù)編碼的示例,感興趣的朋友一起看看吧2024-06-06
spring cloud 配置中心客戶端啟動(dòng)遇到的問題
這篇文章主要介紹了spring cloud 配置中心客戶端啟動(dòng)遇到的問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-09-09
MyBatis學(xué)習(xí)教程(二)—如何使用MyBatis對(duì)users表執(zhí)行CRUD操作
這篇文章主要介紹了MyBatis學(xué)習(xí)教程(二)—如何使用MyBatis對(duì)users表執(zhí)行CRUD操作的相關(guān)資料,需要的朋友可以參考下2016-05-05
Java 實(shí)現(xiàn)Redis存儲(chǔ)復(fù)雜json格式數(shù)據(jù)并返回給前端
這篇文章主要介紹了Java 實(shí)現(xiàn)Redis存儲(chǔ)復(fù)雜json格式數(shù)據(jù)并返回給前端操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過來看看吧2020-07-07
SpringBoot服務(wù)設(shè)置禁止server.point端口的使用
本文主要介紹了SpringBoot服務(wù)設(shè)置禁止server.point端口的使用,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2024-01-01

