分类: Spark

  • Spark中的数据框与数理统计

    Spark中的数据框与数理统计

    DataFrame是R中一个基本结构,俗称数据框。其内部可以由多种数据类型,每一列是一个变量,每行是一个观测记录。在R中数据框是很通用的数据结构,它是一种特殊的列表对象。

    在R中数据框对象包含了大量的基本数理统计运算,从简单的最大最小值和均值,再到协方差等等。

    可以说数据框是了解一系列数据基本特性的快速方法,数据框的存在让R更加易用。

    在Spark 1.3中新增了DataFrame,隶属于Spark-Sql模块下。

    在即将到来的1.4版本中DataFrame的功能进一步扩充,而且在后续版本中会进一步提升。

    具体可以参考

    而1.4版本中主要增加的功能有:

    • 随机数据生成
    • 描述性统计
    • 样本协方差和相关性
    • 列联表
    • 频繁项

    随机数据生成

    随机数据生成对于测试现有算法很有用,也可以用在随机算法的开发中。

    SQLContext sqlContext =new SQLContext(sc);
    DataFrame dataFrame = sqlContext.range(1,10);
    DataFrame frame = dataFrame.select(functions.col("id"), functions.rand(10).alias("uniform"), functions.randn(25).alias("normal"));
    frame.describe("uniform","normal").show();

    随机数据

    描述性统计

    获得一些数据后需要做的第一步就是对数据整体有一个感知。描述性统计是指运用制表和分类,图形以及计算概括性数据来描述数据特征的各项活动。

    描述性统计可以帮助我们了解到数据的大体分布和情况。

    Spark中的描述统计提供了计数,均值,标准差,最大最小值等信息。

    SQLContext sqlContext =new SQLContext(sc);
    
  • Spark中的梯度下降

    Spark中的梯度下降

    梯度下降是一个最优化算法,通常也称为最速下降。
    梯度下降是求解无约束优化问题最简单和最古老的方法之一,许多有效算法都是以它为基础进行改进和修正而得到的。

    最速下降是用负梯度方向为搜索方向的,最速下降法越接近目标值,步长越小,前进越慢。

    虽然古老而简单,但是很多时候依然可以将它用在一些情况中,也很适合快速入门。

    Spark中的优化组件

    Spark中的MLlib下的optimization包含了梯度下降的相关组件。

    从数理基础和很多改进算法来看,梯度下降的主要抽象集中在两个方面,即如何度量梯度和如何更新参数。

    Spark中的梯度抽象为Gradient.scala

    abstractclass Gradient extends Serializable {
    def compute(data: Vector, label: Double, weights: Vector): (Vector, Double) = {
    val gradient = Vectors.zeros(weights.size)
    val loss = compute(data, label, weights, gradient)
    (gradient, loss)
    }
    
    def compute(data: 
  • Spark中的ML Pipelines

    Spark中的ML Pipelines

    Spark生态圈中有一个MLlib,其目标在于使机器学习更加简单和可扩展。
    MLlib的开发非常活跃,其中添加新的算法和提升性能是一个主要的方向,而另一个方面是让MLlib更简单。

    和Spark Core类似,MLlib提供了三种语言用API:Scala,Java和Python。用户可以根据文档和例子开始使用相应的功能,但是由于用户的技术和知识背景不同,学习曲线各不相同。为了让MLlib更好用,Spark 1.2开始引入了ML Pipeline API。

    一个常见的机器学习流程包含了序列数据的预处理,特征提取,模型拟合和验证等。

    以文本分类为例,包含了文本切割,文本清理,特征提取,训练分类模型和交叉检验。
    虽然对于其中的每一个步,都有很多可以选用的第三方库来完成,但是将它们结合起来一起使用却不是那么简单。

    很多库的API外观迥异,而且并没有为分布式数据处理提供支持。

    Spark的ML Pipelines要做的就是从抽象层开始让机器学习更简单(不论是使用Spark以后的算法还是自己实现)。

    数据抽象

    在整体设计中,数据集的表现通过DataFrame来完成。DataFrame是Spark SQL中的组件。

    除了单纯的数据源以外,还有一个或者多个转化器将数据转化后置入Spark甬道中。

    转化器获取到输入数据并给出输出数据,输出数据又会成为下一个步骤的输入数据。
    直接选用Spark SQL中的组件主要是考虑到了数据的输入输出,灵活的数据列类型操作和优化问题。

    数据的输入输出是一个甬道的开始或者结束。
    目前提供的输入输出包含了

    • LabeledPoint (分类和回归)
    • Rating (协同过滤)
    • 等等

    特征转化是一个甬道中的重要部分,比如将文本转为词序列,对词做TF-IDF变化等。

    甬道示例

    甬道

    ML Pipelines相关内容都在spark.ml包下。

    一个甬道包含了多个步骤,最基本的包含转化和评估。

    转化就是上文提到的将文本转为词序列等等,评估更多的是训练对应的模型。

    创建甬道:

    Tokenizer tokenizer =new Tokenizer()
    .setInputCol("text")
    .setOutputCol("words");
    
  • 从Hadoop MapReduce到Spark

    Spark是一款通用集群计算框架,和Hadoop的MapReduce类似。由于其提供的抽象更简单,性能和功能上比Hadoop强不少。
    它已经越来越流行了。

    如果是新的项目,或者是为了学习,那么选择Spark完全没问题。

    不过对于一些使用了MapReduce的项目来说,迁移就稍微复杂一些了。

    Hadoop自身

    Hadoop的使用在不断扩大,但同时越来越多的实践证实了MapReduce并不是通用计算范例。

    Hadoop的架构本身为其他可能的替代方案提供了场所,比如Impala项目等等。

    而对于Hadoop来说,有一部分Hadoop的实现本身和MapReduce本身的抽象并不一致。

    • Mapper和Reducer总是使用键值对作为输入输出
    • Reducer处理的级别是键
    • Mapper和Reducer的对象的生命周期跨越了多个map()和reduce(),同时还支持了setup()和cleanup()

    Spark

    Spark也是众多替代方案中的一种。

    但是对于已经部署在生产环境的项目而言,一句替代方案是不够的,对于一些实时计算的系统更是这样。

    好在利用Spark实现类似MapReduce的模型是完全可行的。同时实现本身还可以更简单,并在大部分情况下更快。

    对于MapReduce模型本身,使用Spark来实现反而显得更亲近,毕竟Scala的编码风格和API对于本身源于LISP的MapReduce抽象更接近。

    键值对和元组

    从最基础的例子来看,如果需要计算一个大型文本文件每一行的长度,对于Hadoop MapReduce来说由于输出是键值对,那么就会使用长度作为键,以1为值。

    publicclass LineLengthMapper extends Mapper<LongWritable,Text,IntWritable,IntWritable> {
    protectedvoidmap(LongWritable lineNumber, Text line, Context context)
    throws IOException, InterruptedException {
    context.write(new IntWritable(line.getLength()),new IntWritable(1));
    }
    }

    LineLengthMapper