2-Storm集群安装(伪分布式)

基础环境:百度网盘,提取码:NIIT
--来自百度网盘超级会员V4的分享)

安装ZooKeeper集群

Storm 使用 Zookeeper 协调管理集群. Zookeeper 并不是 用于消息传递, 所以 Storm 对Zookeeper造成的负载压力非常低. 单节点Zookeeper集群在大多数情况下应该是足够的,但是如果您想要故障转移或部署大型Storm集群,则可能需要较大的Zookeeper集群. 部署Zookeeper的说明是这里.

这里我们安装一个单机伪分布式的Zookeeper集群:

conf/zoo_1.cfg

tickTime=2000
initLimit=10
syncLimit=5
dataDir=/home/hadoop/tmp/zk1/data
dataLogDir=/home/hadoop/tmp/zk1/dataLog
clientPort=2181
4lw.commands.whitelist=*
server.1=localhost:2891:3891
server.2=localhost:2892:3892
server.3=localhost:2893:3893

conf/zoo_2.cfg

tickTime=2000
initLimit=10
syncLimit=5
dataDir=/home/hadoop/tmp/zk2/data
dataLogDir=/home/hadoop/tmp/zk2/dataLog
clientPort=2182
4lw.commands.whitelist=*
server.1=localhost:2891:3891
server.2=localhost:2892:3892
server.3=localhost:2893:3893

conf/zoo_3.cfg

tickTime=2000
initLimit=10
syncLimit=5
dataDir=/home/hadoop/tmp/zk3/data
dataLogDir=/home/hadoop/tmp/zk3/dataLog
clientPort=2183
4lw.commands.whitelist=*
server.1=localhost:2891:3891
server.2=localhost:2892:3892
server.3=localhost:2893:3893

在dataDir中创建文件myid内容是1或2或3

echo 1 > /home/hadoop/tmp/zk1/data/myid
echo 2 > /home/hadoop/tmp/zk2/data/myid
echo 3 > /home/hadoop/tmp/zk3/data/myid

启动服务

bin/zkServer.sh start conf/zoo_1.cfg
bin/zkServer.sh start conf/zoo_2.cfg
bin/zkServer.sh start conf/zoo_3.cfg

关于Zookeeper部署的几点注意事项:

1.在监督(supervision)下运行Zookeeper至关重要,因为Zookeeper是快速失败的,如果遇到任何错误的情况都将退出进程. 有关详细信息,请参阅这里 .

2.建立一个cron定时任务来压缩Zookeeper的数据和事务日志至关重要. Zookeeper守护进程本身不会这样做,如果没有设置cron,Zookeeper将很快耗尽磁盘空间. 有关详细信息,请参阅这里.

Nimbus和worker节点的安装环境

接下来你需要准备Nimbus 和 worker 节点的安装环境:

  1. Java 8+ (Apache Storm 2.x is tested through travisci against a java 8 JDK)
  2. Python 2.6.6 (Python 3.x should work too,but is not tested as part of our CI enviornment)

这些依赖版本是Storm已经测试过的. Storm 在不同的Java 或Python版本上也许会存在问题.

下载解压 Storm

接下来,下载一个Storm版本,并解压zip文件到Nimbus和每个worker机器上的某个目录下. Storm版本可以从这里下载.

当前最新的版本是2.2.0,这里使用2.1.0的版本。

我们需要下载:

也可以使用wget下载,如

wget https://dlcdn.apache.org/storm/apache-storm-2.1.0/apache-storm-2.1.0.tar.gz

在storm.yaml中设置必要的配置

