草庐IT

MQ如何保证消息不丢失

胡尚 2023-12-27 原文

如何保证消息不丢失

哪些环节会造成消息丢失

其实主要就是跨网络的环境中需要考虑消息的丢失,主要是有以下几个方面

  • 生产者往MQ发送消息
  • MQ的Broker是集群有主从的,主节点把消息同步到从节点时也需要考虑消息丢失问题
  • 消息从内存持久化到硬盘时,MQ的消息是工作在内存中的,但是内存是断电就丢失数据,所以需要持久化到磁盘,这一步也需要考虑消息丢失问题
  • 消费者消费MQ的消息

如下图所示的四个步骤都有可能造成消息丢失

如何去防止消息丢失

其实也就是针对上面四个环节来分析,保证每个环节的消息不丢失

  • 生产者发送消息不丢失

    • kafka:消息发送+回调。生产者向MQ发送一个消息之后,MQ会向生产者发送一个请求执行相应的回调函数,如果一直没有执行回调函数Produce就知道消息发送失败了就可以重新发送消息

    • RocketMQ:它是Kafka之后出来的,也支持消息发送+回调的机制。同时它还支持事务消息来保证生产者发送消息不丢失

      RocketMQ的事务消息是保证生产者 本地事务和发送消息两步是原子性的:

      1. Producer向MQ发送一个half消息(对于Consumer是不可见的),首先确定MQ是正常运行的

      2. 执行本地事务

      3. 向MQ发送真正的消息,并携带本地事务的执行结果

      4. 如果本地事务执行结果是成功,那么消费者就可以消费此消息;如果本地事务执行结果是失败,那么MQ就会丢弃此消息;如果本地事务执行结果是未知,那么就会经过以下的步骤

      5. 过一段时间MQ向Producer发送一个请求回查本地事务的执行结果

      6. Producer会执行相应的操作去检查当前 本地事务的执行结果,并将结果发送给MQ

      7. MQ再判断本地事务的执行结果,如果还是未知就继续重复执行5~7步,默认会重复执行15次。

    • RabbitMQ:

      也支持消息发送+回调的机制。

      它还支持手动事务机制。

      RabbitMQ的提供了这些API方法,让我们程序自己去实现事务的逻辑

      channel.txSelect()开启事务 channel.txCommit()提交事务 channel.txRollback()回滚事务

      我们首先开启事务,然后执行本地事务,再执行后面两个Api方法。这种手动事务机制有一个问题就是对channel是会产生阻塞的,会造成吞吐量下降

      在RabbitMQ3.*版本开始,它还支持Publisher Confirm机制,相当于是生产者确认机制,整个处理流程和RocketMQ事务消息的处理流程基本是一样的

  • MQ主从同步消息不丢失

    • RocketMQ

      在RocketMQ中对于普通集群,主从数据复制有两种方式:同步复制和异步复制。同步复制就能保证消息不丢失,异步复制效率高但是可能丢消息

      第二种方式就是Dledger集群,它在主从数据复制时采用两阶段提交来保证消息不丢失

      ​ 普通集群就是我们指定哪个节点是Master,哪个节点是Slave;而Dledger集群会频繁的隔一段时间从至少三个节点中选举一个节点成为Master,其余的为Slave,当Producer发送消息到Master后,Master会直接返回给Producer,然后当前消息会标记为UnCommited,再给Slave进行消息同步,当大部分Slave都同步成功后才会把消息的状态改为Commited。这就是这里的两阶段。

    • RabbitMQ

      普通集群:消息是分散存储在各个节点,节点之间不会主动进行消息同步,只有在消费时才会进行消息同步。就比如Producer生产一个消费发送到了A节点,这个时间集群中各个节点是不同步消息的,Consumer却在B节点上消费这条消息,这个时候才会把A节点的消息同步到B节点来 再进行消费。这种方式是可以丢失消息的

      镜像集群:当Producer生产一个消息后,各个节点之间会主动的进行消息的同步,这样数据安全性会更高

    • Kafka

      它通常是在允许少量消息丢失的场景,它可以通过配置acks,配置为0 1 all 。这就相当于RocketMQ的同步复制/异步复制

  • MQ消息持久化存盘时消息不丢失

    • RocketMQ:提供了一种配置的方式,可以选择同步刷盘,也可以选择异步刷盘
    • RabbitMQ:将队列配置成持久化队列,这样就可以保证消息不丢失。在RabbitMQ3.*版本中还有一个Quorum类型的队列,会采用Raft协议来进行消息同步
  • 消费者消费消息时消息不丢失

    一般情况下,MQ的队列中会有一个offset偏移量指向当前消费消息的位置,Consumer消费消息之后会往MQ返回一个消息,然后MQ就会把offset偏移量往前移动。如果consumer消费失败了,那么MQ的offset也不会移动,下次Consume再重新消费就行了。

    会造成消费时消息丢失才场景是:Consumer消费消息变为异步的方式,刚开始收到需要消费的消息就往MQ发送一个消息,然后再去执行本地事务,这个时候如果本地事务执行失败,可是发送给MQ的确认消息却已经发送成功了,这就造成了消费端消息丢失。

    解决消费者消息消息时不丢失的方式:

    • RocketMQ:使用默认的消费方式就行,不要采用异步方式
    • RabbitMQ:它在消费消息时有一个autocommit自动提交的机制,我们将autocommit关闭,不要让它自动提交,改为手动提交offset
    • Kafka:也是一样的改为手动提交offset

