IDEA开发Spark应用

创建Maven项目

配置Maven路径

在这里插入图片描述

设置编码

在这里插入图片描述

配置POM依赖

  1. 配置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>
  1. 配置编译插件

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路径下生成了文件。
在这里插入图片描述

更多推荐