ack是什么,如何使用Ack机制,如何关闭Ack机制,基本实现,STORM的消息容错机制,Ack机制

1、ack是什么

ack 机制是storm整个技术体系中非常闪亮的一个创新点。

通过Ack机制,spout发送出去的每一条消息,都可以确定是被成功处理或失败处理, 从而可以让开发者采取动作。比如在Meta中,成功被处理,即可更新偏移量,当失败时,重复发送数据。 
因此,通过Ack机制,很容易做到保证所有数据均被处理,一条都不漏。 
另外需要注意的,当spout触发fail动作时,不会自动重发失败的tuple,需要spout自己重新获取数据,手动重新再发送一次

ack机制即, spout发送的每一条消息,

? 在规定的时间内,spout收到Acker的ack响应,即认为该tuple 被后续bolt成功处理 
? 在规定的时间内,没有收到Acker的ack响应tuple,就触发fail动作,即认为该tuple处理失败, 
? 或者收到Acker发送的fail响应tuple,也认为失败,触发fail动作

另外Ack机制还常用于限流作用: 为了避免spout发送数据太快,而bolt处理太慢,常常设置pending数,当spout有等于或超过pending数的tuple没有收到ack或fail响应时,跳过执行nextTuple, 从而限制spout发送数据。

通过conf.put(Config.TOPOLOGY_MAX_SPOUT_PENDING, pending);设置spout pend数。

2、如何使用Ack机制

spout 在发送数据的时候带上msgid

设置acker数至少大于0;Config.setNumAckers(conf, ackerParal); 
在bolt中完成处理tuple时,执行OutputCollector.ack(tuple), 当失败处理时,执行OutputCollector.fail(tuple); 
推荐使用IBasicBolt, 因为IBasicBolt 自动封装了OutputCollector.ack(tuple), 处理失败时,请抛出FailedException,则自动执行OutputCollector.fail(tuple)

3、如何关闭Ack机制

有2种途径

spout发送数据时不带上msgid 
设置acker数等于0

4、基本实现

Storm 系统中有一组叫做”acker”的特殊的任务,它们负责跟踪DAG(有向无环图)中的每个消息。 
acker任务保存了spout id到一对值的映射。第一个值就是spout的任务id,通过这个id,acker就知道消息处理完成时该通知哪个spout任务。第二个值是一个64bit的数字,我们称之为”ack val”, 它是树中所有消息的随机id的异或计算结果。

<TaskId,<RootId,ackValue>>
Spoutid,<系统生成的id,ackValue>
Task-0,64bit,0
  • 1
  • 2
  • 3
  • 1
  • 2
  • 3

ack val表示了整棵树的的状态,无论这棵树多大,只需要这个固定大小的数字就可以跟踪整棵树。当消息被创建和被应答的时候都会有相同的消息id发送过来做异或。 每当acker发现一棵树的ack val值为0的时候,它就知道这棵树已经被完全处理了 
 
 
 

STORM的消息容错机制http://www.woaipu.com/shops/zuzhuan/61406

数据在处理中出现异常时,需要保证消息被完整处理。
-----------------------
SPOUT --A---B---C---D
期望:当其中一个环节出现异常时,Spout能够重新发送一份数据。
-----------------------
问题:SPOUT如何知道一条消息的处理状态
        成功:ack(Object msgid)
        失败:fail(Object msgid)
    :Bolt如何告知Spout消息处理的状态
        collector.emit(new Value())
        collector.ack() //当消息处理成功时
        collector.fail()//当消息处理失败时
------------------

Ack机制http://www.woaipu.com/shops/zuzhuan/61406

    Spout发送一条数据出去,需要知道数据处理成功和失败的状态,如果失败进行消息的重新发送
    1、自定义spout实现BaseRichSpout,覆写ack,fail方法。
    2、在自定义的spout发送数据的时候,需要制定messageid,messageid是一个Object。
    3、当消息处理成功或失败之后,Storm框架会将messageId传回来。
        如果消息要重发,直接通过messageId找到或直接转化成数据内容进行重发。
    4、自定义Bolt实现BaseRichBolt
    5、在bolt的execute中进行两个操作
        5.1、发送数据时,需要指定血缘关系,锚点
            collector.emit(父tuple,new 子Tuple)
        5.2、当execute处理完业务逻辑的时候,需要告诉storm框架当前阶段的处理状态。
            collector.ack(tuple)
如果在编写storm程序时,在bolt环节忘了手动ack或fail,怎么办?
    忘了手动ack或fail,storm框架会等待反馈,达到超时阈值之后,就直接给fail。
如果在编写storm程序时,在bolt环节忘了标识锚点,怎么办?
    忘了标识锚点,就是忘了标识血缘关系。storm会认为你不关心后面阶段的处理状况。

Storm BaseRichBolt API 过于繁琐,就开了另外一个api:BaseBasicBolt
    如果实现了BaseBasicBolt,就不需要锚点,不需要手动ack或fail。

http://www.woaipu.com/shops/zuzhuan/61406http://www.woaipu.com/shops/zuzhuan/61406
时间: 2024-11-03 05:28:30

