SpringBoot操作spark處理hdfs文件的操作方法
SpringBoot操作spark處理hdfs文件

1、導(dǎo)入依賴
<!-- spark依賴-->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_2.12</artifactId>
<version>3.2.2</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_2.12</artifactId>
<version>3.2.2</version>
</dependency>
<!-- https://mvnrepository.com/artifact/org.apache.spark/spark-mllib -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-mllib_2.12</artifactId>
<version>3.2.2</version>
</dependency>2、配置spark信息
建立一個(gè)配置文件,配置spark信息
import org.apache.spark.SparkConf;
import org.apache.spark.sql.SparkSession;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
//將文件交于spring管理
@Configuration
public class SparkConfig {
//使用yml中的配置
@Value("${spark.master}")
private String sparkMaster;
@Value("${spark.appName}")
private String sparkAppName;
@Value("${hdfs.user}")
private String hdfsUser;
@Value("${hdfs.path}")
private String hdfsPath;
@Bean
public SparkConf sparkConf() {
SparkConf conf = new SparkConf();
conf.setMaster(sparkMaster);
conf.setAppName(sparkAppName);
// 添加HDFS配置
conf.set("fs.defaultFS", hdfsPath);
conf.set("spark.hadoop.hdfs.user",hdfsUser);
return conf;
}
@Bean
public SparkSession sparkSession() {
return SparkSession.builder()
.config(sparkConf())
.getOrCreate();
}
}3、controller和service
controller類
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import xyz.zzj.traffic_main_code.service.SparkService;
@RestController
@RequestMapping("/spark")
public class SparkController {
@Autowired
private SparkService sparkService;
@GetMapping("/run")
public String runSparkJob() {
//讀取Hadoop HDFS文件
String filePath = "hdfs://192.168.44.128:9000/subwayData.csv";
sparkService.executeHadoopSparkJob(filePath);
return "Spark job executed successfully!";
}
}處理地鐵數(shù)據(jù)的service
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileStatus;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.types.DataTypes;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import xyz.zzj.traffic_main_code.service.SparkReadHdfs;
import java.io.IOException;
import java.net.URI;
import static org.apache.spark.sql.functions.*;
@Service
public class SparkReadHdfsImpl implements SparkReadHdfs {
private final SparkSession spark;
@Value("${hdfs.user}")
private String hdfsUser;
@Value("${hdfs.path}")
private String hdfsPath;
@Autowired
public SparkReadHdfsImpl(SparkSession spark) {
this.spark = spark;
}
/**
* 讀取HDFS上的CSV文件并上傳到HDFS
* @param filePath
*/
@Override
public void sparkSubway(String filePath) {
try {
// 設(shè)置Hadoop配置
JavaSparkContext jsc = JavaSparkContext.fromSparkContext(spark.sparkContext());
Configuration hadoopConf = jsc.hadoopConfiguration();
hadoopConf.set("fs.defaultFS", hdfsPath);
hadoopConf.set("hadoop.user.name", hdfsUser);
// 讀取HDFS上的文件
Dataset<Row> df = spark.read()
.option("header", "true") // 指定第一行是列名
.option("inferSchema", "true") // 自動(dòng)推斷列的數(shù)據(jù)類型
.csv(filePath);
// 顯示DataFrame的所有數(shù)據(jù)
// df.show(Integer.MAX_VALUE, false);
// 對(duì)DataFrame進(jìn)行清洗和轉(zhuǎn)換操作
// 檢查缺失值
df.select("number", "people", "dateTime").na().drop().show();
// 對(duì)數(shù)據(jù)進(jìn)行類型轉(zhuǎn)換
Dataset<Row> df2 = df.select(
col("number").cast(DataTypes.IntegerType),
col("people").cast(DataTypes.IntegerType),
to_date(col("dateTime"), "yyyy年MM月dd日").alias("dateTime")
);
// 去重
Dataset<Row> df3 = df2.dropDuplicates();
// 數(shù)據(jù)過(guò)濾,確保people列沒(méi)有負(fù)數(shù)
Dataset<Row> df4 = df3.filter(col("people").geq(0));
// df4.show();
// 數(shù)據(jù)聚合,按dateTime分組,統(tǒng)計(jì)每天的總客流量
Dataset<Row> df6 = df4.groupBy("dateTime").agg(sum("people").alias("total_people"));
// df6.show();
sparkForSubway(df6,"/time_subwayData.csv");
//數(shù)據(jù)聚合,獲取每天人數(shù)最多的地鐵number
Dataset<Row> df7 = df4.groupBy("dateTime").agg(max("people").alias("max_people"));
sparkForSubway(df7,"/everyday_max_subwayData.csv");
//數(shù)據(jù)聚合,計(jì)算每天的客流強(qiáng)度:每天總people除以632840
Dataset<Row> df8 = df4.groupBy("dateTime").agg(sum("people").divide(632.84).alias("strength"));
sparkForSubway(df8,"/everyday_strength_subwayData.csv");
} catch (Exception e) {
e.printStackTrace();
}
}
private static void sparkForSubway(Dataset<Row> df6, String hdfsPath) throws IOException {
// 保存處理后的數(shù)據(jù)到HDFS
df6.coalesce(1)
.write().mode("overwrite")
.option("header", "true")
.csv("hdfs://192.168.44.128:9000/time_subwayData");
// 創(chuàng)建Hadoop配置
Configuration conf = new Configuration();
// 獲取FileSystem實(shí)例
FileSystem fs = FileSystem.get(URI.create("hdfs://192.168.44.128:9000"), conf);
// 定義臨時(shí)目錄和目標(biāo)文件路徑
Path tempDir = new Path("/time_subwayData");
FileStatus[] files = fs.listStatus(tempDir);
// 檢查目標(biāo)文件是否存在,如果存在則刪除
Path targetFile1 = new Path(hdfsPath);
if (fs.exists(targetFile1)) {
fs.delete(targetFile1, true); // true 表示遞歸刪除
}
for (FileStatus file : files) {
if (file.isFile() && file.getPath().getName().startsWith("part-")) {
Path targetFile = new Path(hdfsPath);
fs.rename(file.getPath(), targetFile);
}
}
// 刪除臨時(shí)目錄
fs.delete(tempDir, true);
}
}4、運(yùn)行
- 項(xiàng)目運(yùn)行完后,打開瀏覽器
- spark處理地鐵數(shù)據(jù)
- http://localhost:8686/spark/dispose
- 觀察spark和hdfs
- http://192.168.44.128:8099/
- http://192.168.44.128:9870/explorer.html#/