Storm 发布包中在目录conf/storm.yaml 下包含一个默认的配置文件. 你可以在[这里](http://github.com/apache/storm/blob/master /conf/defaults.yaml)查看默认值. storm.yaml 中的存在的配置项会覆盖掉 defaults.yaml中相应的配置项. 下面一些配置是集群运行时所必要的:

  1. storm.zookeeper.servers: 这是一个Storm集群所依赖 Zookeeper 集群的hosts列表. 类似于:
storm.zookeeper.servers:
  - "111.222.333.444"
  - "555.666.777.888"

如果配置的Zookeeper集群不是默认的端口, 你应该设置 storm.zookeeper.port 选项.

  1. storm.local.dir: Nimbus 和 Supervisor 守护进程需要配置一个本地目录来存储少量状态信息(例如jars包,配置文件等等). 您应该在每个机器上创建该目录,给予适当的权限,然后使用此配置填写目录位置. 例如:
storm.local.dir: "/mnt/storm"

如果您在windows下运行Strom,应该如下: yaml storm.local.dir: "C:\\storm-local" 如果您使用相对路径,那么路径是相对于(STORM_HOME). 您也可以使用默认值 $STORM_HOME/storm-local

  1. nimbus.seeds: worker节点需要知道哪些机器是主机的候选者,以便下载 topology jar和confs(nimbus.host 在1.0之后已经废弃,这里实现了HA). 例如:
nimbus.seeds: ["111.222.333.44"]

鼓励您填写机器的FQDN (Fully Qualified Domain Name,全域名)列表. 如果要设置Nimbus HA,则必须解决运行nimbus的所有机器的FQDN.当您只想设置“伪分布式”集群时您可能希望将其保留为默认值,仍然鼓励您填写FQDN.

  1. supervisor.slots.ports: 对于每个worker节点,您可以使用此配置设置在该计算机上运行的worker数量. 每个worker使用单个端口接收消息,并且此设置定义哪些端口打开以供使用. 如果您在此定义五个端口,那么Storm将分配最多五个worker在本机上运行. 如果您定义了三个端口,Storm将只能运行三个worker. 默认情况下,此设置被配置为在端口6700,6701,6702和6703上运行4个worker:
supervisor.slots.ports:
    - 6700
    - 6701
    - 6702
    - 6703 

以上是一些配置的介绍,下面我们的具体配置如下:

 storm.zookeeper.servers:
    - "hadoop000"
 storm.local.dir: "/home/hadoop/app/storm-2.1.0/data"
 nimbus.seeds: ["hadoop000"]
 storm.zookeeper.port: 2181
 supervisor.slots.ports:
    - 6700
    - 6701
    - 6702
    - 6703
 ui.port: 8082

注意:配置storm.zookeeper.servers前面有空格,等等细节很严格。另外ui服务进程的端口8080由于很容易和其他服务冲突,这里改成了8082

启动Storm进程

  • 主节点启动nimbus服务

    启动Nimbus。Run the command bin/storm nimbus under supervision on the master machine.

    storm nimbus &
  • 主节点启动UI服务:

    Run the Storm UI (a site you can access from the browser that gives diagnostics on the cluster and topologies) by running the command “bin/storm ui” under supervision. The UI can be accessed by navigating your web browser to http://{ui host}:8080. by default

    但是由于之前在配置文件中已经将UI端口修改成8082(避免和tomcat的8080冲突), 因此现在的访问方式是http://{ui host}:8082

    bin/storm ui &
  • 在从节点(工作节点)上启动supervisor服务(只不过目前是配置的伪分布式,所有服务都是在一个节点上,没有备用节点):

    storm supervisor &
  • 用jps判断是否启动成功(如果失败,则检查日志)

    img

  • 避免打印信息或者报错信息出现在屏幕上

    为了避免打印信息或者报错信息出现的屏幕上,也可以这样

    bin/storm ui  >/dev/null 2>&1 &  

    其中,2>&1 是将标准出错重定向到标准输出,但是最好的方式是采用nohup方式运行。

    nohup命令:如果你正在运行一个进程,而且你觉得在退出帐户时该进程还不会结束,那么可以使用nohup命令。该命令可以在你退出帐户/关闭终端之后继续运行相应的进程。nohup就是不挂起的意思( no hang up), 像这样:

    nohup ./bin/storm nimbus  > /dev/null  2>&1 & 
    nohup ./bin/storm supervisor > /dev/null 2>&1 & 
    nohup ./bin/storm ui > /dev/null  2>&1 &

查看UI界面

访问http://hadoop000:8082/ (端口是在配置文件中配置过的)

img

img

img

Numbus的配置超级繁杂, 有14页之多

查看log

log日志在storm包的logs文件夹中

如果想要在UI界面点击主机名或者端口号查看日志,则对应的机器需要开启logviewer服务

storm logviewer &

配置集群

假设当前机器hadoop001作为主节点, 如果要增加从节点机器hadoop002和hadoop003`,只需要修改如下:

storm.zookeeper.servers:
     - "hadoop001"
     - "hadoop002"
     - "hadoop003"
 storm.local.dir: "/home/hadoop/app/storm-2.1.0/data"
 nimbus.seeds: ["hadoop001"]
 storm.zookeeper.port: 2181
 supervisor.slots.ports:
   - 6700
   - 6701
   - 6702
   - 6703
 ui.port: 8082  

给bin目录的文件添加可执行权限:

bin>   chmod u+x *

开启nimbus和ui后台进程,把整个安装文件夹使用scp命令复制到其他两台机器(需要提前配置好通信机器之间的ssh免密访问),清空logsdata文件夹里面的内容,给bin目录的文件添加可执行权限,开启supervisor进程,这样就准备好集群环境了。

打开UI界面查看 可以看到Nimbus是hadoop001, 而supervisor是hadoop002和hadoop003

Storm UI

Cluster Summary

Version Supervisors Used slots Free slots Total slots Executors Tasks
2.1.0 2 0 8 8 0 0

Nimbus Summary

Search:

Host Port Status Version Uptime
hadoop001 6627 Leader 2.1.0 9m 16s

Showing 1 to 1 of 1 entries

Owner Summary

Search:

Owner Total Topologies Total Executors Total Workers Memory Usage (MB)
No data available in table

Showing 0 to 0 of 0 entries

Topology Summary

Search:

Name Owner Status Uptime Num workers Num executors Num tasks Replication count Assigned Mem (MB) Scheduler Info Topology Version Storm Version
No data available in table

Showing 0 to 0 of 0 entries

Supervisor Summary

Search:

Host Id Uptime Slots Used slots Avail slots Used Mem (MB) Version
hadoop002 (log) d8412147-0dbb-4933-a2d3-1bb28c806dc2-192.168.186.102 10m 33s 4 0 4 0 2.1.0
hadoop003 (log) 57798559-c338-4844-9e29-eb9aa99fa2fe-192.168.186.103 1m 7s 4 0 4 0 2.1.0

Showing 1 to 2 of 2 entries

Nimbus Configuration

Show entries

Search:

Key Value
blacklist.scheduler.reporter "org.apache.storm.scheduler.blacklist.reporters.LogReporter"
blacklist.scheduler.resume.time.secs 1800
blacklist.scheduler.strategy "org.apache.storm.scheduler.blacklist.strategies.DefaultBlacklistStrategy"
blacklist.scheduler.tolerance.count 3
blacklist.scheduler.tolerance.time.secs 300
client.blobstore.class "org.apache.storm.blobstore.NimbusBlobStore"
dev.zookeeper.path "/tmp/dev-storm-zookeeper"
drpc.authorizer.acl.filename "drpc-auth-acl.yaml"
drpc.authorizer.acl.strict false
drpc.childopts "-Xmx768m"
drpc.disable.http.binding true
drpc.http.creds.plugin "org.apache.storm.security.auth.DefaultHttpCredentialsPlugin"
drpc.http.port 3774
drpc.https.keystore.password ""
drpc.https.keystore.type "JKS"
drpc.https.port -1
drpc.invocations.port 3773
drpc.invocations.threads 64
drpc.max_buffer_size 1048576
drpc.port 3772

Showing 1 to 20 of 266 entries

为nimbus配置高可用

最后, 尽管前面我们使用了三台机器搭建了集群,由于nimbus服务没有配置高可用, 实际上还是算伪分布式。为nimbus配置高可用也非常简单, 比如说让hadoop001hadoop002其中一个作为主节点,另一个作为备用主节点,则配置如下:

nimbus.seeds: ["hadoop001", "hadoop002"]

Nimbus Summary

Search:

Host Port Status Version Uptime
hadoop001 6627 Leader 2.1.0 9m 16s
hadoop002 6627 Follower 2.1.0 ...

作业

使用Centos7系统搭建Storm集群, 将搭建过程形成笔记,并附上截图。

练习

在互联网上查询Twitter是如何使用storm进行情感分析的,了解storm的重要性以及用途。

Views: 895

1- 初识Storm

什么是Storm?

Storm为分布式实时计算提供了一组通用原语,可被用于“流处理”之中,实时处理消息并更新数据库。这是管理队列及工作者集群的另一种方式。 Storm也可被用于“连续计算”(continuous computation),对数据流做连续查询,在计算时就将结果以流的形式输出给用户。它还可被用于“分布式RPC”,以并行的方式计算。

Storm可以方便地在一个计算机集群中编写与扩展复杂的实时计算,Storm用于实时处理,就好比 Hadoop 用于批处理。Storm保证每个消息都会得到处理,而且它很快——在一个小集群中,每秒可以处理数以百万计的消息。更棒的是你可以使用任意编程语言来做开发。

离线计算和流式计算

离线计算

  • 离线计算:批量获取数据、批量传输数据、周期性批量计算数据、数据展示

  • 代表技术:Sqoop批量导入数据、HDFS批量存储数据、MapReduce批量计算、Hive

流式计算

(终极目的:留住用户,提升用户体验)

  • 流式计算:数据实时产生、数据实时传输、数据实时计算、实时展示

  • 代表技术:Flume实时获取数据、Kafka实时数据存储、Storm实时数据计算、Redis实时结果缓存。

  • 一句话总结:将源源不断产生的数据实时收集并实时计算,尽可能快的得到计算结果

image-20210830005200380

Storm与Hadoop的区别

Storm Hadoop
Storm用于实时计算 Hadoop用于离线计算
Storm的数据保存在内存,源源不断 Hadoop的数据保存在文件系统中,一批一批
Storm的数据通过网络传输进来 Hadoop的数据保存在磁盘中
Storm与Hadoop的编程模型相似 Storm与Hadoop的编程模型相似

Storm的体系结构

img

http://blog.chinaunix.net/attachment/201309/13/790245_13790273614l7w.jpg

  • Nimbus:负责资源分配和任务调度。

  • Supervisor:负责接受nimbus分配的任务,启动和停止属于自己管理的worker进程。通过配置文件设置当前supervisor上启动多少个worker。

  • Worker:运行具体处理组件逻辑的进程。Worker运行的任务类型只有两种,一种是Spout任务,一种是Bolt任务。

  • Executor:Storm 0.8之后,Executor为Worker进程中的具体的物理线程,同一个Spout/Bolt的Task可能会共享一个物理线程,一个Executor中只能运行隶属于同一个Spout/Bolt的Task。

  • Task:worker中每一个spout/bolt的线程称为一个task. 在storm0.8之后,task不再与物理线程对应,不同spout/bolt的task可能会共享一个物理线程,该线程称为executor。

storm 简介及单机版安装指南

Storm的运行机制

img

  • 整个处理流程的组织协调不用用户去关心,用户只需要去定义每一个步骤中的具体业务处理逻辑

  • 具体执行任务的角色是Worker,Worker执行任务时具体的行为则有我们定义的业务逻辑决定

img

参考

官方网站: http://storm.apache.org/
官方文档: http://storm.apache.org/releases/2.1.0/index.html

Views: 274

Spark RDD 实战案例

IDEA 创建基于Maven的Spark项目

三台虚拟机搭建的集群
启动集成如下:

------------ hadoop100 ------------
26369 Master
2514 QuorumPeerMain
2898 Kafka
29011 SparkSubmit
30643 HRegionServer
3493 JobHistoryServer
3573 NodeManager
3065 NameNode
30476 HMaster
31036 Jps
28927 SparkSubmit
------------ hadoop101 ------------
5153 QuorumPeerMain
5537 Kafka
5905 ResourceManager
6018 NodeManager
104355 Jps
84258 Worker
104039 HRegionServer
5645 DataNode
104140 HMaster
------------ hadoop102 ------------
71683 Worker
5509 Kafka
104804 Jps
5752 SecondaryNameNode
5897 NodeManager
5643 DataNode
5116 QuorumPeerMain
104542 HRegionServer

本地环境由于需要打包放到虚拟机运行,因此scala版本需要和虚拟机中编译spark所用的scala的版本一致。如何知道虚拟机中spark是用的什么版本scala编译的呢?可以进入虚拟机的spark-shell查看:

Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /___/ .__/\_,_/_/ /_/\_\   version 2.4.8
      /_/

Using Scala version 2.11.12 (Java HotSpot(TM) 64-Bit Server VM, Java 1.8.0_281)

可见2.4.8并不是官网所说的使用scala2.12,实际上是2.11

接下来使用IDEA创建一个MAVEN项目,使用scala模板
file
配置项目名称
file

配置有效的maven环境
file

修改pom文件:

  1. 设置scala的版本,我虚拟机里的Spark是用的scala2.11版本编译的,这里能找到最接近的可用的版本就是2.11.8。
  2. 添加maven-scala-plugin依赖,使用的是2.11的版本。
  3. 添加spark-core_2.11, 版本为2.4.8,和虚拟机安装的Spark版本保持一致
  4. 此外因为下面的示例还要操作hbase,因为又引入了hadoop和hbase相关的一些依赖,其中hadoop-common依赖排除了jackson-databind是因为出现了依赖冲突的问题。
  5. 为了演示官方示例(蒙特卡洛法求Pi值),引入了spark-sql依赖

最终完整的pom.xml文件内容如下:

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/maven-v4_0_0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <groupId>cn.delucia</groupId>
    <artifactId>SparkProject</artifactId>
    <version>1.0-SNAPSHOT</version>
    <inceptionYear>2008</inceptionYear>
    <properties>
        <scala.version>2.11.8</scala.version>
        <spark.version>2.4.8</spark.version>
        <hadoop.version>3.1.4</hadoop.version>
        <hbase.version>2.2.3</hbase.version>
        <maven.compiler.source>1.8</maven.compiler.source>
        <maven.compiler.target>1.8</maven.compiler.target>
        <encoding>UTF-8</encoding>
    </properties>

    <repositories>
        <repository>
            <id>scala-tools.org</id>
            <name>Scala-Tools Maven2 Repository</name>
            <url>http://scala-tools.org/repo-releases</url>
        </repository>
    </repositories>
    <pluginRepositories>
        <pluginRepository>
            <id>scala-tools.org</id>
            <name>Scala-Tools Maven2 Repository</name>
            <url>http://scala-tools.org/repo-releases</url>
        </pluginRepository>
    </pluginRepositories>

    <dependencies>
        <dependency>
            <groupId>org.apache.hadoop</groupId>
            <artifactId>hadoop-common</artifactId>
            <version>${hadoop.version}</version>
            <exclusions>
                <exclusion>
                    <groupId>com.fasterxml.jackson.core</groupId>
                    <artifactId>jackson-databind</artifactId>
                </exclusion>
            </exclusions>
        </dependency>
        <dependency>
            <groupId>org.apache.hadoop</groupId>
            <artifactId>hadoop-client</artifactId>
            <version>${hadoop.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.hbase</groupId>
            <artifactId>hbase-client</artifactId>
            <version>${hbase.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.hbase</groupId>
            <artifactId>hbase-server</artifactId>
            <version>${hbase.version}</version>
            <!--
            解决打包错误:Failure to find org.glassfish:javax.el:pom:3.0.1-b08-SNAPSHOT
            -->
            <exclusions>
                <exclusion>
                    <groupId>org.glassfish</groupId>
                    <artifactId>javax.el</artifactId>
                </exclusion>
            </exclusions>
        </dependency>
        <dependency>
            <groupId>org.apache.hbase</groupId>
            <artifactId>hbase-mapreduce</artifactId>
            <version>${hbase.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.11</artifactId>
            <version>${spark.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-sql_2.11</artifactId>
            <version>${spark.version}</version>
        </dependency>
        <dependency>
            <groupId>org.scala-lang</groupId>
            <artifactId>scala-library</artifactId>
            <version>${scala.version}</version>
        </dependency>
        <dependency>
            <groupId>mysql</groupId>
            <artifactId>mysql-connector-java</artifactId>
            <version>5.1.41</version>
        </dependency>
        <dependency>
            <groupId>junit</groupId>
            <artifactId>junit</artifactId>
            <version>4.4</version>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.specs</groupId>
            <artifactId>specs</artifactId>
            <version>1.2.5</version>
            <scope>test</scope>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.scala-tools/maven-scala-plugin -->
        <dependency>
            <groupId>org.scala-tools</groupId>
            <artifactId>maven-scala-plugin</artifactId>
            <version>2.15.2</version>
        </dependency>
    </dependencies>

    <build>
        <sourceDirectory>src/main/scala</sourceDirectory>
        <testSourceDirectory>src/test/scala</testSourceDirectory>
        <plugins>
            <plugin>
                <groupId>org.scala-tools</groupId>
                <artifactId>maven-scala-plugin</artifactId>
                <version>2.15.2</version>
                <executions>
                    <execution>
                        <goals>
                            <goal>compile</goal>
                            <goal>testCompile</goal>
                        </goals>
                    </execution>
                </executions>
                <configuration>
                    <scalaVersion>${scala.version}</scalaVersion>
                    <args>
                        <arg>-target:jvm-1.8</arg>
                    </args>
                </configuration>
            </plugin>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-eclipse-plugin</artifactId>
                <version>2.5.1</version>
                <configuration>
                    <downloadSources>true</downloadSources>
                    <buildcommands>
                        <buildcommand>ch.epfl.lamp.sdt.core.scalabuilder</buildcommand>
                    </buildcommands>
                    <additionalProjectnatures>
                        <projectnature>ch.epfl.lamp.sdt.core.scalanature</projectnature>
                    </additionalProjectnatures>
                    <classpathContainers>
                        <classpathContainer>org.eclipse.jdt.launching.JRE_CONTAINER</classpathContainer>
                        <classpathContainer>ch.epfl.lamp.sdt.launching.SCALA_CONTAINER</classpathContainer>
                    </classpathContainers>
                </configuration>
            </plugin>
        </plugins>
    </build>
    <reporting>
        <plugins>
            <plugin>
                <groupId>org.scala-tools</groupId>
                <artifactId>maven-scala-plugin</artifactId>
                <configuration>
                    <scalaVersion>${scala.version}</scalaVersion>
                </configuration>
            </plugin>
        </plugins>
    </reporting>
</project>

项目右键选择 maven->Reimport 等待依赖下载完成。
你会发现test包下面有一些自动生成的源文件有错误(依赖版本问题),这里直接把报错的文件删除,只保留一个AppTest.scala文件,可以点击测试以下环境,main包下面自动生成的文件也是直接删除即可。:

file

给项目根目录下创建data和output文件夹,最终项目结构:

file

实现单词计数

在data下面创建文本文件words.txt,内容如下

hello hadoop
hello java
scala

本地模式运行

package cn.delucia.spark.rdd

import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}

/**
  * Spark RDD 单词计数程序
  * Program Args: data\words.txt output\wc
  * 注意:需提前删除output\wc文件夹
  */
object WordCountLocal {

  def main(args: Array[String]): Unit = {
    //创建SparkConf对象
    val conf = new SparkConf()
    //设置应用程序名称,可以在Spark WebUI中显示
    conf.setAppName("WordCountLocal")

    //设置集群Master节点访问地址
    //conf.setMaster("spark://hadoop100:7077"); //远程-提交到集群(standalone模式)
    conf.setMaster("local[2]")  // 本地模式 - 2个CPU核心

    //创建SparkContext对象,该对象是提交Spark应用程序的入口
    val sc = new SparkContext(conf);

    //读取指定路径(取程序执行时传入的第一个参数)中的文件内容,生成一个RDD集合
    val linesRDD:RDD[String] = sc.textFile(args(0))
    //将RDD数据按照空格进行切分并合并为一个新的RDD
    val wordsRDD:RDD[String] = linesRDD.flatMap(_.split(" "))
    //将RDD中的每个单词和数字1放到一个元组里,即(word,1)
    val paresRDD:RDD[(String, Int)] = wordsRDD.map((_,1))
    //对单词根据key进行聚合,对相同的key进行value的累加
    val wordCountsRDD:RDD[(String, Int)] = paresRDD.reduceByKey(_+_)
    //按照单词数量降序排列
    val wordCountsSortRDD:RDD[(String, Int)] = wordCountsRDD.sortBy(_._2,false)
    //保存结果到指定的路径(取程序执行时传入的第二个参数)
    wordCountsSortRDD.saveAsTextFile(args(1))
    //停止SparkContext,结束该任务
    sc.stop()
  }
}

直接点击运行,由于是两个核心,结果有两个文件生成
file

上传到集群运行 - Spark on YARN 模式

前提 集群已经安装好Hadoop集群并运行DFS和YARN
Spark已经配置好Spark on YARN的相关设置

编写程序

package cn.delucia.spark.rdd

import org.apache.spark.{SparkConf, SparkContext}

/**
 Spark RDD单词计数程序 - 打jar包上传到集群 Spark on YARN 模式
*/
object WordCountCluster {

  def main(args: Array[String]): Unit = {

    val conf = new SparkConf().setAppName("WordCountCluster")
    conf.setMaster("spark://hadoop100:7077")
    val sc = new SparkContext(conf)
    sc.textFile(args(0))
      .flatMap(_.split(" "))
      .map((_, 1))
      .reduceByKey(_ + _)
      .sortBy(_._2, ascending = false)
      .saveAsTextFile(args(1))

    sc.stop()
  }
}

使用maven生命周期插件clean然后install打成jar包上传到集群的/opt/data目录下,将jar包改名位WordCountCluster.jar,进入spark的安装目录执行命令:

bin/spark-submit
--master yarn
--class cn.delucia.spark.rdd.WordCountCluster /opt/data/WordCountCluster.jar hdfs://hadoop100:8020/tmp/words.txt hdfs://hadoop100:8020/tmp/output/wc_cluster

即可把任务提交到YARN集群运行,可以去output目录可以查看结果
file

有两个文件生成说明有两个分区,因为设置了2个核心

求平均成绩

data下创建score.txt

Andy,98
Jack,87
Bill,99
Andy,78
Jack,85
Bill,86
Andy,90
Jack,88
Bill,76
Andy,58
Jack,67
Bill,79

编写程序:

package cn.delucia.spark.rdd

import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}

/**
  * 求成绩平均分
  */
object AverageScoreLocal {
   def main(args: Array[String]): Unit = {
      //创建SparkConf对象,存储应用程序的配置信息
      val conf = new SparkConf()
      //设置应用程序名称,可以在Spark WebUI中显示
      conf.setAppName("AverageScore")
      //设置集群Master节点访问地址
      conf.setMaster("local")

      val sc = new SparkContext(conf)
      //1. 加载数据
      val linesRDD: RDD[String] = sc.textFile("data/avg_score.txt")
      //2. 将RDD中的元素转为(key,value)形式,便于后面进行聚合
      val tupleRDD: RDD[(String, Int)] = linesRDD.map(line => {
         val name = line.split("\t")(0)//姓名
         val score = line.split("\t")(1).toInt//成绩
         (name, score)
      })
      //3. 根据姓名进行分组,形成新的RDD
      val groupedRDD: RDD[(String, Iterable[Int])] = tupleRDD.groupByKey()
      //4. 迭代计算RDD中每个学生的平均分
      val resultRDD: RDD[(String, Int)] = groupedRDD.map(line => {
         val name = line._1//姓名
         val iteratorScore: Iterator[Int] = line._2.iterator//成绩迭代器
         var sum = 0//总分
         var count = 0//科目数量

         //迭代累加所有科目成绩
         while (iteratorScore.hasNext) {
            val score = iteratorScore.next()
            sum += score
            count += 1
         }
         //计算平均分
         val averageScore = sum / count
         (name, averageScore)//返回(姓名,平均分)形式的元组
      })
      //保存结果
      resultRDD.saveAsTextFile("output/avg_score")
   }
}

输出结果:

(BETA,45)
(Bill,85)
(Andy,81)
(Jack,81)

统计学生最好的三次成绩并倒序排列

package cn.delucia.spark.rdd

import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}
/**
  * Spark分组取TopN程序
  */
object GroupTopNLocal {
   def main(args: Array[String]): Unit = {
      //创建SparkConf对象,存储应用程序的配置信息
      val conf = new SparkConf()
      //设置应用程序名称,可以在Spark WebUI中显示
      conf.setAppName("RDDGroupTopN")
      //设置集群Master节点访问地址,此处为本地模式,使用尽可能多的核心
      conf.setMaster("local[*]")

      val sc = new SparkContext(conf)
      //1. 加载本地数据
      val linesRDD: RDD[String] = sc.textFile("data/score.txt")

      //2. 将RDD元素转为(String,Int)形式的元组
      val tupleRDD:RDD[(String,Int)]=linesRDD.map(line=>{
         val name=line.split(",")(0)
         val score=line.split(",")(1)
         (name,score.toInt)
      })

      //3. 按照key(姓名)进行分组
      val top3=tupleRDD.groupByKey().map(groupedData=>{
         val name:String=groupedData._1
         //每一组的成绩降序后取前3个
         val scoreTop3:List[Int]=groupedData._2
           .toList.sortWith(_>_).take(3)
         (name,scoreTop3)//返回元组
      })

      //当使用多核心时,分区排序完就各自并行输出结果了,导致输出内容互相干扰
      //这里先调用count,由于count属于Action算子,会触发前面的计算完成
      //这样等top3里面的所有分区都计算完毕再输出,可以让输出格式更好
      top3.count

      //4. 循环打印分组结果
      top3.foreach(tuple=>{
         println("姓名:"+tuple._1)
         val tupleValue=tuple._2.iterator
         while (tupleValue.hasNext){
            val value=tupleValue.next()
            println("成绩:"+value)
         }
         println("*******************")
      })
   }
}

控制台输出:

姓名:Andy
成绩:98
成绩:90
成绩:78
*******************
姓名:Bill
成绩:99
成绩:86
成绩:79
*******************
姓名:Jack
成绩:88
成绩:87
成绩:85
*******************

倒排索引统计每日新增用户

package cn.delucia.spark.rdd

import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}

/**
  * Spark RDD统计每日新增用户
  */
object DayNewUserLocal {
   def main(args: Array[String]): Unit = {
      val conf = new SparkConf()
      conf.setAppName("DayNewUser")
      conf.setMaster("local[*]")

      val sc = new SparkContext(conf)
      //1. 构建测试数据
      val tupleRDD:RDD[(String,String)] = sc.parallelize(
         Array(
            ("2020-01-01", "user1"),
            ("2020-01-01", "user2"),
            ("2020-01-01", "user3"),
            ("2020-01-02", "user1"),
            ("2020-01-02", "user2"),
            ("2020-01-02", "user4"),
            ("2020-01-03", "user2"),
            ("2020-01-03", "user5"),
            ("2020-01-03", "user6")
         )
      )
      //2. 倒排(互换RDD中元组的元素顺序)
      val tupleRDD2:RDD[(String,String)] = tupleRDD.map(
      line => (line._2, line._1)
      )
      //3. 将倒排后的RDD按照key分组
      val groupedRDD: RDD[(String, Iterable[String])] = tupleRDD2.groupByKey()
      //4. 取分组后的每个日期集合中的最小日期,并计数为1  (只有第一次登陆才算做当日新增用户)
      val dateRDD:RDD[(String,Int)] = groupedRDD.map(
         line => (line._2.min, 1)
      )
      //5. 计算所有相同key(即日期)的数量
      val resultMap: collection.Map[String, Long] = dateRDD.countByKey()
      //将结果Map循环打印到控制台
      resultMap.foreach(println)
   }

}

结果:

(2020-01-01,3)
(2020-01-02,1)
(2020-01-03,2)

自定义排序规则(二次排序)

先准备数组,data下创建文本文件

2 98
1 99
2 67
3 75
3 88
2 85
1 90
3 100
1 62

先根据第一列升序排列,如果相同再根据第二列降序排列

package cn.delucia.spark.rdd

import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}

/**
 * 二次排序自定义key类
 * 先根据名字升序排列,如果名字相同再根据成绩降序排列
 * @param first  每一行的第一个字段
 * @param second 每一行的第二个字段
 */
class SecondSortKey(val first: Int, val second: Int)
  extends Ordered[SecondSortKey] with Serializable {
  override def toString: String = first + "->" + second;

  /**
   * 实现compare()方法
   */
  override def compare(that: SecondSortKey): Int = {
    //若第一个字段不相等,按照第一个字段升序排列
    if (this.first - that.first != 0) {
      this.first - that.first
    } else { //否则按照第二个字段降序排列
      that.second - this.second
    }
  }
}

/**
 * 二次排序运行主类
 */
object SecondSort {
  def main(args: Array[String]): Unit = {
    //创建SparkConf对象
    val conf = new SparkConf()
    //设置应用程序名称,可以在Spark WebUI中显示
    conf.setAppName("SecondSort")
    //设置集群Master节点访问地址,此处为本地模式
    conf.setMaster("local")
    //创建SparkContext对象,该对象是提交Spark应用程序的入口
    val sc = new SparkContext(conf);

    //1. 读取指定路径的文件内容,生成一个RDD集合
    val lines: RDD[String] = sc.textFile("data\\sort.txt")
    //2. 将RDD中的元素转为(SecondSortKey, String)形式的元组
    val pair: RDD[(SecondSortKey, String)] = lines.map(line => (
      new SecondSortKey(line.split(" ")(0).toInt, line.split(" ")(1).toInt),
      line)
    )
    //3. 按照元组的key(SecondSortKey的实例)进行排序
    val pairSort: RDD[(SecondSortKey, String)] = pair.sortByKey()
    //取排序后的元组中的第二个值(value值)
    //val result: RDD[String] = pairSort.map(line => line._2)
    //打印最终结果
    pairSort.foreach(line => println(line))
  }
}

控制台输出结果:

(1->99,1 99)
(1->90,1 90)
(1->62,1 62)
(2->98,2 98)
(2->85,2 85)
(2->67,2 67)
(3->100,3 100)
(3->88,3 88)
(3->75,3 75)

读写HBase数据库

HBase是Spark应用程序经常打交道的一个数据源。下面使用Spark程序对其进行读写操作:

使用HBase API向HBase写入数据

package cn.delucia.spark.rdd

import org.apache.hadoop.hbase.HBaseConfiguration
import org.apache.hadoop.hbase.client.Put
import org.apache.spark.{SparkConf, SparkContext}
import org.apache.hadoop.hbase.TableName
import org.apache.hadoop.hbase.util.Bytes
import org.apache.hadoop.hbase.client.ConnectionFactory
/**
  * 向HBase表写入数据 - hbase API
  */
object SparkWriteHBase {
   def main(args: Array[String]): Unit = {
      //创建SparkConf对象,存储应用程序的配置信息
      val conf = new SparkConf()
      conf.setAppName("SparkWriteHBase")
      conf.setMaster("local[*]")
      //创建SparkContext对象
      val sc = new SparkContext(conf)

      //1. 构建需要添加的数据RDD
      val initRDD = sc.makeRDD(
         Array(
            "003,王五,山东,23",
            "004,赵六,河北,20"
         )
      )

      //2. 循环RDD的每个分区
      initRDD.foreachPartition(partition=> {
         //2.1 设置HBase配置信息
         val hbaseConf = HBaseConfiguration.create()
         //设置ZooKeeper集群地址
         hbaseConf.set("hbase.zookeeper.quorum","hadoop100")
         //设置ZooKeeper连接端口,默认2181
         hbaseConf.set("hbase.zookeeper.property.clientPort", "2181")
         //创建数据库连接对象
         val conn = ConnectionFactory.createConnection(hbaseConf)
         //指定表名
         val tableName = TableName.valueOf("student")
         //获取需要添加数据的Table对象
         val table = conn.getTable(tableName)

         //2.2 循环当前分区的每行数据
         partition.foreach(line => {
            //分割每行数据,获取要添加的每个值
            val arr = line.split(",")
            val rowkey = arr(0)
            val name = arr(1)
            val address = arr(2)
            val age = arr(3)

            //创建Put对象
            val put = new Put(Bytes.toBytes(rowkey))
            put.addColumn(
               Bytes.toBytes("info"),//列族名
               Bytes.toBytes("name"),//列名
               Bytes.toBytes(name)//列值
            )
            put.addColumn(
               Bytes.toBytes("info"),//列族名
               Bytes.toBytes("address"),//列名
               Bytes.toBytes(address)) //列值
            put.addColumn(
               Bytes.toBytes("info"),//列族名
               Bytes.toBytes("age"),//列名
               Bytes.toBytes(age)) //列值

            //执行添加
            table.put(put)
         })
      })

   }
}

使用Spark API向HBase写入数据

package cn.delucia.spark.rdd

import org.apache.hadoop.hbase.client.Put
import org.apache.hadoop.hbase.io.ImmutableBytesWritable
import org.apache.hadoop.hbase.mapred.TableOutputFormat
import org.apache.hadoop.hbase.util.Bytes
import org.apache.hadoop.mapred.JobConf
import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}
/**
  * 向HBase表写入数据 使用 Spark API(saveAsHadoopDataset方法)
  */
object SparkWriteHBase2 {
   def main(args: Array[String]): Unit = {
      //创建SparkConf对象,存储应用程序的配置信息
      val conf = new SparkConf()
      conf.setAppName("SparkWriteHBase2")
      conf.setMaster("local[*]")
      //创建SparkContext对象
      val sc = new SparkContext(conf)

      //1. 设置配置信息
      //创建Hadoop JobConf对象
      val jobConf = new JobConf()
      //设置ZooKeeper集群地址
      jobConf.set("hbase.zookeeper.quorum","hadoop100")
      //设置ZooKeeper连接端口,默认2181
      jobConf.set("hbase.zookeeper.property.clientPort", "2181")
      //指定输出格式
      jobConf.setOutputFormat(classOf[TableOutputFormat])
      //指定表名
      jobConf.set(TableOutputFormat.OUTPUT_TABLE,"student")

      //2. 构建需要写入的RDD数据
      val initRDD = sc.makeRDD(
         Array(
            "005,王五,山东,23",
            "006,赵六,河北,20"
         )
      )

      //将RDD转换为(ImmutableBytesWritable, Put)类型
      val resultRDD: RDD[(ImmutableBytesWritable, Put)] = initRDD.map(
         _.split(",")
      ).map(arr => {
         val rowkey = arr(0)
         val name = arr(1)//姓名
         val address = arr(2)//地址
         val age = arr(3)//年龄

         //创建Put对象
         val put = new Put(Bytes.toBytes(rowkey))
         put.addColumn(
            Bytes.toBytes("info"),//列族
            Bytes.toBytes("name"),//列名
            Bytes.toBytes(name)//列值
         )
         put.addColumn(
            Bytes.toBytes("info"),//列族
            Bytes.toBytes("address"),//列名
            Bytes.toBytes(address)) //列值
         put.addColumn(
            Bytes.toBytes("info"),//列族
            Bytes.toBytes("age"),//列名
            Bytes.toBytes(age)) //列值

         //拼接为元组返回
         (new ImmutableBytesWritable, put)
      })

      //3. 写入数据
      resultRDD.saveAsHadoopDataset(jobConf)
      sc.stop()
   }
}

批量向Hbase写入数据

使用BulkLoader加载HFile数据到HBase:
参考阅读:
https://blog.cloudera.com/how-to-use-hbase-bulk-loading-and-why/
https://www.opencore.com/blog/2016/10/efficient-bulk-load-of-hbase-using-spark/

package cn.delucia.spark.rdd

import org.apache.hadoop.conf.Configuration
import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.hadoop.hbase._
import org.apache.hadoop.hbase.client.ConnectionFactory
import org.apache.hadoop.hbase.io.ImmutableBytesWritable
import org.apache.hadoop.hbase.mapred.TableOutputFormat
import org.apache.hadoop.hbase.mapreduce.{HFileOutputFormat2, LoadIncrementalHFiles}
import org.apache.hadoop.hbase.util.Bytes
import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}

/**
  * Spark 批量写入数据到HBase(使用BulkLoader加载HFile数据到HBase)
  */
object SparkWriteHBase3 {
   def main(args: Array[String]): Unit = {
      System.setProperty("HADOOP_USER_NAME", "hadoop")
      //创建SparkConf对象,存储应用程序的配置信息
      val conf = new SparkConf()
      conf.setAppName("SparkWriteHBase3")
      conf.setMaster("local[*]")
      //创建SparkContext对象
      val sc = new SparkContext(conf)

      //1. 设置HDFS和HBase配置信息
      val hadoopConf = new Configuration()
      hadoopConf.set("fs.defaultFS", "hdfs://hadoop100:8020")
      val fileSystem = FileSystem.get(hadoopConf)
      val hbaseConf = HBaseConfiguration.create(hadoopConf)
      //设置ZooKeeper集群地址
      hbaseConf.set("hbase.zookeeper.quorum", "hadoop100")
      //设置ZooKeeper连接端口,默认2181
      hbaseConf.set("hbase.zookeeper.property.clientPort", "2181")
      //创建数据库连接对象
      val conn = ConnectionFactory.createConnection(hbaseConf)
      //指定表名
      val tableName = TableName.valueOf("student")
      //获取需要添加数据的Table对象
      val table = conn.getTable(tableName)
      //获取操作数据库的Admin对象
      val admin = conn.getAdmin()

      //2. 添加数据前的判断
      //如果HBase表不存在,则创建一个新表
      if (!admin.tableExists(tableName)) {
         val desc = new HTableDescriptor(tableName)
         //表名
         val hcd = new HColumnDescriptor("info") //列族
         desc.addFamily(hcd)
         admin.createTable(desc) //创建表
      }
      //如果存放HFile文件的HDFS目录已经存在,则删除
      if (fileSystem.exists(new Path("hdfs://hadoop100:8020/tmp/hbase"))) {
         fileSystem.delete(new Path("hdfs://hadoop100:8020/tmp/hbase"), true)
      }

      //3. 构建需要添加的RDD数据
      //初始数据
      val initRDD = sc.makeRDD(
         Array(
            "rowkey:007,name:王五",
            "rowkey:007,address:山东",
            "rowkey:007,age:23",
            "rowkey:008,name:赵六",
            "rowkey:008,address:河北",
            "rowkey:008,age:20"
         )
      )
      //数据转换
      //转换为(ImmutableBytesWritable, KeyValue)类型的RDD
      val resultRDD: RDD[(ImmutableBytesWritable, KeyValue)] = initRDD.map(
         _.split(",")
      ).map(arr => {
         val rowkey = arr(0).split(":")(1)
         //rowkey
         val qualifier = arr(1).split(":")(0)
         //列名
         val value = arr(1).split(":")(1) //列值

         val kv = new KeyValue(
            Bytes.toBytes(rowkey),
            Bytes.toBytes("info"),
            Bytes.toBytes(qualifier),
            Bytes.toBytes(value)
         )
         //构建(ImmutableBytesWritable, KeyValue)类型的元组返回
         (new ImmutableBytesWritable(Bytes.toBytes(rowkey)), kv)
      })

      //4. 写入数据
      //在HDFS中生成HFile文件
      hbaseConf.set("hbase.mapreduce.hfileoutputformat.table.name",tableName.getNameAsString)
      resultRDD.saveAsNewAPIHadoopFile(
         "hdfs://hadoop100:8020/tmp/hbase",
         classOf[ImmutableBytesWritable], //对应RDD元素中的key
         classOf[KeyValue], //对应RDD元素中的value
         classOf[HFileOutputFormat2],
         hbaseConf
      )
      //加载HFile文件到HBase
      val bulkLoader = new LoadIncrementalHFiles(hbaseConf)
      val regionLocator = conn.getRegionLocator(tableName)
      bulkLoader.doBulkLoad(
         new Path("hdfs://hadoop100:8020/tmp/hbase"), //HFile文件位置
         admin, //操作HBase数据库的Admin对象
         table, //目标Table对象(包含表名)
         regionLocator //RegionLocator对象,用于查看单个HBase表的区域位置信息
      )
      sc.stop()
   }
}

读取Hbase数据

前面学了三钟方法向HBase数据库插入数据,并且一共向Hbase的student表写入了6条记录,下面我们写一个程序从HBase中把它们读取出来

package cn.delucia.spark.rdd

import org.apache.spark.{SparkConf, SparkContext}
import org.apache.hadoop.hbase.HBaseConfiguration
import org.apache.hadoop.hbase.client.Result
import org.apache.hadoop.hbase.io.ImmutableBytesWritable
import org.apache.hadoop.hbase.mapreduce.TableInputFormat
import org.apache.hadoop.hbase.util.Bytes
import org.apache.spark.rdd.RDD
/**
  * Spark读取HBase表数据
  */
object SparkReadHBase {
   def main(args: Array[String]): Unit = {
      //创建SparkConf对象,存储应用程序的配置信息
      val conf = new SparkConf()
      conf.setAppName("SparkReadHBase")
      conf.setMaster("local[*]")
      //创建SparkContext对象
      val sc = new SparkContext(conf)

      //1. 设置HBase配置信息
      val hbaseConf = HBaseConfiguration.create()
      //设置ZooKeeper集群地址
      hbaseConf.set("hbase.zookeeper.quorum","hadoop100")
      //设置ZooKeeper连接端口,默认2181
      hbaseConf.set("hbase.zookeeper.property.clientPort", "2181")
      //指定表名
      hbaseConf.set(TableInputFormat.INPUT_TABLE, "student")

      //2. 读取HBase表数据并转化成RDD
      val hbaseRDD: RDD[(ImmutableBytesWritable, Result)] = sc.newAPIHadoopRDD(
         hbaseConf,
         classOf[TableInputFormat],
         classOf[ImmutableBytesWritable],
         classOf[Result]
      )

      //3. 输出RDD中的数据到控制台
      hbaseRDD.foreach{ case (_ ,result) =>
         //获取行键
         val key = Bytes.toString(result.getRow)
         //通过列族和列名获取列值
         val name = Bytes.toString(result.getValue("info".getBytes,"name".getBytes))
         val gender = Bytes.toString(result.getValue("info".getBytes,"address".getBytes))
         val age = Bytes.toString(result.getValue("info".getBytes,"age".getBytes))
         println("行键:"+key+"\t姓名:"+name+"\t地址:"+gender+"\t年龄:"+age)
      }
   }
}

输出

行键:003  姓名:王五   地址:山东   年龄:23
行键:004  姓名:赵六   地址:河北   年龄:20
行键:005  姓名:王五   地址:山东   年龄:23
行键:006  姓名:赵六   地址:河北   年龄:20
行键:007  姓名:王五   地址:山东   年龄:23
行键:008  姓名:赵六   地址:河北   年龄:20

解决数据倾斜问题

避免数据倾斜的办法还有很多,比如:

  1. 对数据进行预处理
  2. 过滤掉没有意义的数据
  3. 提高shuffle的并行度,增加分区数量
    但这不能解决大量相同key导致的单个分区数据过多的问题

对于存在大量相同的key导致数据集中在某个分区的问题,这里给出一个解决办法:

添加随机前缀进行双重集合:

  1. 首先给数据添加随机前缀,进行分区内的局部聚合
  2. 然后去除随机前缀,进行全局聚合
package cn.delucia.spark.rdd

import org.apache.spark.{SparkConf, SparkContext}
import scala.util.Random

/**
  * Spark RDD解决数据倾斜案例 - 由于存在大量相同的key导致数据集中在某个分区
 *  解决办法:添加随机前缀进行双重集合)
  * 1. 首先给数据添加随机前缀,进行分区内的局部聚合
  * 2. 然后去除随机前缀,进行全局聚合
  */
object DataLean {
   def main(args: Array[String]): Unit = {
      //创建Spark配置对象
      val conf = new SparkConf();
      conf.setAppName("DataLean")
      conf.setMaster("local[*]")

      //创建SparkContext对象
      val sc = new SparkContext(conf)

      //1. 读取测试数据
      val linesRDD = sc.textFile("data/data.txt")
      //2. 统计单词数量
      linesRDD
        .flatMap(_.split(" "))
        .map((_, 1))
        .map(t => {
           val word = t._1
           val random = Random.nextInt(100)//产生0~99的随机数
           //单词加入随机数前缀,格式:(前缀_单词,数量)
           (random + "_" + word, 1)
        })
        .reduceByKey(_ + _)//局部聚合
        .map(t => {
         val word = t._1
         val count = t._2
         val w = word.split("_")(1)//去除前缀
         //单词去除随机数前缀,格式:(单词,数量)
         (w, count)
      })
        .reduceByKey(_ + _)//全局聚合
        //输出结果到指定的HDFS目录
        .saveAsTextFile("output/data_lean")
   }
}

Views: 98

Spark核心源码分析

01 Spark集群启动原理

Spark集群启动时,会在当前节点(脚本执行节点)上启动Master,在配置文件conf/slave中指定的每个节点上启动一个Worker。而Spark集群是通过脚本启动的。
查看启动脚本的内容:

$ cat sbin/start-all.sh

内容:

#!/usr/bin/env bash
#启动所有spark守护进程
#在此节点上启动Master
#在conf/slave中指定的每个节点上启动一个Worker
 #如果SPARK_HOME环境变量为空,则使用export设置
if [ -z "${SPARK_HOME}" ]; then
  export SPARK_HOME="$(cd "<code>dirname "$0""/..; pwd)"
fi
#加载Spark配置
. "${SPARK_HOME}/sbin/spark-config.sh"
#启动Master
"${SPARK_HOME}/sbin"/start-master.sh
#启动Workers
"${SPARK_HOME}/sbin"/start-slaves.sh

从启动脚本的内容和注释可以看出,start-all.sh脚本主要做了四件事:

(1)检查并设置环境变量
(2)加载Spark配置
Spark配置的加载,使用了脚本文件sbin/spark-config.sh。
(3)启动Master进程
Master进程的启动,使用了脚本文件sbin/start-master.sh。
(4)启动Worker进程
Worker进程的启动,使用了脚本文件sbin/start-slaves.sh,该文件负责批量启动集群中各个节点上的Worker进程。

02 Spark应用程序提交原理

Spark应用程序的提交入口是spark-submit脚本。除此之外,也可以使用spark-shell脚本和spark-sql脚本进入交互式命令行执行Spark程序。
查看spark-shell脚本的内容:

$ cat bin/spark-shell

内容关键代码:

#!/usr/bin/env bash
......
function main() {
  if $cygwin; then
    stty -icanon min 1 -echo > /dev/null 2>&1
    export SPARK_SUBMIT_OPTS="$SPARK_SUBMIT_OPTS -Djline.terminal=unix"
"${SPARK_HOME}"/bin/spark-submit --class org.apache.spark.repl.Main --name
 "Spark shell" "$@"
    stty icanon echo > /dev/null 2>&1
  else
    export SPARK_SUBMIT_OPTS
"${SPARK_HOME}"/bin/spark-submit --class org.apache.spark.repl.Main --name
 "Spark shell" "$@"
  fi
}
......

从关键代码可以看出,spark-shell脚本调用了应用程序提交脚本spark-submit,并传入了org.apache.spark.repl.Main类和使用spark-shell命令时的所有参数(例如--master)。通过org.apache.spark.repl.Main类启动的spark-shell将进入REPL模式(交互式编程环境),一直等待用户的输入,除非手动退出。
查看spark-sql脚本的内容:

$ cat bin/spark-sql

代码如下:

#!/usr/bin/env bash
if [ -z "${SPARK_HOME}" ]; then
  source "$(dirname "$0")"/find-spark-home
fi

export _SPARK_CMD_USAGE="Usage: ./bin/spark-sql [options] [cli option]"
exec "${SPARK_HOME}"/bin/spark-submit --class 
org.apache.spark.sql.hive.thriftserver.SparkSQLCLIDriver "$@"

从上述代码可以看出,spark-sql脚本同样调用了应用程序提交脚本spark-submit。这说明,无论采用哪种方式提交应用程序,都间接的调用了spark-submit脚本。查看spark-submit脚本的内容,发现最后一行有如下代码:

exec "${SPARK_HOME}"/bin/spark-class org.apache.spark.deploy.SparkSubmit "$@“

可以看出,最后使用exec命令执行了脚本bin/spark-class,并传入了org.apache.spark.deploy.SparkSubmit类和使用spark-submit命令时的所有参数。脚本bin/spark-class用于执行Spark应用程序并启动JVM进程,而这里执行的正是SparkSubmit类。

从上面的分析可以总结出,若通过spark-shell进行启动,将默认传入以下参数:

org.apache.spark.deploy.SparkSubmit
--class org.apache.spark.repl.Main
--name "Spark shell"

若通过spark-submit提交应用程序,将默认传入以下参数:

org.apache.spark.deploy.SparkSubmit

03 Spark作业工作原理

在学习Spark作业的工作原理时,通常会引入Hadoop MapReduce的工作原理作为入门比较,因为MapReduce与Spark的工作原理有很多相似之处。

1、MapReduce工作原理

MapReduce计算模型主要由三个阶段组成:Map阶段、Shuffle阶段、Reduce阶段。

file

(1)Map阶段

将输入的多个分片(Split)由Map任务以完全并行的方式处理。每个分片由一个Map任务来处理。默认情况下,输入分片的大小与HDFS中数据块(Block)的大小是相同的,即文件有多少个数据块就有多少个输入分片,也就会有多少个Map任务。从而可以调整HDFS数据块的大小来间接改变Map任务的数量。

每个Map任务对输入分片中的记录按照一定的规则解析成多个<key,value>对。默认将文件中的每一行文本内容解析成一个<key,value>对,key为每一行的起始位置,value为本行的文本内容。然后将解析出的所有<key,value>对分别输入到map()方法中进行处理(map()方法一次只处理一个<key,value>对)。map()方法将处理结果仍然是以<key,value>对的形式进行输出。

在数据溢写到磁盘之前,会对数据进行分区(Partition)。分区的数量与设置的Reduce任务的数量相同(默认Reduce任务的数量为1,可以在编写MapReduce程序时对其修改)。这样每个Reduce任务会处理一个分区的数据,可以防止有的Reduce任务分配的数据量太大,而有的Reduce任务分配的数据量太小,从而可以负载均衡,避免数据倾斜。数据分区的划分规则为:取<key,value>对中key的hashCode值,然后除以Reduce任务数量后取余数,余数则是分区编号,分区编号一致的<key,value>对则属于同一个分区。因此,key值相同的<key,value>对一定属于同一个分区,但是同一个分区中可能有多个key值不同的<key,value>对。由于默认Reduce任务的数量为1,而任何数字除以1的余数总是0,因此分区编号从0开始。

(2)Reduce阶段

Reduce阶段首先会对Map阶段的输出结果按照分区进行再一次合并,将同一分区的<key,value>对合并到一起,然后按照key对分区中的<key,value>对进行排序。
每个分区会将排序后的<key,value>对按照key进行分组,key相同的<key,value>对将合并为<key,value-list>对,最终每个分区形成多个<key,value-list>对。例如,key中存储的是用户ID,则同一个用户的<key,value>对会合并到一起。
排序并分组后的分区数据会输入到reduce()方法中进行处理,reduce()方法一次只能处理一个<key,value-list>对。
最后,reduce()方法将处理结果仍然以<key,value>对的形式通过context.write(key,value)进行输出。

(3)Shuffle阶段

Shuffle阶段所处的位置是Map任务输出后,Reduce任务接收前。主要是将Map任务的无规则输出形成一定的有规则数据,以便Reduce任务进行处理。
总结来说,MapReduce的工作原理主要是:通过Map任务读取HDFS中的数据块,这些数据块由Map任务以完全并行的方式处理;然后将Map任务的输出进行排序后输入到Reduce任务中;最后Reduce任务将计算的结果输出到HDFS文件系统中。

2、Spark工作原理

典型的Spark作业(Job)的工作流程如图:

file

(1)从数据源(本地文件、HDFS、HBase等)读取数据并创建RDD。
(2)对RDD进行一系列的转化操作。
(3)对最终RDD执行行动操作,开始一系列的计算,产生计算结果。
(4)将计算结果发送到Driver端,进行查看和输出。

04 Spark检查点原理

通过调用SparkContext的setCheckpointDir()方法指定检查点数据的存储路径,即把RDD数据放到什么位置,通常是放到HDFS中(如果在集群中运行,则必须是HDFS目录)。

总结来说,RDD检查点的运行流程如下:
(1)通过SparkContext的setCheckpointDir()方法设置RDD检查点数据的存储路径。
(2)调用RDD的checkpoint()方法将RDD标记为检查点。
(3)当在RDD上运行一个作业后,会立即触发RDDCheckpointData中的checkpoint()方法。
(4)在checkpoint()方法中又执行了doCheckpoint()方法。
(5)doCheckpoint()方法中执行了writeRDDToCheckpointDirectory()方法。
(6)writeRDDToCheckpointDirectory()方法内部通过调用runJob()方法运行一个作业,真正将RDD数据写入到检查点目录中,写入完成后返回一个ReliableCheckpointRDD实例。

Views: 105

Spark RDD弹性分布式数据集

01 什么是RDD

Spark提供了一种对数据的核心抽象,称为弹性分布式数据集(Resilient Distributed Dataset,简称RDD)。这个数据集的全部或部分可以缓存在内存中,并且可以在多次计算时重用。RDD其实就是一个分布在多个节点上的数据集合。

RDD的弹性主要是指:当内存不够时,数据可以持久化到磁盘,并且RDD具有高效的容错能力。分布式数据集是指:一个数据集存储在不同的节点上,每个节点存储数据集的一部分。

例如,将数据集(hello,world,scala,spark,love,spark,happy)存储在三个节点上,节点一存储(hello,world),节点二存储(scala,spark,love),节点三存储(spark,happy),这样对三个节点的数据可以并行计算,并且三个节点的数据共同组成了一个RDD。

file

分布式数据集类似于HDFS中的文件分块,不同的块存储在不同的节点上;而并行计算类似使用MapReduce读取HDFS中的数据并进行Map和Reduce操作。Spark则包含了这两种功能,并且计算更加灵活。

RDD的主要特征如下:

  • RDD是不可变的,但是可以将RDD转换成新的RDD进行操作
  • RDD是可分区的,一个RDD可以由很多分区组成(数据存在多个节点上),每个分区对应一个Task来执行
  • 对RDD的操作是作用在RDD的每个分区上
  • RDD拥有一系列对分区进行计算的函数,成为算子
  • RDD之间存在依赖关系,可以实现管道化,避免中间数据的存储

在编程时,可以把RDD看作是一个数据操作的基本单位,而不必关心数据的分布式特性,Spark会自动将RDD的数据分发到集群的各个节点。Spark中对数据的操作主要是对RDD的操作(创建、转化、求值)。

1、从对象集合创建RDD

Spark可以通过parallelize()或makeRDD()方法将一个对象集合转化为RDD。
例如,将一个List集合转化为RDD:

scala> val rdd=sc.parallelize(List(1,2,3,4,5,6))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0]

scala> val rdd=sc.makeRDD(List(1,2,3,4,5,6))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[1]

从返回信息可以看出,上述创建的RDD中存储的是Int类型数据。实际上RDD也是一个集合(并行集合),与常用的List集合不同的是,RDD集合的数据分布于多台机器上。

2、从外部存储创建RDD

Spark的textFile()方法可以读取本地文件系统或外部其它系统中的数据,并创建RDD。所不同的是,数据的来源路径不同。

首先在本地准备文件"/tmp/words.txt,内容如下:

hello hadoop
hello java
scala

进入Spark Shell读取本地数据(加载本地文件,必须以file:///开头):

//读取本地数据
scala> val rdd=sc.textFile("file:///tmp/words.txt")
rdd: org.apache.spark.rdd.RDD[String] = /tmp/words.txt MapPartitionsRDD[1]

scala> rdd.collect
res1: Array[String] = Array("hello hadoop ", "hello java ", "scala ")

异常问题

Spark Shell 调用 rdd.collect报错

spark Compression codec com.hadoop.compression.lzo.LzoCodec not found.

这是因为在hadoop中配置了编解码器lzo,所以当使用yarn模式时,spark自身没有lzo的jar包所以报错无法找到!
解决办法:
配置spark-default.conf文件!
修改spark-defaults.sh, 添加:

spark.jars /opt/pkg/hadoop/share/hadoop/common/hadoop-lzo-0.4.20.jar
读取本地数据时报错路径不存在

相似问题

这个问题发生在集群方式提交本地文件的时候。

这是因为文件的读取是发生在Executor所在的节点(即Worker节点),只有将文件tmp/words.txt同步到所有Worker节点才能避免路径不存在的问题发生。

有意思的地方是,即使不同步到所有Worker节点,先调用rdd.first拉取第一条数据之后再调用rdd.collect就可以获取到数据了(中间也会出现文件找不到问题,但是仍可以获取到结果)。原因就是RDD的collect方法属于转化算子(算子的概念后面会讲)是有惰性的(不是一开始马上收集,需要等待动作算子(如first或者count)操作数据时才会收集)才会主动收集数据

解决办法就是:

  1. 改为Local模式运行 $ spark-shell --master local
  2. 把文件file:///tmp/words.txt同步到所有salve节点。
  3. 把文件先上传到HDFS,然后再读取HDFS上的文件

读取HDFS数据

先将数据上传到HDFS

$ hadoop fs -put /tmp/words.txt /tmp/

再读取HDFS数据

//读取HDFS数据
scala> val rdd=sc.textFile("hdfs://centos01:8082/tmp/words.txt")
rdd: org.apache.spark.rdd.RDD[String] = hdfs://centos01:8082/tmp/words.txt MapPartitionsRDD[2]

scala> rdd.collect
res2: Array[String] = Array("hello hadoop ", "hello java ", "scala ")

02 RDD常用算子

RDD被创建后是只读的,不允许修改。Spark提供了丰富的用于操作RDD的方法,这些方法被称为算子。一个创建完成的RDD只支持两种算子:转化(Transformation)算子和行动(Action)算子。

1、转化算子

Spark中的转化算子主要时对RDD进行操作并返回新的RDD。

(1)map(func)

对rdd1应用map()算子,将rdd1中的每个元素加1并返回一个名为rdd2的新RDD:

scala> val rdd1=sc.parallelize(List(1,2,3,4,5,6))
scala> val rdd2=rdd1.map(x => x+1)

(2)filter(func)

对源RDD的每个元素应用函数func进行过滤,并返回一个新的RDD。
例如,过滤出rdd1中大于3的所有元素,并输出结果:

scala> val rdd1=sc.parallelize(List(1,2,3,4,5,6))
scala> val rdd2=rdd1.filter(_>3)
scala> rdd2.collect
res1: Array[Int] = Array(4, 5, 6)

上述代码中的下划线“_”代表rdd1中的每个元素。

(3)flatMap(func)

与map()算子类似,但是每个传入给函数func的RDD元素会返回0到多个元素,最终会将返回的所有元素合并到一个RDD。
例如,将集合List转为rdd1,然后调用rdd1的flatMap()算子将rdd1的每个元素按照空格进行分割成多个元素,最终合并所有元素到一个新的RDD。

scala> val rdd1=sc.parallelize(List("hadoop hello scala","spark hello"))
scala> val rdd2=rdd1.flatMap(_.split(" "))
scala> rdd2.collect
res3: Array[String] = Array(hadoop, hello, scala, spark, hello)

file

(4)reduceByKey()

reduceByKey()算子的作用对像是元素为(key,value)形式(Scala元组)的RDD,使用该算子可以将相同key的元素聚集到一起,最终把所有相同key的元素合并成为一个元素。该元素的key不变,value可以聚合成一个列表或者进行求和等操作。最终返回的RDD的元素类型和原有类型保持一致。
例如,有两个同学zhangsan和lisi,zhangsan的语文和数学成绩分别为98、78,lisi的语文和数学成绩分别为88、79,现需要分别求zhangsan和lisi的总成绩,代码如下:

scala> val list=List(("zhangsan",98),("zhangsan",78),("lisi",88),("lisi",79))
scala> val rdd1=sc.parallelize(list)
scala> val rdd2=rdd1.reduceByKey((x,y)=>x+y) //可以简化为rdd1.reduceByKey(_+_)
scala> rdd2.collect
res5: Array[(String, Int)] = Array((zhangsan,176), (lisi,167))

file

(5)groupByKey()

groupByKey()算子的作用对像是元素为(key,value)形式(Scala元组)的RDD,使用该算子可以将相同key的元素聚集到一起,最终把所有相同key的元素合并成为一个元素。该元素的key不变,value则聚集到一个集合中。
仍然以上述求学生zhangsan和lisi的总成绩为例,使用groupByKey()算子的代码如下:

scala> val list=List(("zhangsan",98),("zhangsan",78),("lisi",88),("lisi",79))
scala> val rdd1=sc.parallelize(list)
scala> val rdd2=rdd1.groupByKey()
scala> rdd2.map(x => (x._1,x._2.sum)).collect
res0: Array[(String, Int)] = Array((zhangsan,176), (lisi,167))

file

(6)union()

union()算子将两个RDD合并为一个新的RDD,主要用于对不同的数据来源进行合并,两个RDD中的数据类型要保持一致。
例如以下代码,通过集合创建了两个RDD,然后将两个RDD合并成了一个RDD:

scala> val rdd1=sc.parallelize(Array(1,2,3))
scala> val rdd2=sc.parallelize(Array(4,5,6))
scala> val rdd3=rdd1.union(rdd2)
scala> rdd3.collect
res8: Array[Int] = Array(1, 2, 3, 4, 5, 6) 

(7)sortBy()

sortBy()算子将RDD中的元素按照某个规则进行排序。该算子的第一个参数为排序函数,第二个参数是一个布尔值,指定升序(默认)或降序。若需要降序排列,需将第二个参数置为false。
例如,一个数组中存放了三个元组,将该数组转为RDD集合,然后对该RDD按照每个元素中的第二个值进行降序排列,代码如下:

scala> val rdd1=sc.parallelize(Array(("hadoop",12),("java",32),("spark",22)))
scala> val rdd2=rdd1.sortBy(x=>x._2,false)
scala> rdd2.collect
res2: Array[(String, Int)] = Array((java,32),(spark,22),(hadoop,12))

(8)sortByKey()

sortByKey()算子将(key,value)形式的RDD按照key进行排序。默认升序,若需降序排列,可以传入参数false,代码如下:
rdd.sortByKey(false)

(9)join ()

join()算子将两个(key,value)形式的RDD根据key进行连接操作,相当于数据库的内连接(inner join),只返回两个RDD都匹配的内容。例如,将rdd1和rdd2进行内连接,代码如下:

scala> val arr1= Array(("A","a1"),("B","b1"),("C","c1"),("D","d1"),("E","e1")
scala> val rdd1 = sc.parallelize(arr1))
scala> val arr2= Array(("A","A1"),("B","B1"),("C","C1"),("C","C2"),("C","C3"),("E","E1"))
scala> val rdd2 = sc.parallelize(arr2)
scala> rdd1.join(rdd2).collect
res0: Array[(String, (String, String))] = Array((B,(b1,B1)), (A,(a1,A1)), (C,(c1,C1)), (C,(c1,C2)), (C,(c1,C3)), (E,(e1,E1)))
scala> rdd2.join(rdd1).collect
res1: Array[(String, (String, String))] = Array((B,(B1,b1)), (A,(A1,a1)), (C,(C1,c1)), (C,(C2,c1)), (C,(C3,c1)), (E,(E1,e1)))

file

2、行动算子

Spark中的转化算子并不会马上进行运算,而是在遇到行动算子时才会执行相应的语句,触发Spark的任务调度。

file

2、行动算子

(1)reduce()

将数字1到100所组成的集合转为RDD,然后对该RDD进行reduce()算子计算,统计RDD中所有元素值的总和,代码如下:

scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.reduce(_+_)
res2: Int = 5050
scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.count
res3: Long = 100

(2)count()

统计RDD集合中元素的数量:

scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.reduce(_+_)
res2: Int = 5050
scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.count
res3: Long = 100

(3)countByKey()

例如,List集合中存储的是键值对形式的元组,使用该List集合创建一个RDD,然后对其进行countByKey()的计算:

scala> val rdd1 = sc.parallelize(List(("zhang",87),("zhang",79),("li",90)))
scala> rdd1.countByKey
res1: scala.collection.Map[String,Long] = Map(zhang -> 2, li -> 1) 

scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.take(5)
res4: Array[Int] = Array(1, 2, 3, 4, 5)

(4)take(n)

返回集合中前5个元素组成的数组:

scala> val rdd1 = sc.parallelize(List(("zhang",87),("zhang",79),("li",90)))
scala> rdd1.countByKey
res1: scala.collection.Map[String,Long] = Map(zhang -> 2, li -> 1) 

scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.take(5)
res4: Array[Int] = Array(1, 2, 3, 4, 5)

03 RDD的分区

RDD是一个大的数据集合,该集合被划分成多个子集合分布到了不同的节点上,而每一个子集合就称为分区(Partition)。因此也可以说,RDD是由若干个分区组成的。
RDD各个分区中的数据可以并行计算,因此分区的数量决定了并行计算的粒度。Spark会给每一个分区分配一个单独的Task任务对其进行计算,因此并行Task的数量是由分区的数量决定的。RDD分区的一个分区原则是使得分区的数量尽量等于集群中CPU核心数量。
file

04 RDD的依赖

Spark中,对RDD的每一次转化操作都会生成一个新的RDD,由于RDD的懒加载特性,新的RDD会依赖原有RDD,因此RDD之间存在类似流水线的前后依赖关系。这种依赖关系分为两种:窄依赖和宽依赖。

1、窄依赖

窄依赖是指父RDD的一个分区最多被子RDD的一个分区所用。也就是说,父RDD的分区与子RDD的分区的对应关系为一对一或多对一。例如map()、filter()、union()等操作都会产生窄依赖。
file

2、宽依赖

宽依赖是指,父RDD的一个分区被子RDD的多个分区所用。也就是说,父RDD的分区与子RDD的分区的对应关系为多对多。例如groupByKey()、reduceByKey()、sortByKey()等操作都会产生宽依赖。
file

Stage 划分

在Spark中,对每一个RDD的操作都会生成一个新的RDD,将这些RDD用带方向的直线连接起来(从父RDD连接到子RDD)会形成一个关于计算路径的有向无环图,称为DAG(Directed Acyclic Graph)。
file
Spark会根据DAG将整个计算划分为多个阶段,每个阶段称为一个Stage。每个Stage由多个Task任务并行进行计算,每个Task任务作用在一个分区上,一个Stage的总Task任务数量是由Stage中最后一个RDD的分区个数决定。

Stage的划分依据为是否有宽依赖,即是否有Shuffle。Spark调度器会从DAG图的末端向前进行递归划分,遇到Shuffle则进行划分,Shuffle之前的所有RDD组成一个Stage,整个DAG图为一个Stage。经典的单词计数执行流程的Stage划分:
file
比较复杂一点的Stage划分:
file

为说明要根据宽依赖划分Stage呢?这是因为窄依赖对优化很有利。逻辑上,每个RDD的算子都是一个fork/join操作(此join非上文的join算子,而是指同步多个并行任务的barrier): 把计算fork到每个分区,算完后join,然后fork/join下一个RDD的算子。如果子RDD的分区到 父RDD的分区是窄依赖,就可以实施优化把两个fork/join合为一个;如果连续的变换算子序列都是窄依赖,就可以把很多个 fork/join并为一个,这将极大地提升性能。Spark把这个叫做 流水线(pipeline)优化。

05 RDD的持久化

Spark中的RDD是懒加载的,只有当遇到行动算子时才会从头计算所有RDD,而且当同一个RDD被多次使用时,每次都需要重新计算一遍,这样会严重增加消耗。为了避免重复计算同一个RDD,可以将RDD进行持久化。
Spark中最重要的功能之一是可以将某个RDD中的数据保存到内存或者磁盘中,每次需要对这个RDD进行算子操作时,可以直接从内存或磁盘中取出该RDD的持久化数据,而不需要从头计算才能得到这个RDD。
例如有多个RDD,它们的依赖关系如图。若RDD3没有持久化保存,则每次对RDD3进行操作时都需要从textFile()开始计算,将文件数据转化为RDD1,再转化为RDD2,最终才得到RDD3。
可以在RDD上使用persist()或cache()方法来标记要持久化的RDD(cache()方法实际上底层调用的是persist()方法)。

file

RDD的检查点

RDD的检查点机制(Checkpoint)相当于对RDD数据进行快照,可以将经常使用的RDD快照到指定的文件系统中,最好是共享文件系统,例如HDFS。当机器发生故障导致内存或磁盘中的RDD数据丢失时可以快速从快照中对指定的RDD进行恢复,而不需要根据RDD的依赖关系从头进行计算,大大提高了计算效率。与cache()或者persist()将RDD数据存放到内存或者磁盘中的不同有以下几点:

(1)cache()或者persist()是将数据存储于机器本地的内存或磁盘,当机器故障时无法进行数据恢复,而检查点是将RDD数据存储于外部的共享文件系统(例如HDFS),共享文件系统的副本机制保证了数据的可靠性。

(2)在Spark应用程序执行结束后,cache()或者persist()存储的数据将被清空,而检查点存储的数据不会受影响,将永久存在,除非手动将其移除。因此,检查点数据可以被下一个Spark应用程序使用,而cache()或者persist()数据只能被当前Spark应用程序使用。

val sc = new SparkContext(conf);
//设置检查点数据存储路径
sc.setCheckpointDir("hdfs://centos01:8020/spark-ck")

06 共享变量

通常情况下,Spark应用程序运行的时候,Spark算子(例如map(func)或filter(func))中的函数func会被发送到远程的多个Worker节点上执行,如果一个算子中使用了某个外部变量,则该变量会拷贝到Worker节点的每一个Task任务中,各个Task任务对变量的操作相互独立。当变量所存储的数据量非常大时(例如一个大型集合)将增加网络传输及内存的开销。因此,Spark提供了两种共享变量:广播变量和累加器。

广播变量

广播变量是将一个变量通过广播的形式发送到每个Worker节点的缓存中,而不是发送到每个Task任务中,各个Task任务可以共享该变量的数据。因此广播变量是只读的。
1、默认情况下变量的传递
例如,map()算子中使用了外部变量arr:

val arr=Array(1,2,3,4,5);
val lines:RDD[String] = sc.textFile("file:///tmp/data.txt")
val result = lines.map(line =>
   (line, arr)
)

变量传递流程如图。
file
2、使用广播变量时变量的传递
例如,使用广播变量将数组arr传递给了map()算子:

val arr=Array(1,2,3,4,5);
val broadcastVar = sc.broadcast(arr)
val result = lines.map(line =>
   (line, broadcastVar) // broadcastVar为广播变量
)

传递流程如图。

file

在分布式函数中可以通过Broadcast对象的value方法访问
广播变量的值:

scala> val broadcastVar = sc.broadcast(Array(1, 2, 3))
scala> broadcastVar.value
res0: Array[Int] = Array(1, 2, 3)

累加器

累加器提供了将Worker节点的值聚合到Driver的功能,可以用于实现计数和求和。
例如,对一个整型数组进行求和,若不使用累加器,以下代码的输出结果不正确:

var sum=0 //在Driver中声明
val rdd=sc.makeRDD(Array(1,2,3,4,5))
rdd.foreach(x=>
  //在Executor中执行
  sum+=x
)
println(sum)//输出0

使用累加器对数组进行求和:

//声明一个累加器,默认初始值为0(只能在Driver端定义)
val myacc=sc.longAccumulator("My Accumulator")
val rdd=sc.makeRDD(Array(1,2,3,4,5))
rdd.foreach(x=>
   myacc.add(x)//向累加器中添加值
)
println(myacc.value)//输出15(只能在Driver端读取)

注意,累加器只能在Driver端定义,在Executor端更新。Executor端不能读取累加器的值,需要在Driver端使用value属性读取。

Views: 90

Index