2017年12月29日星期五

5.什么是RDD

 RDD是个抽象类,全称为Resilient Distributed Datasets,是一个容错的、并行的数据结构,可以让用户显式地将数据存储到磁盘和内存中,并能控制数据的分区。同时,RDD还提供了一组丰富的操作来操作这些数据诸如mapflatMapfilter等转换操作除此之外,RDD还提供了诸如joingroupByreduceByKey等更为方便的操作以支持常见的数据运算。但实际上继承RDD的派生类一般只要实现两个方法:
1getPartitions()用来告知怎么将input分片;
2compute()用来输出每个Partition被函数处理的一个单元);

RDD的特点:

1它是在集群节点上的不可变的、已分区的集合对象。
2通过并行转换的方式来创建如(map, filter, join, etc)。
3失败自动重建。
4可以控制存储级别(内存、磁盘等)来进行重用。
5必须是可序列化的。
6是静态类型的。

RDD的好处

1RDD只能从持久存储或通过Transformation操作产生,相比于分布式共享内存(DSM)可以更高效实现容错,对于丢失部分数据分区只需根据它的lineage就可重新计算出来,而不需要做特定的Checkpoint( RDD实现了基于Lineage的容错机制。RDD的转换关系,构成了compute chain,可以把这个compute chain认为是RDD之间演化的Lineage。在部分计算结果丢失时,只需要根据这个Lineage重算即可。)
2RDD的不变性,可以实现类似Hadoop MapReduce的推测式执行。
3RDD的数据分区特性,可以通过数据的本地性来提高性能,这与Hadoop MapReduce是一样的。
4RDD都是可序列化的,在内存不足时可自动降级为磁盘存储,把RDD存储于磁盘上,这时性能会有大的下降但不会差于现在的MapReduce

RDD的存储与分区

1用户可以选择不同的存储级别存储RDD以便重用。
2当前RDD默认是存储于内存,但当内存不足时,RDDspilldisk
3RDD在需要进行分区把数据分布于集群中时会根据每条记录Key进行分区(如Hash 分区),以此保证两个数据集在Join时能高效。

RDD的内部表示