ack是什么,如何使用Ack机制,如何关闭Ack机制,基本实现,STORM的消息容错机制,Ack机制的相关文章

storm 消息的可靠处理机制——Ack整个tuple树异或

消息的可靠处理机制 Storm内部通过一种巧妙的异或算法判读每个tuple是否被正确完整的处理. Spout的一个Task创建一个Tuple时,即在Spout的nextTuple()方法中实现从特定数据源读取数据的处理逻辑中,会与Acker进行通信,向Acker发送消息,Acker保存该Tuple对应信息:{:spout-task task-id :val ack-val)}. Bolt在emit一个新的子Tuple时,会保存子Tuple与父Tuple的关系. 在Bolt中进行ack时,会计算出

ActiveMQ源码解析(四):聊聊消息的可靠传输机制和事务控制

在消息传递的过程中,某些情况下比如网络闪断.丢包等会导致消息永久性丢失,这时消费者是接收不到消息的,这样就会造成数据不一致的问题.那么我们怎么才能保证消息一定能发送给消费者呢?怎么才能避免数据不一致呢?又比如我们发送多条消息,有时候我们期望都发送成功但实际上其中一部分发送成功,另一部分发送失败了,没达到我们的预期效果,那么我们怎么解决这个问题呢? 前一种问题我们通过消息确认机制来解决,它分为几种模式,需要在创建session时指定是否开启事务和确认模式,像下面这样: <span style=&quo

RabbitMQ消息队列:ACK机制

每个Consumer可能需要一段时间才能处理完收到的数据.如果在这个过程中,Consumer出错了,异常退出了,而数据还没有处理完成,那么 非常不幸,这段数据就丢失了. 因为我们采用no-ack的方式进行确认,也就是说,每次Consumer接到数据后,而不管是否处理完 成,RabbitMQ Server会立即把这个Message标记为完成,然后从queue中删除了. 如果一个Consumer异常退出了,它处理的数据能够被另外的Consumer处理,这样数据在这种情况下就不会丢失了(注意是这种情况

Storm消息可靠性的保障机制

参考[并发编程网]的Storm官方教程翻译 以WordCountToPology为例: // 构造Topology TopologyBuilder builder = new TopologyBuilder(); builder.setSpout(SPOUT_ID,new SentenceSpout(), 2)// 指定 Spout ,2 指的是使用2个executor来运行spout .setNumTasks(4);//指定tasks的数量 // 指定 SentenceSpout 向Split

springboot项目整合rabbitMq涉及消息的发送确认,消息的消费确认机制

1.引入maven依赖 <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency>2.在application.yml的配置: spring: rabbitmq: host: 106.52.82.241 port: 5672 username: yang

重磅消息!AppCan扩容机制上线,扩大空间随心所欲!

亲爱的AppCan开发者: AppCan在线打包空间扩容机制已经正式上线啦! 开发的项目越来越多,可用的空间却越来越小,怎么办?[扩容空间]让您不再为空间"斤斤计较"! VS 扩容规则: 1.请关注AppCan官方微信号,并在申请原因中注明您的微信号. AppCan微信二维码 2.有至少一个上线的AppCan应用,可申请50M空间,经官方人员审核后即可扩容. 3.任何一个应用安装量超过1000次可申请一次200M空间. 4.3个月内不可重复申请,且3个月后需凭借其他符合上述要求的应用才

Objective-C 消息发送与转发机制原理

消息发送和转发流程可以概括为:消息发送是 Runtime 通过 selector 快速查找 IMP 的过程,有了函数指针就可以执行对应的方法实现:消息转发是在查找 IMP 失败后执行一系列转发流程的慢速通道,如果不作转发处理,则会打日志和抛出异常. http://www.huanbohailawyer.com/e/space/?userid=52858?feed_filter=ks&lk20160609=&85 http://www.huanbohailawyer.com/e/space/

(总结)高并发消息队列常用通知机制

最近在研究一个高性能的无锁共享内存消息队列,使用的fifo来通知.结合之前<基于管道通知的百万并发长连接server模型>文章,这里总结一下常用的通知机制. 常用的通知机制中比较典型的有以下几种: 1.signal 这种机制下,我们向被通知进程发送一个特殊的signal(比如SIGUSR1),这样正在睡眠的读进程就会被信号中断,然后醒来. 该方法的优点是:读进程不需要监听一个额外的eventfd,适合一些不方便使用eventfd的场景:另外,用户可以选择是使用实时信号(SIGRTMIN+1),

macOS 的 rootless 机制的关闭与打开

一.现象 升级系统 经常遇到文件 operation not permitted 问题. 二.原因 macOS 在升级系统之后,电脑启用了SIP(System Integrity Protection),增加了rootless机制,导致即使在root权限下依然无法修改文件,在必要时候为了能够修改下面的文件,我们只能关闭该保护机制. 三.解决方案 重启mac,在重启黑屏的时候,按住 Command + R 进入恢复模式,注意查找到[菜单]里面的[终端]工具. 1. 关闭 终端输入 csrutil