有关MQ如何保证消息不丢失的更多相关文章

  1. ruby - 如何使用 Nokogiri 的 xpath 和 at_xpath 方法 - 2

    我正在学习如何使用Nokogiri,根据这段代码我遇到了一些问题:require'rubygems'require'mechanize'post_agent=WWW::Mechanize.newpost_page=post_agent.get('http://www.vbulletin.org/forum/showthread.php?t=230708')puts"\nabsolutepathwithtbodygivesnil"putspost_page.parser.xpath('/html/body/div/div/div/div/div/table/tbody/tr/td/div

  2. ruby - 如何从 ruby​​ 中的字符串运行任意对象方法? - 2

    总的来说,我对ruby​​还比较陌生,我正在为我正在创建的对象编写一些rspec测试用例。许多测试用例都非常基础,我只是想确保正确填充和返回值。我想知道是否有办法使用循环结构来执行此操作。不必为我要测试的每个方法都设置一个assertEquals。例如:describeitem,"TestingtheItem"doit"willhaveanullvaluetostart"doitem=Item.new#HereIcoulddotheitem.name.shouldbe_nil#thenIcoulddoitem.category.shouldbe_nilendend但我想要一些方法来使用

  3. python - 如何使用 Ruby 或 Python 创建一系列高音调和低音调的蜂鸣声? - 2

    关闭。这个问题是opinion-based.它目前不接受答案。想要改进这个问题?更新问题,以便editingthispost可以用事实和引用来回答它.关闭4年前。Improvethisquestion我想在固定时间创建一系列低音和高音调的哔哔声。例如:在150毫秒时发出高音调的蜂鸣声在151毫秒时发出低音调的蜂鸣声200毫秒时发出低音调的蜂鸣声250毫秒的高音调蜂鸣声有没有办法在Ruby或Python中做到这一点?我真的不在乎输出编码是什么(.wav、.mp3、.ogg等等),但我确实想创建一个输出文件。

  4. ruby-on-rails - 如何验证 update_all 是否实际在 Rails 中更新 - 2

    给定这段代码defcreate@upgrades=User.update_all(["role=?","upgraded"],:id=>params[:upgrade])redirect_toadmin_upgrades_path,:notice=>"Successfullyupgradeduser."end我如何在该操作中实际验证它们是否已保存或未重定向到适当的页面和消息? 最佳答案 在Rails3中,update_all不返回任何有意义的信息,除了已更新的记录数(这可能取决于您的DBMS是否返回该信息)。http://ar.ru

  5. ruby-on-rails - 'compass watch' 是如何工作的/它是如何与 rails 一起使用的 - 2

    我在我的项目目录中完成了compasscreate.和compassinitrails。几个问题:我已将我的.sass文件放在public/stylesheets中。这是放置它们的正确位置吗?当我运行compasswatch时,它不会自动编译这些.sass文件。我必须手动指定文件:compasswatchpublic/stylesheets/myfile.sass等。如何让它自动运行?文件ie.css、print.css和screen.css已放在stylesheets/compiled。如何在编译后不让它们重新出现的情况下删除它们?我自己编译的.sass文件编译成compiled/t

  6. ruby - 如何将脚本文件的末尾读取为数据文件(Perl 或任何其他语言) - 2

    我正在寻找执行以下操作的正确语法(在Perl、Shell或Ruby中):#variabletoaccessthedatalinesappendedasafileEND_OF_SCRIPT_MARKERrawdatastartshereanditcontinues. 最佳答案 Perl用__DATA__做这个:#!/usr/bin/perlusestrict;usewarnings;while(){print;}__DATA__Texttoprintgoeshere 关于ruby-如何将脚

  7. ruby - 如何指定 Rack 处理程序 - 2

    Rackup通过Rack的默认处理程序成功运行任何Rack应用程序。例如:classRackAppdefcall(environment)['200',{'Content-Type'=>'text/html'},["Helloworld"]]endendrunRackApp.new但是当最后一行更改为使用Rack的内置CGI处理程序时,rackup给出“NoMethodErrorat/undefinedmethod`call'fornil:NilClass”:Rack::Handler::CGI.runRackApp.newRack的其他内置处理程序也提出了同样的反对意见。例如Rack

  8. ruby - 如何每月在 Heroku 运行一次 Scheduler 插件? - 2

    在选择我想要运行操作的频率时,唯一的选项是“每天”、“每小时”和“每10分钟”。谢谢!我想为我的Rails3.1应用程序运行调度程序。 最佳答案 这不是一个优雅的解决方案,但您可以安排它每天运行,并在实际开始工作之前检查日期是否为当月的第一天。 关于ruby-如何每月在Heroku运行一次Scheduler插件?,我们在StackOverflow上找到一个类似的问题: https://stackoverflow.com/questions/8692687/

  9. ruby-on-rails - 如何从 format.xml 中删除 <hash></hash> - 2

    我有一个对象has_many应呈现为xml的子对象。这不是问题。我的问题是我创建了一个Hash包含此数据,就像解析器需要它一样。但是rails自动将整个文件包含在.........我需要摆脱type="array"和我该如何处理?我没有在文档中找到任何内容。 最佳答案 我遇到了同样的问题;这是我的XML:我在用这个:entries.to_xml将散列数据转换为XML,但这会将条目的数据包装到中所以我修改了:entries.to_xml(root:"Contacts")但这仍然将转换后的XML包装在“联系人”中,将我的XML代码修改为

  10. ruby - 如何使用文字标量样式在 YAML 中转储字符串? - 2

    我有一大串格式化数据(例如JSON),我想使用Psychinruby​​同时保留格式转储到YAML。基本上,我希望JSON使用literalstyle出现在YAML中:---json:|{"page":1,"results":["item","another"],"total_pages":0}但是,当我使用YAML.dump时,它不使用文字样式。我得到这样的东西:---json:!"{\n\"page\":1,\n\"results\":[\n\"item\",\"another\"\n],\n\"total_pages\":0\n}\n"我如何告诉Psych以想要的样式转储标量?解

随机推荐