到此這篇關(guān)于SpringBoot操作spark處理hdfs文件的文章就介紹到這了,更多相關(guān)SpringBoot spark處理hdfs文件內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
java中的Arrays這個(gè)工具類你真的會(huì)用嗎(一文秒懂)
這篇文章主要介紹了java中的Arrays這個(gè)工具類你真的會(huì)用嗎,本文通過(guò)實(shí)例代碼給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2020-06-06
Java中如何動(dòng)態(tài)創(chuàng)建接口的實(shí)現(xiàn)方法
這篇文章主要介紹了Java中如何動(dòng)態(tài)創(chuàng)建接口的實(shí)現(xiàn)方法的相關(guān)資料,需要的朋友可以參考下2017-09-09
Spring實(shí)戰(zhàn)之使用@Resource配置依賴操作示例
這篇文章主要介紹了Spring實(shí)戰(zhàn)之使用@Resource配置依賴操作,結(jié)合實(shí)例形式分析了Spring使用@Resource配置依賴具體步驟、實(shí)現(xiàn)及測(cè)試案例,需要的朋友可以參考下2019-12-12
淺談JAVA版本號(hào)的問(wèn)題 Java版本號(hào)與JDk版本
這篇文章主要介紹了淺談JAVA版本號(hào)的問(wèn)題 Java版本號(hào)與JDk版本,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧2020-08-08
Java實(shí)現(xiàn)彩色圖片轉(zhuǎn)換為灰度圖片的示例代碼
將彩色圖片轉(zhuǎn)換為灰度圖片是圖像處理中的常見操作,通常用于簡(jiǎn)化圖像、增強(qiáng)對(duì)比度、或者進(jìn)行后續(xù)的圖像分析,本項(xiàng)目的目標(biāo)是通過(guò)Java實(shí)現(xiàn)將彩色圖片轉(zhuǎn)換為灰度圖片,需要的朋友可以參考下2025-02-02
Java多線程編程小實(shí)例模擬停車場(chǎng)系統(tǒng)
這是一個(gè)關(guān)于Java多線程編程的例子,用多線程的思想模擬停車場(chǎng)管理系統(tǒng),這里分享給大家,供需要的朋友參考。2017-10-10

