暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

jstorm的恰好一次事务处理

考拉苑 2017-01-13
399

一、新版jstorm2.1.1

概述

在原有的Storm设计中,Trident支持了只处理一次,Acker支持至少处理一次场景。但是Trident和Acker在这两种消息保证机制中,都面临着同样的性能问题, 如果用户需要exactly-once 保证, 性能会急剧下降。

同时由于Trident和Acker机制的API完全不同,用户很难用一套代码去支持两个场景。

针对这两个问题,我们对开始考虑做一套新的框架来支持者两个消息保证机制。

基本原理:

设计参考了Flink的只处理一次方案,barrier和stream align(流对齐)的设计,去做batch的划分和如何保证每个节点只会处理当前批次的消息。具体可以参考Flink的官方文档。 https://ci.apache.org/projects/flink/flink-docs-release-1.0/internals/stream_checkpointing.html 对于barrier和流对齐,本文不详细展开,主要介绍JStorm的实现和相关接口的使用。

基本介绍

  • 状态维护: JStorm把topology任务的节点分为了四类节点,数据源spout,状态无关节点non-stateful bolt,状态节点stateful bolt,结束节点end bolt。


    • 数据源节点和状态节点只维护自己节点的状态提交顺序,在处理完每个batch后,会马上处理下一个batch的消息。


    • Topology master负责每个batch的全局状态维护,既是否所有节点都完成了当前batch的处理。如果完成,topology master开始做全局snapshot状态的持久化存储。

  • 如何回滚:


    • Barrier和stream align让每个节点在做checkpoint时,保证batch的顺序性和一致性,而topology master维护了全局状态。如果失败topology master会通知所有状态节点 回滚到最后一次全局成功的checkpoint,然后数据源开始从最后一次全局成功的位置开始重新拉取消息。


    • JStorm里回滚以数据源的种类为最小单位。比如说我们现在有两个spout component,TT spout和metaq spout。如果metaq spout的下游节点挂了,我们只会回滚,metaq spout这个流的数据。

接口介绍

Spout: ITransactionSpoutExecutor

Bolt: ITransactionBoltExecutor, ITransactionStatefulBoltExecutor

应用的消息处理接口和JStorm的基础API一致,如果用户从原有的JStorm任务迁移过来,消息的处理逻辑不用改变。只需要在spout和状态bolt节点实现以下状态处理相关的接口。

        /**
          *  Init state from user checkpoint
          *  @param userState user checkpoint state 
          */
        public  void  initState(Object  userState);
        /**
          *  Called when current batch is finished
          *  @return user state to be committed
          */
        public Object finishBatch();
        /**
          *  Commit a batch  
          *  @param The user state data
          *  @return snapshot state which is used to retrieve the persistent user state
          */
        public Object commit(BatchGroupId id, Object state);
        /**
          *  Rollback from last success checkpoint state  
          *  @param user state for rollback
          */
        public void rollBack(Object userState);
        /**
          *  Called when the whole topology finishes committing
          */
        public void ackCommit(BatchGroupId id);
  • initState和rollback: 是在任务起来和回滚时,用于让当前节点回到最后一次成功batch的checkpoint。

  • finishBatch: 是用于通知用户,当前batch接收完毕。返回值为当前batch对应的业务checkpoint。

  • commit: 在用户返回对应的checkpoint后,我们会在一个异步线程里调用commit接口,对这个checkpoint进行提交操作,以不阻塞task execute处理线程。如果checkpoint不大,可以直接在这个接口中返回checkpoint,JStorm最终会在topology master里统一做持久化。如果checkpoint比较大。用户可以在这个接口中,把对应的checkpoint缓存一份到相应外部存储里,然后返回对应的key。在回滚时,JStorm会把该key值返回给用户,用于找回对应的checkpoint。

  • ackCommit: 当topology任务在所有节点成功后,该接口会被回调。如果之前有checkpoint的缓存,用户在这个接口里可以去完成一些清理工作。比如batch-10完成了,那我们可以删除batch-9和之前的所有checkpoint缓存。

参考例子: sequence-split-merge模块中TransactionTestTopology.java类。

架构优势

相对于Trident和Acker机制。这套架构的最大优势在于: 1. 每个节点只关心自己节点的提交状态。不用等待所有节点都成功后,再开始下一个batch的提交和计算。 2. 不需要再依赖acker节点,减少额外系统计算和带宽消耗。原有的acker模型下,每条用户消息都会产生一条对应的系统消息到acker,同时有额外的计算消耗,并且acker消息会消耗大量的网络带宽。

相关配置

    # true:  只处理一次;  false:  至少处理一次
   transaction.exactly.once.mode:  true

二、历史版本 trident的处理

概述

老的storm的事务主要用于对数据准确性要求非常高的环境中,尤其是在计算交易金额或笔数,数据库同步的场景中。

老的storm 事务逻辑(需要使用trident)是挺复杂的,而且坦白讲,代码写的挺烂的。 JStorm的新事务请参考JStorm新事务

Storm 事务的核心设计思想:

Transaction 还是基于api的属性之上,做的一层封装,从而满足transaction

其实,相当于把一个batch当做一个原子tuple来处理,只是中间计算的过程,可以并发。

核心设计1

提供一个strong order,也就是,如果一个tuple没有被完整的处理完,就不会处理下一个tuple,说简单一些,就是,采用同步方式。并对每一个tuple赋予一个transaction ID,这个transaction ID是递增属性(强顺序性),如果每个bolt在处理tuple时,记录了上次的tupleID,这样即使在failure replay时能保证处理仅且处理一次

核心设计2

如果一次处理一个tuple,性能不够好,可以考虑,一次处理一批(batch tuples) 这个时候,一个batch为一个transaction处理单元,当一个batch处理完毕,才能处理下一个batch,每个batch赋予一个transaction ID。

核心思想3

如果在计算任务中,并不是所有步骤需要强顺序性,因此将一个计算任务拆分为2个阶段: 1. processing 阶段:这个阶段可以并发 2. commit阶段:这个阶段必须强顺序性,因此,一个时刻,只要一个batch在被处理 任何一个阶段发生错误,都会完整重发batch

结果

一次性从Meta或Kafka 中取出一批数据,然后一条一条将数据发送出去,当所有数据均被正确处理后, 触发一个commit 流,这个commit流是严格排序,通常在commit流中进行flush动作或刷数据库动作,如果commit流最后返回也成功,spout 就更新Meta或kafka的偏移量,否则,任何一个环节出错,都不会更新偏移量,也就最终重复消费这批数据。

具体代码逻辑

代码逻辑参考下图



详细内容建议移步到: http://jstorm.io/index_cn.html

源码:https://github.com/alibaba/jstorm/blob/master/history_cn.md(版本升级说明)

         https://github.com/alibaba/jstorm

文章转载自考拉苑,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论