RDD的内部实现中每个RDD都可以使用5个方面的特性来表示:
1分区列表(数据块列表)
2计算每个分片的函数(根据父RDD计算出此RDD
3对父RDD的依赖列表
4key-value RDDPartitioner(可选)
5每个数据分片的预定义地址列表(HDFS上的数据块的地址)(可选)

RDD创建方式:

1、从Hadoop文件系统(或与Hadoop兼容的其它存储系统)输入(例如HDFS)创建。
2、从父RDD转换得到新RDD
3、通过parallelize将单机数据创建为分布式RDD

4.Spark集群搭建

安装scala环境

下载地址http://www.scala-lang.org/download/
上传scala-2.10.5.tgzmasterslave机器的hadoop用户installer目录下
两台机器都要做

[hadoop@master installer]$ ls
hadoop2  hadoop-2.6.0.tar.gz  scala-2.10.5.tgz
解压
[hadoop@master installer]$ tar -zxvf scala-2.10.5.tgz
[hadoop@master installer]$ mv scala-2.10.5 scala
[hadoop@master installer]$ cd scala
[hadoop@master scala]$ pwd
/home/hadoop/installer/scala

配置环境变量:
[hadoop@master ~]$ vim .bashrc

# .bashrc

# Source global definitions
if [ -f /etc/bashrc ]; then
        . /etc/bashrc
fi

# User specific aliases and functions
export JAVA_HOME=/usr/java/jdk1.7.0_79
export HADOOP_HOME=/home/hadoop/installer/hadoop2
export SCALA_HOME=/home/hadoop/installer/scala
export HADOOP_COMMON_LIB_NATIVE_DIR=${HADOOP_HOME}/lib/native
export HADOOP_OPTS="-Djava.library.path=$HADOOP_HOME/lib"
export CLASSPATH=$CLASSPATH:$HADOOP_HOME/lib:$JAVA_HOME/lib:$SCALA_HOME/lib
export PATH=$PATH:$JAVA_HOME/bin:$HADOOP_HOME/bin:$HADOOP_HOME/sbin:$SCALA_HOME/bin
[hadoop@master ~]$ . .bashrc

安装python

安装gcc

[root@master ~]# mkdir /RHEL5U4
[root@master ~]# mount /dev/cdrom /media/
[root@master media]# cp -r * /RHEL5U4/
[root@master ~]vim /etc/yum.repos.d/iso.repo

[rhel-Server]
Name=5u4_Server
Baseurl=file:///RHEL5U4/Server
Enable=1
Gpgcheck=0
Gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-redhat-release

yum clean all
yum install gcc

Python安装

[root@master installer]# tar -zxvf Python-2.7.12
上传zlib-1.2.8.tar.gz
替换/root/installer/Python-2.7.12/Moduleszlib
[root@master Python-2.7.12]# ./configure --prefix=/usr/local/python27
[root@master Python-2.7.12]# make
[root@master Python-2.7.12]# make install
[root@master Python-2.7.12]# mv /usr/bin/python /usr/bin/python_old
[root@master Python-2.7.12]# ln -s /usr/local/python27/bin/python /usr/bin/
[root@master Python-2.7.12]# python
Python 2.7.12 (default, Nov  7 2016, 21:42:16)
[GCC 4.1.2 20080704 (Red Hat 4.1.2-46)] on linux2
Type "help", "copyright", "credits" or "license" for more information.
>>>


安装spark环境

下载地址http://spark.apache.org/downloads.html
上传spark-2.0.0-bin-hadoop2.6.tgzmasterhadoop用户installer目录下
解压缩
[hadoop@master installer]$ tar -zxvf spark-2.0.0-bin-hadoop2.6.tgz
[hadoop@master installer]$ mv spark-2.0.0-bin-hadoop2.6 spark2
[hadoop@master installer]$ cd spark2/
[hadoop@master spark2]$ ls
bin  conf  data  examples  jars  LICENSE  licenses  NOTICE  python  R  README.md  RELEASE  sbin  yarn
[hadoop@master spark2]$ pwd
/home/hadoop/installer/spark2

[hadoop@master ~]$ vim .bashrc

# .bashrc

# Source global definitions
if [ -f /etc/bashrc ]; then
        . /etc/bashrc
fi

# User specific aliases and functions
export JAVA_HOME=/usr/java/jdk1.7.0_79
export HADOOP_HOME=/home/hadoop/installer/hadoop2
export SCALA_HOME=/home/hadoop/installer/scala
export SPARK_HOME=/home/hadoop/installer/spark2
export HADOOP_COMMON_LIB_NATIVE_DIR=${HADOOP_HOME}/lib/native
export HADOOP_OPTS="-Djava.library.path=$HADOOP_HOME/lib"
export CLASSPATH=$CLASSPATH:$HADOOP_HOME/lib:$JAVA_HOME/lib:$SCALA_HOME/lib
export PATH=$PATH:$JAVA_HOME/bin:$HADOOP_HOME/bin:$HADOOP_HOME/sbin:$SCALA_HOME/bin:$SPARK_HOME/bin:$SPARK_HOME/sbin

[hadoop@master ~]$ . .bashrc

[hadoop@master ~]$ scp .bashrc slave:~
.bashrc                                                                                            100%  621     0.6KB/s   00:00
slave机器上执行
[hadoop@slave ~]$ . .bashrc

配置spark

[hadoop@master conf]$ cp spark-env.sh.template spark-env.sh
[hadoop@slave conf]$ vim spark-env.sh

#!/usr/bin/env bash

#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements.  See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License.  You may obtain a copy of the License at
#
#    http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
#
export JAVA_HOME=/usr/java/jdk1.7.0_79
export SCALA_HOME=/home/hadoop/installer/scala
export SPARK_MASTER_HOST=master
export HADOOP_CONF_DIR=$HADOOP_HOME/etc/hadoop
export SPARK_EXECUTOR_MEMORY=600M
export SPARK_DRIVER_MEMORY=600M

[hadoop@slave conf]$ vim slaves
master
slave

[hadoop@master installer]$ scp -r spark2 slave:~/installer/
启动spark集群

[hadoop@master ~]$ start-master.sh
[hadoop@master ~]$ start-slaves.sh

[hadoop@master ~]$ jps
17769 ResourceManager
20192 Master
20275 Worker
17443 NameNode
20521 Jps
17631 SecondaryNameNode

[hadoop@slave ~]$ jps
13297 DataNode
15367 Worker
13408 NodeManager
16245 Jps

Spark wordcount


[hadoop@master ~]$ spark-shell
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel).
16/11/04 11:05:07 WARN util.NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
16/11/04 11:05:09 WARN spark.SparkContext: Use an existing SparkContext, some configuration may not take effect.
Spark context Web UI available at http://192.168.3.100:4040
Spark context available as 'sc' (master = local[*], app id = local-1478228709028).
Spark session available as 'spark'.
Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /___/ .__/\_,_/_/ /_/\_\   version 2.0.0
      /_/
         
