一、搭建Spark项目结构





在SparkProject模块的pom.xml文件中增加一下依赖,并等待依赖包下载完毕,如上图。
-
-
-
- <properties>
-
- <scala.version>2.11scala.version>
-
- <spark.version>2.1.1spark.version>
-
- properties>
-
-
-
- <dependencies>
-
-
- <dependency>
-
- <groupId>com.thoughtworks.paranamergroupId>
-
- <artifactId>paranamerartifactId>
-
- <version>2.8version>
-
- dependency>
-
- <dependency>
-
- <groupId>org.apache.sparkgroupId>
-
- <artifactId>spark-core_${scala.version}artifactId>
-
- <version>${spark.version}version>
-
- dependency>
-
- <dependency>
-
- <groupId>org.apache.sparkgroupId>
-
- <artifactId>spark-sql_${scala.version}artifactId>
-
- <version>${spark.version}version>
-
- dependency>
-
- <dependency>
-
- <groupId>org.apache.sparkgroupId>
-
- <artifactId>spark-streaming_2.11artifactId>
-
- <version>${spark.version}version>
-
- dependency>
-
- <dependency>
-
- <groupId>org.apache.sparkgroupId>
-
- <artifactId>spark-mllib_2.11artifactId>
-
- <version>2.1.1version>
-
- dependency>
-
- <dependency>
-
- <groupId>org.apache.sparkgroupId>
-
- <artifactId>spark-streaming-kafka-0-10_2.11artifactId>
-
- <version>2.3.0version>
-
- dependency>
-
- <dependency>
-
- <groupId>org.apache.sparkgroupId>
-
- <artifactId>spark-streaming-kafka-0-8_${scala.version}artifactId>
-
- <version>2.3.0version>
-
- dependency>
-
- <dependency>
-
- <groupId>net.jpountz.lz4groupId>
-
- <artifactId>lz4artifactId>
-
- <version>1.3.0version>
-
- dependency>
-
- <dependency>
-
- <groupId>mysqlgroupId>
-
- <artifactId>mysql-connector-javaartifactId>
-
- <version>8.0.18version>
-
- dependency>
-
- <dependency>
-
- <groupId>org.apache.flume.flume-ng-clientsgroupId>
-
- <artifactId>flume-ng-log4jappenderartifactId>
-
- <version>1.7.0version>
-
- dependency>
-
-
-
-
-
-
- <dependency>
-
- <groupId>org.apache.sparkgroupId>
-
- <artifactId>spark-hive_2.12artifactId>
-
- <version>2.4.8version>
-
- dependency>
-
- dependencies>
-
-
- <build>
-
- <plugins>
-
- <plugin>
-
- <groupId>org.apache.maven.pluginsgroupId>
-
- <artifactId>maven-compiler-pluginartifactId>
-
- <version>3.8.1version>
-
- <configuration>
-
- <source>1.8source>
-
- <target>1.8target>
-
- configuration>
-
- plugin>
-
- <plugin>
-
- <groupId>org.apache.maven.pluginsgroupId>
-
- <artifactId>maven-assembly-pluginartifactId>
-
- <configuration>
-
- <descriptorRefs>
-
- <descriptorRef>jar-with-dependenciesdescriptorRef>
-
- descriptorRefs>
-
- configuration>
-
- plugin>
-
- plugins>
-
- build>
-
-








二、解决无法创建scala文件问题













三、编写LoggerLevel特质


在特质 下增加如下代码
-
- Logger.getLogger("org").setLevel(Level.ERROR)
-
-
这个时候需要导包




完整代码如下:
-
- import org.apache.log4j.{Level, Logger}
-
-
-
- trait LoggerLevel {
-
- Logger.getLogger("org").setLevel(Level.ERROR)
-
- }
-
-
-
四、编写getLocalSparkSession方法


以下是完整代码:
-
- import org.apache.spark.sql.SparkSession
-
-
-
- object SparkUnit {
-
- /**
- * 一个class参数
- **/
-
- def getLocalSparkSession(appName: String): SparkSession = {
-
- SparkSession.builder().appName(appName).master("local[2]").getOrCreate()
-
- }
-
-
-
- def getLocalSparkSession(appName: String, support: Boolean): SparkSession = {
-
- if (support) SparkSession.builder().master("local[2]").appName(appName).enableHiveSupport().getOrCreate()
-
- else getLocalSparkSession(appName)
-
- }
-
-
-
- def getLocalSparkSession(appName: String, master: String): SparkSession = {
-
- SparkSession.builder().appName(appName).master(master).getOrCreate()
-
- }
-
-
-
- def getLocalSparkSession(appName: String, master: String, support: Boolean): SparkSession = {
-
- if (support) SparkSession.builder().appName(appName).master(master).enableHiveSupport().getOrCreate()
-
- else getLocalSparkSession(appName, master)
-
- }
-
-
-
- def stopSpark(ss: SparkSession) = {
-
- if (ss != null) {
-
- ss.stop()
-
- }
-
- }
-
-
-
- }