三、消息的可靠处理

2024-09-05 08:58
文章标签 处理 消息 可靠

本文主要是介绍三、消息的可靠处理,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

1、消息被完全处理的含义


当树创建完毕,并且树中的每一个消息都已经被处理时,Storm认为来自Spout的元组是“完全处理”的。当一个元组的消息树在指定的超时范围内不能被完全处理,则元组被认为是失败的。


2、如果一个消息被完全处理或完全处理失败会发生什么

首先,让我们看看Spout的元组的生命周期。ISpout接口的定义如下:

public interface ISpout extends Serializable {void open(Map conf, TopologyContext context, SpoutOutputCollector colle- ctor);void close();void nextTuple();void ack(Object msgId);void fail(Object msgId);
}

首先,Storm通过调用Spout的nextTuple()方法从Spout请求一个元组。Spout使用open()方法提供的SpoutOutputCollector对象发射一个元组到它的输出流。当发射元组时,Spout会提供一个“消息id”,以便用来识别元组。例如,KestrelSpout从Kestrel消息队列中读取一个消息时,会发射Kestrel提供的“消息id”。下面发射一个消息到SpoutOutputCollector对象:

_collector.emit(new Values("field1", "field2", 3) , msgId);

接下来,元组被发送到Bolt,同时Storm负责跟踪创建的消息树。如果Storm检测到一个元组是完全处理的,Storm将调用原Spout任务的ack()方法,把Spout提供给Storm的消息id作为输入参数。同样,如果元组超时,Storm将调用Spoutfail()方法。注意,一个元组将由Spout任务来确认成功或失败,这个Spout任务是创建此元组的完全相同的Spout任务。如果一个Spout跨集群执行很多任务,元组是不会被创建它的那个任务外的其他任务确认成功或失败的。

注意:bolt是没有ack()和fail()函数的,任何消息出错了,都是由根spout重发


3、Storm如何保证可靠性