Using Scala version 2.11.8 (Java HotSpot(TM) Client VM, Java 1.7.0_79)
Type in expressions to have them evaluated.
Type :help for more information.

scala> val file = sc.textFile("hdfs://master:9000/data/wordcount")
16/11/04 11:05:14 WARN util.SizeEstimator: Failed to check whether UseCompressedOops is set; assuming yes
file: org.apache.spark.rdd.RDD[String] = hdfs://master:9000/data/input/wordcount MapPartitionsRDD[1] at textFile at <console>:24

scala> val count=file.flatMap(line => line.split(" ")).map(word => (word,1)).reduceByKey(_+_)
count: org.apache.spark.rdd.RDD[(String, Int)] = ShuffledRDD[4] at reduceByKey at <console>:26
scala> count.collect()
res0: Array[(String, Int)] = Array((package,1), (this,1), (Version"](http://spark.apache.org/docs/latest/building-spark.html#specifying-the-hadoop-version),1), (Because,1), (Python,2), (cluster.,1), (its,1), ([run,1), (general,2), (have,1), (pre-built,1), (YARN,,1), (locally,2), (changed,1), (locally.,1), (sc.parallelize(1,1), (only,1), (Configuration,1), (This,2), (basic,1), (first,1), (learning,,1), ([Eclipse](https://cwiki.apache.org/confluence/display/SPARK/Useful+Developer+Tools#UsefulDeveloperTools-Eclipse),1), (documentation,3), (graph,1), (Hive,2), (several,1), (["Specifying,1), ("yarn",1), (page](http://spark.apache.org/documentation.html),1), ([params]`.,1), ([project,2), (prefer,1), (SparkPi,2), (<http://spark.apache.org/>,1), (engine,1), (version,1), (file,1), (documentation...
scala>

2.Spark-Linux环境准备

硬软件环境

1、虚拟机:VMware Workstation 12
2、虚拟机操作系统:RedHat5u4,单核,1G内存,2两台
3、虚拟机运行环境:
java version "1.7.0_79" 64
Scala version 2.10.5
hadoop-2.6.0
spark-2.0.0-bin-hadoop2.6
Python 2.7.12

配置IP     

192.168.3.100 master
192.168.3.101 slave

配置主机名

[root@localhost ~]# vim /etc/sysconfig/network
NETWORKING=yes
NETWORKING_IPV6=no
HOSTNAME=master
[root@localhost ~]# hostname master
三台主机都配置主机名
[root@localhost ~]# vim /etc/sysconfig/network
NETWORKING=yes
NETWORKING_IPV6=no
HOSTNAME= slave
[root@localhost ~]# hostname slave
重启生效

关闭防火墙和selinux

[root@master ~]# vim /etc/sysconfig/selinux
# This file controls the state of SELinux on the system.
# SELINUX= can take one of these three values:
#       enforcing - SELinux security policy is enforced.
#       permissive - SELinux prints warnings instead of enforcing.
#       disabled - SELinux is fully disabled.
SELINUX=disabled
# SELINUXTYPE= type of policy in use. Possible values are:
#       targeted - Only targeted network daemons are protected.
#       strict - Full SELinux protection.
SELINUXTYPE=targeted

[root@master ~]# service iptables status
表格:filter
Chain INPUT (policy ACCEPT)
num  target     prot opt source               destination         

Chain FORWARD (policy ACCEPT)
num  target     prot opt source               destination         

Chain OUTPUT (policy ACCEPT)
num  target     prot opt source               destination    

域名解析

使得masterslave互相连接对方
[root@master ~]# vim /etc/hosts
# Do not remove the following line, or various programs
# that require network functionality will fail.
127.0.0.1               localhost.localdomain localhost
::1             localhost6.localdomain6 localhost6

192.168.3.100 master
192.168.3.101 slave

[root@slave ~]# vim /etc/hosts
# Do not remove the following line, or various programs
# that require network functionality will fail.
127.0.0.1               localhost.localdomain localhost
::1             localhost6.localdomain6 localhost6

192.168.3.100 master
192.168.3.101 slave

1.Spark简介

简介:

  Spark是加州大学伯克利分校AMP实验室,开发的通用内存并行计算框架。Spark在2013年6月进入Apache成为孵化项目,8个月后成为Apache顶级项目Spark以其先进的设计理念,迅速成为社区的热门项目,围绕着Spark推出了Spark SQL、Spark Streaming、MLLib和GraphX等组件,也就是BDAS(伯克利数据分析栈),这些组件逐渐形成大数据处理一站式解决平台。
  Spark使用Scala语言实现,它是一种面向对象、函数式编程语言,能够像操作本地集合对象一样轻松的操作分布式数据集。

Spark特点:

1、运行速度快

  Spark拥有DAG执行引擎,支持在内存中对数据进行迭代计算。官方提供的数据表明,如果数据由磁盘读取,速度是Hadoop MapReduce的10倍以上,如果数据从内存中读取,速度可以高达100多倍。


2、易用性好

  Spark不仅支持Scala编写应用程序,而且支持Java和Python等语言进行编写,特别是Scala是一种高效、可拓展的语言,能够用简洁的代码处理较为复杂的处理工作。

3、通用性强

  Spark生态圈即BDAS(伯克利数据分析栈)包含了Spark Core、Spark SQL、SparkStreaming、MLLib和GraphX等组件,这些组件分别处理:Spark Core提供内存计算框架、SparkStreaming的实时处理应用、Spark SQL的即时查询、MLlib的机器学习和GraphX的图处理,它们都是由AMP实验室提供,能够无缝的集成并提供一站式解决平台。

4、随处运行
  Spark具有很强的适应性,能够读取HDFS、Cassandra、HBase、S3为持久层读写原生数据,能够以Mesos、YARN和自身携带的Standalone作为资源管理器调度job,来完成Spark应用程序的计算。

Spark与Hadoop的对比(Spark的优势


1、Spark的中间数据放到内存中,对于迭代运算效率更高
2、Spark比Hadoop更通用
3、Spark提供了统一的编程接口
4、容错性– 在分布式数据集计算时通过checkpoint来实现容错
5、可用性– Spark通过提供丰富的Scala, Java,Python API及交互式Shell来提高可用性

1、目前大数据处理场景
1. 复杂的批量处理(Batch Data Processing),偏重点在于处理海量数据的能力,至于处理速度可忍受,通常的时间可能是在数十分钟到数小时;
2. 基于历史数据的交互式查询(Interactive Query),通常的时间在数十秒到数十分钟之间
3. 基于实时数据流的数据处理(Streaming Data Processing),通常在数百毫秒到数秒之间

2、以上三种场景都有比较成熟的处理框架

1————————Hadoop的MapReduce来进行海量数据的批处理
2————————Impala进行交互式查询
3————————Storm分布式处理框架处理实时流式数据

3、总结:

1、Spark是基于内存的迭代计算框架,适用于需要多次操作特定数据集的应用场合。需要反复操作的次数越多,所需读取的数据量越大,受益越大,数据量小但是计算密集度较大的场合,受益就相对较小。
2、由于RDD的特性,Spark不适用那种异步细粒度更新状态的应用,例如web服务的存储或者是增量的web爬虫和索引。就是对于那种增量修改的应用模型不适合。
3、数据量不是特别大,但是要求实时统计分析需求。

Spark成功案例:

1、腾讯

  广点通是最早使用Spark的应用之一。腾讯大数据精准推荐借助Spark快速迭代的优势,围绕“数据+算法+系统”这套技术方案,实现了在“数据实时采集、算法实时训练、系统实时预测”的全流程实时并行高维算法,最终成功应用于广点通投放系统上,支持每天上百亿的请求量。

2、Yahoo

  Yahoo将Spark用在Audience Expansion中的应用。Audience Expansion是广告中寻找目标用户的一种方法:首先广告者提供一些观看了广告并且购买产品的样本客户,据此进行
学习,寻找更多可能转化的用户,对他们定向广告。Yahoo采用的算法是logistic regression。目前在Yahoo部署的Spark集群有112台节点,9.2TB内存。

3、淘宝

  阿里搜索和广告业务,最初使用Mahout或者自己写的MR来解决复杂的机器学习,导致效率低而且代码不易维护。淘宝技术团队使用了Spark来解决多次迭代的机器学习算法、高计算复杂度的算法等。将Spark运用于淘宝的推荐相关算法上,同时还利用Graphx解决了许多生产问题,包括:基于度分布的中枢节点发现、基于最大连通图的社区发现、基于三角形计数的关系衡量、基于随机游走的用户属性传播等。

Spark运行模式

Local
本地模式
用于本地开发测试,本地还分为local单线程和local-cluster多线程。
On yarn
集群模式
运行在yarn框架之上,由yarn负责资源管理,Spark负责任务调度和计算。
Standalone
集群模式
典型的Mater-slave模式,Spark自带的模式。

Concepts:

Application:

Spark Application的概念和Hadoop MapReduce类似,指的是用户编写的Spark应用程序,包含了一个Driver 功能的代码和分布在集群中多个节点上运行的Executor代码。

Driver:

Spark中的Driver即运行上述Application的main()函数并且创建SparkContext,其中创建SparkContext的目的是为了准备Spark应用程序的运行环境。在Spark中由SparkContext负责和ClusterManager通信,进行资源的申请、任务的分配和监控等;当Executor部分运行完毕后,Driver负责将SparkContext关闭。通常用SparkContext代表Drive。

Executor:

Application运行在Worker 节点上的一个进程,该进程负责运行Task,并且负责将数据存在内存或者磁盘上,每个Application都有各自独立的一批Executor。在Spark on Yarn模式下,其进程名称为CoarseGrainedExecutorBackend,类似于Hadoop MapReduce中的YarnChild。一个CoarseGrainedExecutorBackend进程有且仅有一个executor对象,它负责将Task包装成taskRunner,并从线程池中抽取出一个空闲线程运行Task。每个CoarseGrainedExecutorBackend能并行运行Task的数量就取决于分配给它的CPU的个数了。

Cluster Manager:

指的是在集群上获取资源的外部服务,目前有:Standalone、Mesos、Yarn

Worker

集群中任何可以运行Application代码的节点,类似于YARN中的NodeManager节点。在Standalone模式中指的就是通过Slave文件配置的Worker节点,在Spark on Yarn模式中指的就是NodeManager节点。 

作业(Job):

包含多个Task组成的并行计算,往往由Spark Action催生,一个JOB包含多个RDD及作用于相应RDD上的各种Operation。

阶段(Stage):

每个Job会被拆分很多组Task,每组任务被称为Stage,也可称TaskSet,一个作业分为多个阶段。

任务(Task):

被送到某个Executor上的工作任务。

RDD(*):

是Resilient distributed datasets的简称,中文为弹性分布式数据集;是Spark最核心的模块和类。

DAGScheduler:

根据Job构建基于Stage的DAG,并提交Stage给TaskScheduler。

TaskScheduler:

将Taskset提交给Worker node集群运行并返回结果。

Transformations:

是Spark API的一种类型,Transformation返回值还是一个RDD, 所有的Transformation采用的都是懒策略,如果只是将Transformation提交是不会执行计算的。

Action:

是Spark API的一种类型,Action返回值不是一个RDD,而是一个scala集合;计算只有在Action被提交的时候计算才被触发。

Jurassic World 3" opens in theaters this Friday, 27 dinosaurs set to come, 10 first appearance

 The annual mega-production "Jurassic World 3" will be officially released in China on June 10, and simultaneously landed in IMAX ...