分类: Spark

  • 使用Docker快速启动一个Spark集群

    docker-compose 文件如下

    version: '3'
    services:
      master:
        image: bde2020/spark-master:2.3.1-hadoop2.7
        ports:
          - "8080:8080"
          - "7077:7077"
          - "6066:6066"
        environment:
          - INIT_DAEMON_STEP=setup_spark
      worker:
        image: bde2020/spark-worker:2.3.1-hadoop2.7
        depends_on:
          - master
        ports:
          - "8081"
        environment:
          - "SPARK_MASTER=spark://master:7077"
    

     

    然后运行docker-compose up -d –scale worker=3 就可以启动一个有三个节点的集群了.…

  • Spark上可用的自然语言处理框架JSL NLP

    Spark上可用的自然语言处理框架JSL NLP

    Spark ML Pipeline是一个非常方便的结构,只需要提供其中相应的部件就可以做出很多可以重用的Pipeline。官方自带的ML包中包含了很多常用部件,但是唯独缺少对于自然语言处理的支持。今天群友介绍了一款专门处理自然语言的Spark支持库John Snow Labs NLP。

    John Snow Labs NLP和Spark一样,遵循Apache协议。而且一开始定位就是基于Spark的,没有其他第三方依赖。所有组件都是基于Spark ML Pipeline API的,使用上也没有问题。

    主要的内容包括:

    1. Tokenizer
    2. Normalizer
    3. Stemmer
    4. Lemmatizer
    5. Entity Extractor
    6. Date Extractor
    7. Part of Speech Tagger
    8. Named Entity Recognition
    9. Sentence boundary detection
    10. Sentiment analysis
    11. Spell checker

    这些组件有部分功能和Spark自带的有重复,比如Stop Word Remover,但是多一份选择不是坏处。

     …

  • 用Spark MLlib 2.X来驱动你的机器学习工作流

    Spark的机器学习模块在2.x版本正式移动到ml包下,也就是说旧有的包只做维护不在添加新的功能。新的ml包中最大的改变就是使用用统一的工作流模型来囊括所有机器学习相关的东西。主要是以下三个大的目标:

    1. 线性可扩展
    2. 容错
    3. 内建所有的常用算法

    我们的最终产物当然还是模型本身,也就是数学函数。工作流本身只是概念上的东西,并没有影响到本质。如果你喜欢,你也可以把工作流的所有东西拆开使用。

    我们经常说学术界和工业界有很大的区别,比如你看到的ML Pipeline(文档或者文献)大部分长这个样子

    然而真实世界更为复杂,比如长这样

    这之中的主要区别在于数据来源,你在文档中看到的更多是完美的数据输入,干净的数据,一次性的运行,而工业界的数据来源完全不能保证。

    所以更多的时候我们把数据工作分为数据科学和数据工程。数据科学使用R或者Python去构建原型系统,而数据工程使用Java重写对应的实现。

    Spark MLlib 2.x的一个重大改动就是对于模型序列化的优化,也就是说Python或者R输出的模型可以直接被Java载入。

    ML Pipeline的另外一个好处就是隐藏了ML本身的代码,只要内建了足够的常用算法,或者模型能够不同语言通用,那么关于机器学习的核心代码其实是可以被Pipeline组合并隐藏掉具体的实现。

    而工业界还有很多活需要干,比如配置,数据来源,可视化,应用健康监控等等。机器学习确实很酷,但是关注点是不同。

    如果在两边在工作室都遵循工作流模型,并且尽可能重用内建的算法,那么渐渐的工程中就可以集中在数据抽取和最终的输出上。同时最开始输出的原型系统也能够和最终上线的产品系统有一个更高的一致性。…

  • Apache Spark

    Apache Spark

    spark-logo-trademark

    Apache Spark是一个围绕速度、易用性和复杂分析构建的大数据处理框架。最初在2009年由加州大学伯克利分校的AMPLab开发,并于2010年成为Apache的开源项目之一。

    Spark推出时的一个特点是快,对比的对象自然是Hadoop。

    Hadoop这项大数据处理技术大概已有十年历史,而且被看做是首选的大数据集合处理的解决方案。MapReduce是一路计算的优秀解决方案,不过对于需要多路计算和算法的用例来说,并非十分高效。数据处理流程中的每一步都需要一个Map阶段和一个Reduce阶段,而且如果要利用这一解决方案,需要将所有用例都转换成MapReduce模式。

    而Spark则允许程序开发者使用有向无环图开发复杂的多步数据管道。而且还支持跨有向无环图的内存数据共享,以便不同的作业可以共同处理同一个数据。

    Spark的目的并不是代替Hadoop,相反Spark是可以运行于Hadoop之上的,包括了使用Hadoop文件系统,Yarn调度器等等。

    Spark之上还有四个模块,分别对应了SQL,Streaming,ML还有Graph。spark-stack

    Spark近期的发力方向主要在于Machine Learning的Pipeline模式之上。

    这种设计理念将机器学习统一抽象为由DataFrame,Transformer,Estimator和Parameter组成的Pipeline,并提供了大量可以重用的组件。ml-pipeline

    Spark提供了多种语言的API,你可以使用Scala,Java,Python或者R来开发你自己的应用。…

  • 停止词和StopWordsRemover

    停止词简单来说是指在一种语言中广泛使用的词。在各种需要处理文本的地方,我们对这些停止词做出一些特殊处理,以方便我们更关注在更重要的一些词上。

    对于不同类型的需求而言,对停止词的处理是不同的。

    1. 有监督的机器学习 – 将停止词从特征空间剔除
    2. 聚类– 降低停止词的权重
    3. 信息检索– 不对停止词做索引
    4. 自动摘要- 计分时不处理停止词

    对于不同语言,停止词的类型都可能有出入,但是一般而言有这简单的三类

    1. 限定词
    2. 并列连词
    3. 介词

    停止词的词表一般不需要自己制作,有很多可选项可以自己下载选用。

    Spark中提供了StopWordsRemover类处理停止词,它可以用作Machine learning Pipeline的一部分。

    StopWordsRemover的功能是直接移除,所有从inputCol输入的量都会被它检查,然后再outputCol中,这些停止词都会去掉了。

    默认的话会加载/org/apache/spark/ml/feature/stopwords/english.txt

    这是一个简单的停止词表,包含153个词。

    默认还提供了其他几种语言的停止词,遗憾的是没有中文默认停止词表,所以对于中文停止词需要自己提供。…

  • Spark快速获得CrossValidator的最佳模型参数

    Spark快速获得CrossValidator的最佳模型参数

    Spark提供了便利的Pipeline模型,可以轻松的创建自己的学习模型。

    但是大部分模型都是需要提供参数的,如果不提供就是默认参数,那么怎么选择参数就是一个比较常见的问题。Spark提供在org.apache.spark.ml.tuning包下提供了模型选择器,可以替换参数然后比较模型输出。

    目前有CrossValidator和TrainValidationSplit两种,比如一个文本情感预测模型。

    Pipeline只有三步,第一步切词,第二部Hashing TF,第三部NB分类

    Pipeline pipeline = new Pipeline()
                    .setStages(new PipelineStage[]{tokenizer, hashingTF, naiveBayes});
    
    ParamMap[] paramMaps = new ParamGridBuilder()
                    .addGrid(hashingTF.numFeatures(), new int[]{10000, 100000, 500000, 1000000})
                    .build();
    CrossValidator cv = new CrossValidator()
                    .setEstimator(pipeline)
                    .setEvaluator(new BinaryClassificationEvaluator())
                    .setEstimatorParamMaps(paramMaps);

    其中Hashing TF的参数选择非常重要,我们这里就随便尝试几种,然后放在CrossValidator中去。

    最后我们会获得一个CrossValidatorModel类,这里有两种选择。

    第一种是自己手动获取其中的参数,因为bestModel的参数就是我们最后选择的参数

    Pipeline 
  • 利用Docker快速搭建Spark本地环境

    Spark模式是直接local直接开发的,也就是在SparkConf中直接设定为local[*]之类,就可以在本地启动Spark然后开始工作。

    但是有时候还是希望将这些分开,也就是说有一个独立的Spark Master和一些Workers,本地开发,但是运行还是在简单的集群上。

    因为需求很简单,所以直接自己写一个Dockerfile

    FROM java:8
    WORKDIR /opt
    RUN wget -q http://apache.fayea.com/spark/spark-1.6.2/spark-1.6.2-bin-hadoop2.6.tgz -O spark.tgz
    RUN tar xfz spark.tgz && mv spark-1.6.2-bin-hadoop2.6 spark
    EXPOSE 8080 7077 6066
    ENTRYPOINT ./spark/sbin/start-master.sh && ./spark/sbin/start-slave.sh spark://$(ip addr show eth0 | grep "inet\b" | awk '{print 
  • Java项目中混合Scala

    虽然我并不怎么用Scala,但是经常接触到一些Scala的开源库。由于Scala本身的特性,所以对于使用者而言,懂不懂Scala并不重要。

    Spark是由Scala编写,可以只用Java调用,但是有时候需要自定义其中的一些组件的时候可能Java并不能做到,这个时候就需要写一些Scala的代码。

    原有的项目是Gradle管理的,而Gradle本身提供了对于Scala的支持,简单来看看Gradle Scala插件。在build.gradle中添加两行

    apply plugin: 'scala'
    ...
    compile 'com.databricks:spark-csv_2.11:1.4.0'

    然后在目录结构中添加一个scala,位置在这里scala-project

    然后直接开始写就行了,不得不说IDEA和Gradle工作的都很好,运行还是直接运行,其他什么都不用改,直接运行就行了。

    scala-spark

    参考资料

    https://docs.gradle.org/current/userguide/scala_plugin.html…

  • 快速将csv转为Spark Dataframe

    CSV,或者叫逗号分隔值,是以逗号为分隔符,简单而使用。虽然并没有真正的标准,但是RFC 4180中有一个大致的表述。

    很多时候我们拿到的原始数据都是csv的,而快速将其转为Spark的Dataframe做进一步分析就是一个经常遇到的问题。

    先来一个简单的例子,这里以手淘的数据为例子

    spark-dataframe-1

    共有六列。

    先建立一个简单对象Record,然后直接用Spark的createDataFrame方法

    JavaRDD<Record> list = sc.textFile(userFile).map(new Function<String, Record>() {
    	@Override
    	public Record call(String v1) throws Exception {
    		String[] split = v1.split(",");
    		return new Record(split[0], split[1], split[2], split[3], split[4], split[5]);
    	}
    });
    DataFrame dataFrame = sqlContext.createDataFrame(list, Record.class);
  • Spark的Datasets

    对于Spark的使用者来说,越简单易用的API越好。所以在原有的RDD之上,Spark陆续添加了DataFrames和Spark SQL。特别是DataFrames,对于小数据集的快速上手非常简单,其API也是从R中借鉴的,很有亲切感。当然在这之下,底层的实现还是RDD,只是对于使用者,抽象程度更高而已。

    DataSets是为了解决类型安全和面向对象思维而添加的新API,包含在Spark 1.6中,其也是下个版本集中的方向。

    DataSets并不是从关系型数据结构抽象出来的,所以它的API反而和RDD类似,比如

    val lines = sqlContext.read.text("/wikipedia").as[String]
    val words = lines
      .flatMap(_.split(" "))
      .filter(_ != "")
    val counts = words 
        .groupBy(_.toLowerCase)
        .count()

    在直接使用的效率和资源消耗上要好一些(图来自DataBricks)

    Memory-Usage-when-Caching-Chart-1024x359Distributed-Wordcount-Chart-1024x371

    当然,另外一个增强面向对象的特性使用很简单, 主要基于Encoders

    public class Person implements Serializable {
        private String name;
    
        get/set function{}
    }
    
    Dataset<Person>