博客

  • 云应用与Lattice

    快,是现在一个很喜欢提的概念。快速发现用户需求,快速的产品原因,快速的提供服务,快速的迭代,快速的部署,快速的试错。

    为了满足这些,软件架构和基础服务也在发生改变。

    从独立机房到云主机,再到所谓的容器,快速的部署一个具有极大扩展性的应用变得越来越容易,当然费用也是越来越低。

    为了适应这些变化,软件架构也在向云应用转变。

    Cloud Native Application

    Cloud Native Application,说实话,我不知道中文名字叫什么,先称为云应用(我相信应该有一个更正式的说法)。

    从中衍生的还有所谓的 Cloud Native Application Architectures,即云应用体系结构。

    简单说来无非是一下几个问题:

    • 以一个或多个无状态进程运行应用
    • 通过端口绑定提供服务
    • 快速启动和优雅终止可最大化健壮性
    • 把日志当作事件流
    • 其他。。。

    应用由小的模块组成,每个模块尽可能无状态,模块能够快速的终止和启动。

    更多的相关概念和理论可以参考文木末的免费电子书。

    Spring-Boot和Spring-Cloud

    以Web应用说明,如果构建成war包然后发布到容器中,那么就很难做到以端口绑定提供服务,多个无状态进程等等。

    Spring-Boot提供了一个很简单的思路,如果发布到容器不行,那么就把容器包含在自身,使用内置的jetty或者tomcat来作为服务的承载。

    当然Spring-Boot还提供了快速的开发的若干便利,从配置到依赖管理,进而到调试和发布。

    Spring-Cloud提供了功能更进一步的“云特性”。

    • 配置管理
    • 路由
    • 负载均衡
    • 服务注册和发现
    • 集群决策和状态管理

    可以说在Spring的大旗之下有超过10个项目或者子项目再为云应用这个理念服务。

    Lattice

    有了理念,有了工具和框架,自然需要一个平台。

    Lattice就是一个为云应用而生的管理平台,主要的功能包括:

    • Http负载均衡
    • 日志
    • 监控
    • 集群管理和部署

    当然也有一个spring-cloud-lattice

    Lattice目前支持 AWS, Digital Ocean, Google Cloud和Openstack。国内使用较多的阿里云还没有支持。

    不过这个支持指的是部署Lattice的功能,你依然可以手动安装在任何主机上使用。

    Lattice的部署标准是Docker,意味着你的服务需要适配在Docker中。

    对于Spring-Boot的应用而言其实这是非常简单的,即使不使用spring-docker插件,我们也可以自己写Dockerfile。

    FROM dockerfile/java:oracle-java7
    MAINTAINER [email protected]
    EXPOSE8080
    CMD java -jarspring-boot-restful-service.jar
    ADD build/spring-boot-restful-service.jar /data/spring-boot-restful-service.jar

    有了Docker以后我们就可以直接创建应用(群)了。

    ltc createlattice-app cloudfoundry/lattice-app

    如果需要扩展应用的时候我们只需要

    ltcscalelattice-app10

    这样部署了10个服务,由前端的负载均衡负责转发。

    当应用出现故障的时候,Lattice也会自动重启。

    Lattice提供了命令行工具,大部分时候需要和现有平台接入的话可以直接使用Lattice的API。

    参考链接

    迁移到云应用架构
    12factor
    Spring-Cloud
    microservices
    API-based collaboration
    Lattice API
    Spring-Docker

  • Spring boot中的Info Endpoint

    为了监控应用,Spring Boot提供了EndPoint的支持。

    目前提供了十一种,其中大部分都包含默认实现,只有其中的Info需要用户自己提供。

    Info的存在很大程度上可以替代AppCheck功能。

    最简单添加Info信息的方式就是自己在配置中写入

    info.name=app

    这样访问/info就可以获得类似这样的结果

    {
    name: 'app'
    }

    但是有一些信息是不固定的,比如版本号。

    这种可以使用打包工具去完成信息的填充。比如写成

    info.version=${version}

    在Gradle配置中加上

    processResources{
    expand(project.properties)
    }

    当然这样写有一些问题,如果你是用了Flyway等数据库版本管理工具,那么你原声的SQL文件也会被处理,视情况而定,有很高几率信息填充会失败。

    所以先过滤一下

    processResources {
    filesMatching('**/*.properties') { expand(project.properties) }
    }

    当然,有时候我们还需要一些版本库的信息,比如Git相关信息。

    可以Gradle的Git插件来获取数据,然后提供给Info。

    importorg.ajoberstar.grgit.Grgit
    
    Grgit repo = Grgit.open(project.file('.'))
    ext.git = [
    name :repo.branch.current.name,
    gitMessage:repo.head().shortMessage,
    gitTime :repo.head().getDate().toLocaleString(),
    gitId :repo.head().id
    ];

    然后在配置文件中修改

    info.build.version=${project.version}
    info.build.name=${project.name}
    info.git.name=${git.name}
    info.git.commit.message=${git.gitMessage}
    info.git.commit.time=${git.gitTime}
    info.git.commit.id=${git.gitId}

    最终效果:

    {
    git: {
    commit: {
    message:"ÇåÀígradleÎÞÓÃÅäÖÃ",
    time:"2015-7-25 20:45:44",
    id:"24d127abb1401437e752498a4878de5d2670f6e4"
    },
    name: "master"
    },
    build:{
    name: "acgmo",
    version:"0.0.2-SNAPSHOT"
    }
    }

    这个乱码问题应该是提交中有中文,但是编译服务器是英文的。

  • Spring MVC和HtmlUnit测试

    如果有一个Spring MVC项目,那么我们的测试一般有两种,一种是单元测试,第二种是端对端测试。

    单元测试可以选用Spring MVC Test框架,当然也可以当作一般的单元测试不使用SpringJUnit4ClassRunner来运行。

    而端对端测试一般是起一个服务,然后使用测试框架启动一个浏览器来测试。这两者相互结合倒也融洽。

    整个Spring项目最近多了一个新的项目,提供了一种介于其中的测试支持。

    Spring Test Htmlunit

    这个项目是Spring原有测试框架和HtmlUnit的一个结合。

    主要在于解决以下三个问题:

    • 集成常见的测试工具同时不启动服务器(当作单元测试对待)
    • 支持Javascript
    • 可以Mock一些组件来加快测试

    先来看看原有的单元测试框架是如何测试的

    MockHttpServletRequestBuilder retrieveProfile = post("/profile/")
    .param("userid","1"));
    
    mockMvc.perform(retrieveProfile)
    .andExpect(status().isOk());

    当然测试中还可以去测试Model。如果需要测试页面渲染,就需要借助xpath了。

    mockMvc.perform(get("/profile/create"))
    .andExpect(xpath("//input[@name='name']").exists())
    .andExpect(xpath("//textarea[@name='introduction']").exists());

    但是页面的交互是很负责,所以更多部分的测试是在端对端中。

    示例

    来看看新的工具是怎样解决问题的。

    首先创建一个WebClient

    WebClient webClient;
    
    @Before
    publicvoidsetup() {
    webClient = MockMvcWebClientBuilder
    .webAppContextSetup(context
    .contextPath("")
    .createWebClient();
    }

    然后直接使用WebClient去操作

    HtmlForm form = createMsgFormPage.getHtmlElementById("profile");
    HtmlTextInput summaryInput = createMsgFormPage.getHtmlElementById("name");
    summaryInput.setValueAttribute("Spring");
    HtmlTextArea textInput = createMsgFormPage.getHtmlElementById("introduction");
    textInput.setText("Do you know Spring?");
    HtmlSubmitInput submit = form.getOneHtmlElementByAttribute("input","type","submit");
    HtmlPage newProfile = submit.click();

    WebDriver

    HtmlUnit的API稍微有点繁琐,如果你习惯了Selenium的使用,那么可以考虑使用WebDriver。

    WebDriver driver;
    
    @Before
    publicvoidsetup() {
    driver = MockMvcHtmlUnitDriverBuilder
    .webAppContextSetup(context)
    .createDriver();
    }

    当然,依照惯例还是搭配Page Object Pattern使用。

    创建相关的页面对象

    publicclass CreateProfilePage
    extendsAbstractPage {
    
    
    private WebElement name;
    
    private WebElement introduction;
    
    
    @FindBy(css ="input[type=submit]")
    private WebElement submit;
    
    publicCreateMessagePage(WebDriver driver) {
    super(driver);
    }
    
    public <T> TcreateMessage(Class<T> resultPage, String name, String introduction) {
    this.name.sendKeys(name);
    this.introduction.sendKeys(introduction);
    this.submit.click();
    return PageFactory.initElements(driver, resultPage);
    }
    
    publicstatic CreateMessagePageto(WebDriver driver) {
    get(driver,"/profile/create");
    return PageFactory.initElements(driver, CreateProfilePage.class);
    }
    }

    其他支持

    该工具还支持Geb,虽然我没有使用过它,但是从文件上看确实简约了不少。

    目前该项目还没有正式释出,最新版本是1.0.0.BUILD-SNAPSHOT,可以在Spring的快照库中找到。

  • Spring Boot新模块devtools

    Spring Boot 1.3中引入了一个新的模块,devtools。

    顾名思义,这个模块是为开发者构建的,目的在于加快开发速度。

    这个模块包含在最新释出的1.3.M1中。

    dependencies {
    compile("org.springframework.boot:spring-boot-devtools")
    }

    自动禁用模板缓存

    一般情况下,View层都会应用诸如Thymeleaf之类的模版引擎,这些引擎一般会在启动或者第一次加载时编译自己,所以应用启动以后再修改它们就不会立刻生效。

    当然,这种情况下你可以禁用掉缓存已达到快速调试的目的,比如对于Thymeleaf,你需要设置spring.thymeleaf.cache为false。

    devtools会自动帮你做到这些,禁用所有模板的缓存,包括Thymeleaf, Freemarker, Groovy Templates, Velocity, Mustache等。

    自动重加载

    如果你修改了Controller类的代码,那么你只有手动重启来观察修改效果。

    当然也可以配合其他工具来达到自动重加载的目的,比如 JRebel 或者 Spring Loaded。

    现在只需要引入devtools就可以了,它会自动进行重加载。重加载时服务无法访问,下一部分的精力将会放在加快重加载速度,并尽可能自动侦测需要重新加载的类,减少不必要的开销。

    在浏览器方面,devtools内置了一个LiveReload服务,可以自动刷新浏览器。

    其他

    还有个重要的改进是远程调试,主要针对Docker和Pass平台,调试使用的是JDWP。

    目前1.3正式版还没有释出,M2版还有一个改进就是对于默认日志格式的覆盖,这也是一个直接期待的小加强。

  • 使用RxJava简化编程

    使用RxJava简化编程

    最近在做一个很有趣的项目,需求明确应用界面是console的gui。

    长成这样的

    在现在这个年代,这种朴素的需求实在少见。

    项目大量消费自有的或者外部的API,自身逻辑还是很简单。但是其中大量的异步消费远程API,解析数据,更新UI还是让人很烦。

    想起来了RxJava,代码上得到了很多简化,不敢说是最佳实践,但是换一种风格总有一些新鲜的体验。

    普通方法

    先来个简单的,一个和比特币有关的API,https://blockchain.info/ticker
    这个API返回的是当前比特币交易价格,返回内容如下:

    {
    USD: {
    15m: 246.7,
    last: 246.7,
    buy: 246.65,
    sell: 246.7,
    symbol: "$"
    },
    ISK: {
    15m: 32551.57,
    last: 32551.57,
    buy: 32544.97,
    sell: 32551.57,
    symbol: "kr"
    },
    HKD: {
    15m: 1912.37,
    last: 1912.37,
    buy: 1911.98,
    sell: 1912.37,
    symbol: "$"
    },
    TWD: {
    15m: 7628.44,
    last: 7628.44,
    buy: 7626.89,
    sell: 7628.44,
    symbol: "NT$"
    },
    CHF: {
    15m: 230.35,
    last: 230.35,
    buy: 230.31,
    sell: 230.35,
    symbol: "CHF"
    },
    EUR: {
    15m: 220.23,
    last: 220.23,
    buy: 220.19,
    sell: 220.23,
    symbol: "€"
    },
    DKK: {
    15m: 1643.43,
    last: 1643.43,
    buy: 1643.1,
    sell: 1643.43,
    symbol: "kr"
    },
    CLP: {
    15m: 156199.81,
    last: 156199.81,
    buy: 156168.15,
    sell: 156199.81,
    symbol: "$"
    },
    CAD: {
    15m: 305.33,
    last: 305.33,
    buy: 305.27,
    sell: 305.33,
    symbol: "$"
    },
    CNY: {
    15m: 1526.92,
    last: 1526.92,
    buy: 1526.62,
    sell: 1526.92,
    symbol: "¥"
    },
    THB: {
    15m: 8334.08,
    last: 8334.08,
    buy: 8332.39,
    sell: 8334.08,
    symbol: "฿"
    },
    AUD: {
    15m: 318.78,
    last: 318.78,
    buy: 318.72,
    sell: 318.78,
    symbol: "$"
    },
    SGD: {
    15m: 331.35,
    last: 331.35,
    buy: 331.28,
    sell: 331.35,
    symbol: "$"
    },
    KRW: {
    15m: 273699.26,
    last: 273699.26,
    buy: 273643.78,
    sell: 273699.26,
    symbol: "₩"
    },
    JPY: {
    15m: 30518.05,
    last: 30518.05,
    buy: 30511.86,
    sell: 30518.05,
    symbol: "¥"
    },
    PLN: {
    15m: 918.57,
    last: 918.57,
    buy: 918.39,
    sell: 918.57,
    symbol: "zł"
    },
    GBP: {
    15m: 157.26,
    last: 157.26,
    buy: 157.23,
    sell: 157.26,
    symbol: "£"
    },
    SEK: {
    15m: 2033.62,
    last: 2033.62,
    buy: 2033.2,
    sell: 2033.62,
    symbol: "kr"
    },
    NZD: {
    15m: 356.96,
    last: 356.96,
    buy: 356.89,
    sell: 356.96,
    symbol: "$"
    },
    BRL: {
    15m: 763.82,
    last: 763.82,
    buy: 763.66,
    sell: 763.82,
    symbol: "R$"
    },
    RUB: {
    15m: 13387.32,
    last: 13387.32,
    buy: 13384.61,
    sell: 13387.32,
    symbol: "RUB"
    }
    }

    处理的思路也简化一下,这里我们需要最新的以美元表示的交易价格,代码如下:

    new Runnable() {@Overridepublicvoidrun() {
    try {
    String responseAsString = Request.Get("https://blockchain.info/ticker").execute().returnContent().asString();
    Map < String,
    Object > responseAsJsonMap =new JacksonJsonParser().parseMap(responseAsString);
    Object price = ((LinkedHashMap) responseAsJsonMap.get("USD")).get("last");
    usd.setText(String.valueOf(price));
    }catch(IOException e) {
    e.printStackTrace();
    }
    }
    }.run();

    看着有些繁琐,但是功能是实现了。

    如果需要获取其他货币表示直接在其中添加即可。

    但是整个项目有很多类似的代码,都是起一个线程来请求,请求以后处理并更新UI。而且处理过程非常类似,转为JsonMap,从Map中取值。

    引入RxJava

    RxJava背后是Netflix,对于并发模型支持较好,能够利用服务器侧的并发,而无需触及典型的线程安全和同步问题。其API的实现控制了并发原语,能追求系统性能的提升,而不必担心破坏客户端代码。

    当然,还有一个名字,反应性编程,infoq中有一个专栏 http://www.infoq.com/reactive-extensions/

    具体的理念啥的就不细说了,都是能Google到的。

    来看下改进后的效果。

    首先需要一个Observable

    public Observable < String >Get(final String uri) {
    return Observable.create(new Observable.OnSubscribe < String > () {@Overridepublicvoidcall(Subscriber < ?super String > subscriber) {
    try {
    String responseAsString = Request.Get(uri).execute().returnContent().asString();
    logger.trace(responseAsString);
    subscriber.onNext(responseAsString);
    subscriber.onCompleted();
    }catch(IOException e) {
    subscriber.onError(e);
    }
    }
    });
    }

    其次需要多个map

    public Func1<String, Map<String, Object>> mapJson() {
    return json -> jsonParser.parseMap(json);
    }
    public Func1<Object, Object> expression(String l) {
    Expression expression = expressionParser.parseExpression(l);
    return o -> expression.getValue(o);
    }

    一个map是对Json的转换,一个map是对SpEL的支持,方便灵活的取值。

    最后来一个subscribe来更新UI

    public Action1<Object>updateComponent(AbstractComponent component) {
    return o -> {
    Method setText = ReflectionUtils.findMethod(component.getClass(),"setText", String.class);
    if (setText !=null) {
    String s = String.valueOf(o);
    ReflectionUtils.invokeMethod(setText, component, s);
    }
    };
    }

    最后组合在一起

    helper.Get("https://blockchain.info/ticker")
    .map(helper.mapJson())
    .map(helper.expression("[USD][last]"))
    .subscribe(helper.updateComponent(usd));

    效果如下:

    比特币美元价

    定时更新

    这里是一个单纯的数据展示,而且只会执行一次,如果需要定时执行,RxJava也有相关的支持。

    RxJava提供了Schedulers来处理这种情况。具体又包括了三种

    • computationScheduler
    • ioScheduler
    • newThreadScheduler

    这里选用newThreadScheduler。

    worker = Schedulers.newThread().createWorker();
    worker.schedulePeriodically(new Action0() {@Overridepublicvoidcall() {
    //原有的代码
    }
    },
    0,2000, TimeUnit.MILLISECONDS);

    多个订阅者

    如果我们需要根据统一结果做出不同的操作,比如获取API数据以后,获取两个不同的值,然后更新到两个不同的地方,那么我们就需要两个订阅者了。

    使用publish方法可以声明该结果对外公布,然后添加所有所需的订阅者后connect。

    worker = Schedulers.newThread().createWorker();
    worker.schedulePeriodically(new Action0() {@Overridepublicvoidcall() {
    ConnectableObservable < String > publish = helper.Get("https://blockchain.info/ticker").publish();
    publish.map(helper.mapJson()).map(helper.expression("[USD][last]")).subscribe(helper.updateComponent(usd));
    publish.map(helper.mapJson()).map(helper.expression("[CNY][last]")).subscribe(helper.updateComponent(cny));
    publish.connect();
    }
    },
    0,2000, TimeUnit.MILLISECONDS);

    最后效果

    写在最后

    RxJava的API无疑是简单的,背后的功能和思路更是值得学习的。

    所谓的反应性编程其实本质是拉模型和推模型的思考。程序=数据+算法,那么数据源是由数据源将数据推给算法,算法被动地等待数据的到来还是数据源等待算法的访问,算法主动地将数据从数据源中拉出?

    回调函数是典型的推模型的应用,设计模式中的 Command,Visitor 中都有回调函数的影子。

    拉模型符合人的思考逻辑,状态由运行时的栈自动维护,理解容易;推模型性能好,但写起来不容易,而且需要自己维护相关状态。

  • 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);
    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().show();

    描述性统计

    当然,可以看到直接调用describe会返回所有列的数据,如果只需要特定的数据,可以直接指定

          
    1
    2
    3
    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();

    局部描述性统计

    describe实质上返回的是另外一个数据框,其中具体计算的量由expr指定

    val statistics = List[(String, Expression => Expression)](
    "count" -> Count,
    "mean" -> Average,
    "stddev" -> stddevExpr,
    "min" -> Min,
    "max" -> Max)

    样本协方差和相关性

    在概率论和统计学中,协方差用于衡量两个变量的总体误差。
    如果两个变量的变化趋势一致,也就是说如果其中一个大于自身的期望值,另外一个也大于自身的期望值,那么两个变量之间的协方差就是正值。
    如果两个变量的变化趋势相反,即其中一个大于自身的期望值,另外一个却小于自身的期望值,那么两个变量之间的协方差就是负值。
    如果两个量是统计独立的,那么二者之间的协方差就是0。

    从理论上讲两列随机数的相关性应该为零,用Spark来计算一下。

    SQLContext sqlContext =new SQLContext(sc);
    DataFrame dataFrame = sqlContext.range(1,10);
    DataFrame frame = dataFrame.select(functions.col("id"), functions.rand(10).alias("randX"), functions.randn(10).alias("randY"));
    double cov= frame.stat().cov("randX","randY");
    System.out.printf("cov: %s%n", cov);

    协方差

    而相关性相对来说更容易理解,id是固定值,所以id和id自身是强相关的。

    SQLContext sqlContext =new SQLContext(sc);
    DataFrame dataFrame = sqlContext.range(1,10);
    DataFrame frame = dataFrame.select(functions.col("id"), functions.rand(10).alias("randX"), functions.rand(10).alias("randY"));
    double id_corr = frame.stat().corr("id","id");
    System.out.printf("id corr: %s%n", id_corr);

    相关系数

    列联表和频繁项

    从理论上讲这二者更像是一种操作而不是指标。

    Spark数据框中提供的这两个功能更多的是为了快速的数据处理。

    以列联表为例,

    首先生成一个姓名和购买物品的两列表

    SQLContext sqlContext =new SQLContext(sc);
    String[] names = {
    "Xiaoming",
    "Xiaohong",
    "Xiaogang"
    };
    String[] items = {
    "milk",
    "bread",
    "butter",
    "apples",
    "oranges"
    };
    
    Random random =new Random();
    List < Consumption > consumptions = Lists.newArrayList();
    for (String name: names) {
    for (int i =0; i <5; i++) {
    consumptions.add(new Consumption(name, items[random.nextInt(items.length)]));
    }
    }
    JavaRDD < Consumption > parallelize = sc.parallelize(consumptions);
    DataFrame dataFrame = sqlContext.createDataFrame(parallelize, Consumption.class);
    dataFrame.registerTempTable("consumption");
    dataFrame.show()

    原表格

    然后联合

    dataFrame.stat().crosstab("name","item").show();

    列联表

    写在最后

    本文提到Spark对于数理统计的支持都在1.4版本中,目前还在rc4,可以通过快照版体验。

    新的Spark MLlib包中也由于数据框的提升有了不少变动,特别是机器学习甬道的概念的融合,数据框将承担越来越多基本数据分析的功能。

  • 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: Vector, label: Double, weights: Vector, cumGradient: Vector): Double
    }

    而更新抽象为Updater.scala

    abstractclass Updater extends Serializable {
    def compute(
    weightsOld: Vector,
    gradient: Vector,
    stepSize: Double,
    iter: Int,
    regParam: Double): (Vector, Double)
    }

    而具体的实现梯度有

    • LogisticGradient
    • LeastSquaresGradient
    • HingeGradient

    对于更新也是三种实现

    • SimpleUpdater
    • L1Updater
    • SquaredL2Updater

    Spark本身包含了很多内置的机器学习算法,比如线性回归,分类等等,它们都有依赖梯度下降算法,所有单纯从使用角度并不需要太关心这个,因为使用的时候一般不需要具体指定。

    比如线性回归:

    IsotonicRegressionModel model =new IsotonicRegression().setIsotonic(true).run(training);

    直接使用梯度下降

    出于很多原因,大部分时候需要直接调用梯度下降,或者自己实现更新器,达到自定义的目的。

    这里以简单的y=3*x+1为例来简单使用一下

    测试数据就随意

    1 0 1
    7 2 1
    10 3 1
    4 1 1
    19 6 1

    先将它们转为RDD对象

    List < Tuple2 < Object,Vector >> list = Lists.newArrayList(
    new Tuple2 < >(1d, Vectors.dense(0.0d,1d)),
    new Tuple2 < >(7d, Vectors.dense(2.0d,1d)),
    new Tuple2 < >(10d, Vectors.dense(3.0d,1d)),
    new Tuple2 < >(4d, Vectors.dense(1.0d,1d)),
    new Tuple2 < >(19d, Vectors.dense(6.0d,1d))
    );
    JavaPairRDD < Object,Vector > rdd = jsc.parallelizePairs(list);

    构造求解器

    Gradient gradient =new LeastSquaresGradient();
    Updater updater =new L1Updater();
    Optimizer descent =new GradientDescent(gradient, updater);

    调用求解

    Vector optimize = descent.optimize(rdd.rdd(), Vectors.dense(0d,0d));
    System.out.println(optimize);

    输出如图:

    输出

  • 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");
    HashingTF hashingTF =new HashingTF()
    .setNumFeatures(1000)
    .setInputCol(tokenizer.getOutputCol())
    .setOutputCol("features");
    LogisticRegression lr =new LogisticRegression()
    .setMaxIter(10)
    .setRegParam(0.01);
    Pipeline pipeline =new Pipeline()
    .setStages(new PipelineStage[] {tokenizer, hashingTF, lr});

    然后评估模型

    PipelineModel model = pipeline.fit(training);

    结构如图:

    文本分类甬道

    主要实现了对应的接口,就可以很容易自定义转化和评估。
    对应的例子可以参考JavaCrossValidatorExample。

    变动

    Spark 1.2开始引入了ML包,目前最新版本是1.3.1,1.4.0候选版本已经释出。

    从1.2到1.3有一个比较主要大的变动是SchemaRDD被DataFrame替换,其他的主要是相关库的变动。

  • 使用Gradle注册SpringXD的module

    使用Gradle注册SpringXD的module

    SpringXD是Pivotal的大数据产品,提供了一个抽象的数据处理平台。

    SpringXD将数据解决方案抽象为数据吸纳,分析,流调度和输出四大块。

    作为源头的数据吸纳可以从各种数据源中获取需要的数据,基于Spring的另外一个项目spring-integration,这一部分的大部分实现都可以使用简单的dsl实现。

    在SpringXD中每一个部分的组件都可以自己编写并注册到服务中,方便之后的使用。

    但是在编写阶段每次都需要打包,上传服务器,然后注册,整个过程还是有点繁琐的。

    Gradle作为一个构建工具,自然可以通过自定义任务完成这个任务。

    微博数据吸纳

    这里以微博的数据源为例。

    微博的数据源因为新浪微博提供了sdk,所以自己编写稍微方便一些。

    SpringXD项目默认也提供了twitter的两个source方便测试。

    先看一下build.gradle的配置

    ext {
    xdVersion ='1.1.1.RELEASE'
    springVersion ='4.1.3.RELEASE'
    moduleType ='source'
    moduleName ='weibo'
    xdServer ="http://192.168.0.12:9393"
    version ='0.0.1-SNAPSHOT'
    }

    最关键的配置变量是moduleType,moduleName和xdServer。

    为了简单明了,WeiboSource并不具有配置参数(一般情况下都需要提供一些配置参数的,比如apiKey),主配置文件如下

    <?xml version="1.0" encoding="UTF-8"?>
    <beans:beans xmlns="http://www.springframework.org/schema/integration"
    xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
    xmlns:beans="http://www.springframework.org/schema/beans"
    xsi:schemaLocation="http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd
    http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
    
    <channel id="output"/>
    
    <beans:bean class="com.huangyunkun.xd.WeiboSource">
    <beans:property name="autoStartup" value="false"/>
    <beans:property name="outputChannel" ref="output"/>
    </beans:bean>
    
    </beans:beans>

    因为打包后配置需要也在jar中,在jar配置中增加

    jar {
    from("modules/${moduleType}/${moduleName}/config") {
    into"config"
    }
    }

    注册到SpringXD

    SpringXD提供了一套Restful的api用于常用的操作,比如stream的管理,容器状态等。

    xd shell也是调用了这个接口,也有java的实现。

    在build.gradle中加入

    buildscript {
    repositories {
    jcenter()
    }
    dependencies {
    classpath('org.springframework.xd:spring-xd-rest-client:+')
    }
    }

    新定义一个任务reg

    task reg(dependsOn: jar) << {
    SpringXDTemplatetemplate = newSpringXDTemplate(newURI(xdServer));
    ModuleOperations moduleOperations =template.moduleOperations();
    moduleOperations.uploadModule(moduleName,RESTModuleType.valueOf(moduleType), new
    FileSystemResource(jar.archivePath.path),true);
    }

    然后执行gradlew reg即可注册成功。

    注册成功

    注册成功后就可以正常使用了

    创建新的stream
    stream输出

    当然管理界面也可以看到相关信息

    管理界面

    其他问题

    • 这种注册方式是强制的,也就是说如果有同名source就会被覆盖。

    • 因为Module的类型是通过RESTModuleType.valueOf(moduleType)获取的,所以也可以这样注册sink等其他类型的模块。

    • 如果有对应的module正在被使用,那么是没法重新注册的。

  • 在Gradle中限制对jar包签名时机

    Maven仓库是一个包含大量依赖库的地方,有时候我们需要发布自己的库到仓库。

    仓库虽然对于发布的库的具体功能和作用没有太多要求,但是有一些强制要求是必须的。

    发布的内容物可以是jar,aar等,但是都必须满足一下条件:

    • Metadata(pom.xml)
    • 签名
    • source jar
    • javadoc jar

    Gradle自身包含了mvn插件和signing插件,所以整个工作还是比较简单的,大体配置如下:

    signing {
    required {
    isReleaseVersion && gradle.taskGraph.hasTask("uploadArchives")
    }
    signconfigurations.archives
    }
    task javadocJar(type: Jar, dependsOn: javadoc) {
    classifier ='javadoc'from'build/docs/javadoc'
    }
    task sourcesJar(type: Jar) {
    classifier ='sources'fromsourceSets.main.allSource
    }

    如果是快照版本就会发布到snapshot仓库,否则就是staging库。

    对于快照版本是不需要签名的,所以在操作中尽量快照版本不签名,一来是节约时间,二来是快照版本如果通过CI发布的,就可以省去很多事。

    signing插件的配置中通过指定required来决定是否跳过签名。

    可以通过以下代码来判断

    ext.isReleaseVersion = !version.endsWith("SNAPSHOT") signing {
    required {
    isReleaseVersion && gradle.taskGraph.hasTask("uploadArchives")
    }
    sign configurations.archives
    }

    当然如果你的CI中有特殊的环境变量,也可以加入判断中。

    比如SnapCI中的环境变量有

    $ exportSNAP_CI=true
    $ exportCI=true
    $ exportSNAP_CACHE_DIR=/var/go
    $ exportLANG=en_US.UTF-8
    $ exportLC_ALL=en_US.UTF-8
    $ exportSNAP_WORKING_DIR=/var/snap-ci/repo
    $ exportSNAP_PIPELINE_COUNTER=12
    $ exportSNAP_STAGE_NAME=upload
    $ exportSNAP_BRANCH=master
    $ exportSNAP_COMMIT=768980f107b8757812f23a934979549e40977c3d
    $ exportSNAP_COMMIT_SHORT=768980f
    $ exportSNAP_TRACKING_PIPELINE=true
    $ exportSNAP_INTEGRATION_PIPELINE=false