在元组树中指定一个链接,此链接被称为锚定(Anchoring。Anchoring在发射一个新的元组的同一时间完成。让我们使用以下Bolt为例进行介绍,这个Bolt将包含一个句子的元组划分为一个包含每个单词的锚定:

public class SplitSentence extends BaseRichBolt {OutputCollector _collector;public void prepare(Map conf, TopologyContext context, OutputCollector 
collector) {_collector = collector;}public void execute(Tuple tuple) {String sentence = tuple.getString(0);for(String word: sentence.split(" ")) {_collector.emit(tuple, new Values(word));//锚定+发射}_collector.ack(tuple);//确认}public void declareOutputFields(OutputFieldsDeclarer declarer) {declarer.declare(new Fields("word"));}
}<

通过指定输入元组作为第一个参数来发射每个单词元组被锚定anchored。因为这个单词元组是被锚定的如果单词元组未能被下游处理,树的根的Spout元组将在稍后重发。相反,如果单词元组的发射操作如下,让我们看看会发生什么:

_collector.emit(new Values(word));

这种方式发射的单词元组导致未被锚定unanchored。如果元组未被下游处理,根元组将不会重发。这取决于你需要的Topology的容错保证,有时候需要相应地发射一个未被锚定的元组。

很多Bolt遵循一个读取一个输入元组,发射元组,在execute方法确认元组的通用模式。这些Bolt具有类别过滤器和简单的功能。Storm有一个接口称为BasicBolt,为你封装这个模式。SplitSentence的例子可以使用BasicBolt写成:

public class SplitSentence extends BaseBasicBolt {public void execute(Tuple tuple, BasicOutputCollector collector) {String sentence = tuple.getString(0);for(String word: sentence.split(" ")) {collector.emit(new Values(word));}}public void declareOutputFields(OutputFieldsDeclarer declarer) {declarer.declare(new Fields("word"));}        
}

这个实现与之前的实现相比,语义上是相同的,但是更简单。元组发射到BasicOutputCollectorare是自动Anchoring到输入元组execute方法完成时,输入元组自动为你确认。


4、Storm如何实现可靠性

       storm系统中有一组叫做acker(可以设置并行度为多个)的特殊任务,它负责跟踪DAG中的每个消息。每当发现一个DAG被完全处理,它就向创建这个根消息的Spout任务发送一个信号,该tuple tree已经被完全处理成功。

      系统使用一种哈希算法根据Spout消息的id来确定由哪个acker跟踪此消息派生出来的tuple tree。因为每个消息都知道与之对应的根消息的id(每当Bolt新生成一个消息,对应的tuple tree中的根消息id就复制到这个消息中),所以它知道应该与哪个acker通信。当这个消息被应答的时候,它就把关于tuple tree变化的信息发送给跟踪这棵树的acker。例如,它会告诉acker:“本消息已经处理完毕,但是我派生出来了一些新的消息,帮忙跟踪一下吧”

       一个Acker任务存储来自Spout元组id到一对值的映射。第一个值是创建Spout元组的任务id,通过这个ID,acker就知道消息处理完成时该通知哪个spout任务。第二个值是一个64位的值,称为ack val。它是所有消息的随机id的异或结果当一个Acker任务看到ack val已经成为0它就知道元组树已经完成了。



这篇关于三、消息的可靠处理的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1138509

相关文章

如何通过Python实现一个消息队列

《如何通过Python实现一个消息队列》这篇文章主要为大家详细介绍了如何通过Python实现一个简单的消息队列,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下... 目录如何通过 python 实现消息队列如何把 http 请求放在队列中执行1. 使用 queue.Queue 和 reque

使用C++将处理后的信号保存为PNG和TIFF格式

《使用C++将处理后的信号保存为PNG和TIFF格式》在信号处理领域,我们常常需要将处理结果以图像的形式保存下来,方便后续分析和展示,C++提供了多种库来处理图像数据,本文将介绍如何使用stb_ima... 目录1. PNG格式保存使用stb_imagephp_write库1.1 安装和包含库1.2 代码解

C#使用DeepSeek API实现自然语言处理,文本分类和情感分析

《C#使用DeepSeekAPI实现自然语言处理,文本分类和情感分析》在C#中使用DeepSeekAPI可以实现多种功能,例如自然语言处理、文本分类、情感分析等,本文主要为大家介绍了具体实现步骤,... 目录准备工作文本生成文本分类问答系统代码生成翻译功能文本摘要文本校对图像描述生成总结在C#中使用Deep

Spring Boot 整合 ShedLock 处理定时任务重复执行的问题小结

《SpringBoot整合ShedLock处理定时任务重复执行的问题小结》ShedLock是解决分布式系统中定时任务重复执行问题的Java库,通过在数据库中加锁,确保只有一个节点在指定时间执行... 目录前言什么是 ShedLock?ShedLock 的工作原理:定时任务重复执行China编程的问题使用 Shed

解读Redis秒杀优化方案(阻塞队列+基于Stream流的消息队列)

《解读Redis秒杀优化方案(阻塞队列+基于Stream流的消息队列)》该文章介绍了使用Redis的阻塞队列和Stream流的消息队列来优化秒杀系统的方案,通过将秒杀流程拆分为两条流水线,使用Redi... 目录Redis秒杀优化方案(阻塞队列+Stream流的消息队列)什么是消息队列?消费者组的工作方式每

使用C/C++调用libcurl调试消息的方式

《使用C/C++调用libcurl调试消息的方式》在使用C/C++调用libcurl进行HTTP请求时,有时我们需要查看请求的/应答消息的内容(包括请求头和请求体)以方便调试,libcurl提供了多种... 目录1. libcurl 调试工具简介2. 输出请求消息使用 CURLOPT_VERBOSE使用 C

Redis如何使用zset处理排行榜和计数问题

《Redis如何使用zset处理排行榜和计数问题》Redis的ZSET数据结构非常适合处理排行榜和计数问题,它可以在高并发的点赞业务中高效地管理点赞的排名,并且由于ZSET的排序特性,可以轻松实现根据... 目录Redis使用zset处理排行榜和计数业务逻辑ZSET 数据结构优化高并发的点赞操作ZSET 结

微服务架构之使用RabbitMQ进行异步处理方式

《微服务架构之使用RabbitMQ进行异步处理方式》本文介绍了RabbitMQ的基本概念、异步调用处理逻辑、RabbitMQ的基本使用方法以及在SpringBoot项目中使用RabbitMQ解决高并发... 目录一.什么是RabbitMQ?二.异步调用处理逻辑:三.RabbitMQ的基本使用1.安装2.架构

一文详解Python中数据清洗与处理的常用方法

《一文详解Python中数据清洗与处理的常用方法》在数据处理与分析过程中,缺失值、重复值、异常值等问题是常见的挑战,本文总结了多种数据清洗与处理方法,文中的示例代码简洁易懂,有需要的小伙伴可以参考下... 目录缺失值处理重复值处理异常值处理数据类型转换文本清洗数据分组统计数据分箱数据标准化在数据处理与分析过

mysql外键创建不成功/失效如何处理

《mysql外键创建不成功/失效如何处理》文章介绍了在MySQL5.5.40版本中,创建带有外键约束的`stu`和`grade`表时遇到的问题,发现`grade`表的`id`字段没有随着`studen... 当前mysql版本:SELECT VERSION();结果为:5.5.40。在复习mysql外键约