IDEA开发Spark应用
·
IDEA开发Spark应用
创建Maven项目
配置Maven路径

设置编码

配置POM依赖
- 配置spark依赖
为保证编译环境和生产环境一致,根据生产集群所使用的scala、spark、hadoop版本添加依赖spark-core、spark-yarn、scala-library。
<properties>
<spark.version>3.1.2</spark.version>
<scala.version>2.12.10</scala.version>
<scala.binary.version>2.12</scala.binary.version>
</properties>
<dependencies>
<!-- https://mvnrepository.com/artifact/org.apache.spark/spark-core -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_${scala.binary.version}</artifactId>
<version>${spark.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-yarn_${scala.binary.version}</artifactId>
<version>${spark.version}</version>
<scope>provided</scope>
</dependency>
<!-- https://mvnrepository.com/artifact/org.scala-lang/scala-library -->
<dependency>
<groupId>org.scala-lang</groupId>
<artifactId>scala-library</artifactId>
<version>${scala.version}</version>
<scope>provided</scope>
</dependency>
</dependencies>
- 配置编译插件
Maven默认不支持编译scala,需要添加编译插件。
<properties>
<scala.maven.plugin.version>4.5.6</scala.maven.plugin.version>
</properties>
<build>
<plugins>
<!-- Scala Compiler Plugin -->
<plugin>
<groupId>net.alchim31.maven</groupId>
<artifactId>scala-maven-plugin</artifactId>
<version>${scala.maven.plugin.version}</version>
<executions>
<execution>
<goals>
<goal>compile</goal> <!-- 编译主代码 -->
<goal>testCompile</goal> <!-- 编译测试代码 -->
</goals>
</execution>
</executions>
<configuration>
<scalaVersion>${scala.version}</scalaVersion>
</configuration>
</plugin>
<!-- Maven Compiler Plugin for Java -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.6.1</version>
<configuration>
<source>1.8</source>
<target>1.8</target>
</configuration>
</plugin>
<!-->加入maven-assembly-plugin 打包插件,则会将依赖的jar包都进行打包<-->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-assembly-plugin</artifactId>
<version>3.6.0</version>
<configuration>
<descriptorRefs>
<descriptorRef>jar-with-dependencies</descriptorRef>
</descriptorRefs>
</configuration>
<executions>
<execution>
<id>make-assembly</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
代码开发
1、新建一个名为 scala 的目录与工程默认的 java 目录结构上保持一致,这样做是为了可以在 scalca 中调用 Java 方法,实现混合开发的目的,并标识为source root。

2、创建scala object,逻辑是经典的统计单词数
object WordCount {
def main(args: Array[String]): Unit = {
if (args.length < 2) {
println("请输入两个参数")
System.exit((1))
}
// 1. 创建SparkContext
val conf = new SparkConf().setAppName("word_count")
val sc = new SparkContext(conf)
// 2. 读取文件
val rdd1: RDD[String] = sc.textFile(args(0))
// 3. 处理
// 1. 把整句话拆分为多个单词
val rdd2: RDD[String] = rdd1.flatMap(item => item.split(" "))
// 2. 把每个单词指定一个词频1
val rdd3: RDD[(String, Int)] = rdd2.map(item => (item, 1))
// 3. 聚合
val result: RDD[(String, Int)] = rdd3.reduceByKey(_ + _)
// 4. 得到结果,保存文件
result.saveAsTextFile(args(1))
sc.stop()
}
}
package maven项目,打包好的jar包在target目录下,jar-with-dependencies后缀的是携带依赖包的,另一个是不携带依赖的,由于pom中设置的依赖范围都是provided,因此这两个文件大小一样。

Spark On Yarn运行jar包
将生成的spark_word_count-1.0-SNAPSHOTjar包重命名为wc_args.jar,上传至集群客户端或者hdfs,本文选择上传到集群客户端路径下。

提交应用到Spark集群:
/opt/module/spark-3.1.2/bin/spark-submit \
--master yarn \
--deploy-mode client \
--class com.oppo.spark.WordCount \
./wc_args.jar \
hdfs://hadoop100:8020/input/wcinput hdfs://hadoop100:8020/output
注:输出路径不能存在。
打开yarn application,可以看到提交的word_count应用。


打开hdfs页面,查看/output路径下生成了文件。

更多推荐
所有评论(0)