RocketMQ 投递消息方式以及消息体结构分析:Message、MessageQueueSelector

news2025/1/16 8:43:19

在这里插入图片描述

🔭 嗨,您好 👋 我是 vnjohn,在互联网企业担任 Java 开发,CSDN 优质创作者
📖 推荐专栏:Spring、MySQL、Nacos、Java,后续其他专栏会持续优化更新迭代
🌲文章所在专栏:RocketMQ
🤔 我当前正在学习微服务领域、云原生领域、消息中间件等架构、原理知识
💬 向我询问任何您想要的东西,ID:vnjohn
🔥觉得博主文章写的还 OK,能够帮助到您的,感谢三连支持博客🙏
😄 代词: vnjohn
⚡ 有趣的事实:音乐、跑步、电影、游戏

目录

  • 前言
  • Message
    • Properties
    • Topic、Tag
    • Keys
    • Queue
  • SendCallback
  • MessageQueue
  • MessageQueueSelector
  • FAQ
  • 总结

前言

在踏入投递消息 send 方法源码解读开始前,要先搞清楚消息的结构以及消息的不同的方式实现,以便于在阅读底层源码时不至于要重头过来再梳理各自的类、属性结构起到了什么作用,在这里再啰嗦啰嗦,让 RocketMQ 源码能够更加完整.

RocketMQ 专栏篇:

从零开始:手把手搭建 RocketMQ 单节点、集群节点实例

保护数据完整性:探索 RocketMQ 分布式事务消息的力量

RocketMQ 分布式事务消息实战指南:确保数据一致性的关键设计

RocketMQ 生产者源码分析:DefaultMQProducer、DefaultMQProducerImpl

RocketMQ MQClientInstance、生产者实例启动源码分析

Message

RocketMQ 消息构成的结构如下:

  1. Topic:表示要发送消息的主题
  2. body:表示消息的存储内容
  3. Properties:表示消息属性
  4. transactionId:在事务消息中使用

Keys、Tag:无论是 RocketMQ Tag 过滤还是延迟消息等都会利用 Properties 消息属性机制,这些特殊信息都使用了系统保留的属性 Key,设置自定义属性时要避免和系统属性 Key 相冲突

服务端会根据 Keys 创建哈希索引,设置以后,可以在 Console 系统通过 Topic、Keys 来查询消息,由于是哈希索引,尽可能保证 Key 唯一,例如:订单号、会员 Id 等

RocketMQ 系统保留的属性 Key 集合,在使用消息属性机制时尽量避免:org.apache.rocketmq.common.message.MessageConst

Properties

在 Message 结构体中可以设置的属性如下:

字段名默认值必要性说明
Topicnull必填消息所属 Topic 名称
Bodynull必填消息体
Tagsnull选填消息标签:方便服务器过滤消息使用
消息属性:TAGS
Keysnull选填代表这条消息的业务关键词
消息属性:KEYS
Flag0选填由应用程序设置,RocketMQ 不作干预
DelayTimeLevel0选填消息延迟级别,0 表示不延时,
大于 0 会延迟特定的时间才会被消息
消息属性:DELAY
WaitStoreMsgOKtrue选填表示消息是否在服务器落盘后才返回 SEND_OK
消息属性:WAIT

Topic、Tag

Topic、Tag 都是用来在业务上分类的标识,区别在于 Topic 是一级分类,Tag 理解为是二级分类,使用 Tag 可以实现对 Topic 中的消息进行过滤

Topic:消息主题,通过 Topic 对不同的业务消息进行分类

Tag:消息标签,用来进一步区分某个 Topic 下的消息分类,消息从生产者发出时就带上的属性

Tag 属性可以为 Topic 做另外一层业务上再细粒化地区分子业务

在这里插入图片描述

什么时候该用 Topic,什么该用 Tag?

一般可以从以下几方面进行判别:

  1. 消息类型是否一致:普通消息、事务消息、定时(延时)消息、顺序消息,不同的消息类型使用不同的 Topic,无法通过 Tag 进行区分

    除了普通消息以外,其余的消息类型应该使用对应的后缀拼接,如:事务 _transaction、延迟 _delay、顺序 _order

  2. 业务是否相关联:没有直接关联的消息,如淘宝交易消息与京东物流消息使用不同的 Topic 区分;若同样的是京东物流消息,电器类订单、男装类订单、鞋类订单的消息可以使用 Tag 进行划分

  3. 消息优先级是否一致:如同样是物流消息,盒马必须小时内送达,天猫超市 24 小时内送达,淘宝物流则相对会慢一些,不同优先级的消息用不同的 Topic 进行区分

  4. 消息量级是否相当:有些业务消息虽然量小但是实时性要求高,如果跟某些万亿量级的消息使用同一个 Topic,则有可能会因为过长的等待时间而“饿死”,此时需要将不同量级的消息进行拆分,使用不同的 Topic

总的来说,针对消息分类,可以选择创建多个 Topic,或者在同一个 Topic 下创建多个 Tag;通常情况下,不同的 Topic 之间消息没有必然的联系,而 Tag 则用来区分同一个 Topic 下相互关联的消息,例如:全集、子集的关系,流程先后的关系.

Keys

RocketMQ 每个消息可以在业务层面上设置唯一标识码 Keys 字段,方便将来定位消息丢失的问题,Broker 会为每个消息创建哈希索引,应用可以通过 Topic、Key 来查询这条消息的内容,以及消息被谁消费,由于是哈希索引,请务必保证 Key 尽可能唯一,这样可以避免潜在的哈希冲突

// 订单Id
String orderId = "2024016163739";
message.setKeys(orderId);

Queue

为了支持高并发以及水平扩展,需要对 Topic 进行分区「类比于 Kafka Partition 机制」在 RocketMQ 中被称之为队列,一个 Topic 下可能有多个队列,并且可能分布在不同的 Broker 上

在这里插入图片描述

一般来说一条消息,若没有重复发送(比如:因为客户端没有响应而进行重试)则只会存在 Topic 其中的一个队列,消息在队列中按照先进先出的原则存储,每条消息会有自己的位点,每个队列会统计当前消息的总条数,这个称为最大位点:MaxOffset;队列的起始位置对应的位置叫做起始位点 MinOffset,队列可以提升消息发送和消费的并发度.

在这里谈到了 Topic、Queue 之间的关联,也就有必要说说 RocketMQ 中的 AKF 原则了

在这里插入图片描述

  • X 轴:处理的服务节点的单点故障问题,支持横向扩展、全量镜像
  • Y 轴:在 RocketMQ 服务节点基础上根据业务来划分出不同的 Topic
  • Z 轴:基于 Topic 分配出不同的 Queue 分区,每个 Queue 可能分散到不同的 Broker 服务节点上

SendCallback

当 send 方法 SendCallback 参数不为空,说明它是属于异步发送的方式,异步发送:发送方发出一条消息以后,不等待服务端返回响应,会接着发送下一条消息

SendCallback 是基于异步发送时需要的指定的回调接口,它提供了两个方法:

public interface SendCallback {
    void onSuccess(final SendResult sendResult);

    void onException(final Throwable e);
}

消息发送方在发送了一条消息后,不需要等待服务端响应即可发送第二条消息,发送方通过回调接口接收服务端响应,并处理响应结果。异步发送一般用于链路耗时较长,对响应时间较为敏感的业务场景。例如:视频上传后通知启动转码服务,转码完成后通知推送转码结果等

同步发送示例代码:

DefaultMQProducer mqProducer = new DefaultMQProducer();
mqProducer.setProducerGroup("vnjohn");
mqProducer.setNamesrvAddr("172.16.249.10:9876");
mqProducer.start();
Message message = new Message();
message.setTopic("vnjohn");
message.setWaitStoreMsgOK(Boolean.FALSE);
message.setKeys("2024016163739");
message.setBody("Hello RocketMQ".getBytes(StandardCharsets.UTF_8));
message.setTags("blog");
mqProducer.send(message);

异步发送示例代码:

mqProducer.send(message, new SendCallback() {
  @Override
  public void onSuccess(SendResult sendResult) {
    System.out.println("message send success...");
  }

  @Override
  public void onException(Throwable e) {
    System.out.println("message send exception...");
  }
});

异步发送、同步发送代码唯一区别在于调用 send 接口的参数不同,异步发送不会等待发送返回,取而代之的是 send 方法需要传入 SendCallback 的实现,SendCallback 接口主要有 onSuccess、onException 两个方法,表示消息发送成功和消息发送失败

MessageQueue

当 send 方法时,指定了 MessageQueue 参数,说明要将消息指定投递给 Topic 其中的一个队列:Queue

// MessageQueue:指定 Topic 主题、Broker、Topic-Queue
private String topic;
private String brokerName;
private int queueId;

比如我要将消息投递给 「Topic:vnjohn」 下的 「Queue:0」所属的 「brokerName:broker-0」

mqProducer.send(message, new MessageQueue("vnjohn", "broker-0", 0));

MessageQueueSelector

在生产者投递消息时,通过该类可以指定将需要保证顺序消费的消息「通过 args 来指定」放入到同一个队列中,确保在消费时可以通过消费这一个队列保证消费的顺序性

生产顺序性:

RocketMQ 通过生产者和服务端的协议保障 单个生产者以串行的方式 发送消息,并按序存储和持久化,如需保证消息生产的有序性,则必须满足以下条件:

  1. 单一生产者:消息的生产的顺序仅支持单一生产者,不同生产者分布在不同的系统,即使设置为相同的分区键-arg,不同生产者之间产生的消息也无法判定其顺序
  2. 串行发送:生产者客户端支持多线程安全访问,但如果生产者使用了多线程并行发送,则不同线程间产生的消息将无法判定其先后顺序

满足以上条件的生产者,将顺序消息发送至服务端以后,会保证设置了同一分区键的消息,按照发送的顺序存储在同一个队列中

在 RocketMQ 中,Topic 下只存在一个 Queue 时,生产者可以保证 全局有序性,否则生产者就只能保证 局部有序性

对于一个指定的 Topic,消息严格按照先进先出(FIFO)的原则进行消息发布和消费,即先发布的消息先消费,后发布的消息后消费。在 RocketMQ 中支持分区顺序消息,如下图所示:我们可以按照某一个标准对消息进行分区(比如图中的 ShardingKey),同一个 ShardingKey 消息会被分配到同一个队列中,并按照顺序被消费

在这里插入图片描述

顺序消息的应用场景非常广泛,在有序事件处理、撮合交易、数据实时增量同步等场景下,异构系统之间需要维持强一致的状态同步,上游的事件变更需要按照顺序传递到下游进行处理

例如:创建订单的场景,需要保证同一个订单的生成、付款和发货,这三个操作被顺序执行;如果是普通消息,订单-A 的消息可能会被轮询发送到不同的队列中,不同队列的消息将无法保证顺序,而顺序消息发送时将 ShardingKey 相同(同一个订单号)的消息路由到同一个逻辑队列中

提到的 SharingKey 指的就是在往服务端有序发送消息时,必须指定的一个参数:Object arg,假设它就是订单 id

在这里插入图片描述

MessageQueueSelector 是一个接口,它的定义如下:

public interface MessageQueueSelector {
    MessageQueue select(final List<MessageQueue> mqs, final Message msg, final Object arg);
}

其中 mqs 是可以发送的队列,msg 是消息,arg 是 send 方法时传递的 Object 对象,返回的是该消息需要发送到的队列

通过 MessageQueueSelector 顺序发送消息的示例代码如下:

MessageQueueSelector customMessageQueueSelector = (mqs, msg, arg) -> {
  int shardingKey = arg.hashCode();
  int messageQueueIndex = mqs.size() % shardingKey;
  return mqs.get(messageQueueIndex);
};
mqProducer.send(message, customMessageQueueSelector, "20240107015444");

生产环境中建议选择最细粒度的分区键进行拆分,例如:将订单 ID、用户 ID 作为分区键关键字,可实现同一个终端用户的消息按照顺序处理,不同用户的消息无须保证顺序.

FAQ

若一个 Broker 掉线,队列总数发送变化会如何?

若发送变化,那么同一个 ShardingKey 的消息就会发送到不同的队列上,造成乱序;若不发生变化,那消息将会发送到掉线 Broker 队列上,必然是失败的;因此 RocketMQ 提供了两种模式:严格顺序、顺序可用性,如果要保证严格顺序而不是可用性,创建 Topic 时要指定 -o (-order)参数,表示顺序消息主题

$ sh bin/mqadmin updateTopic -c DefaultCluster -t TopicTest -o true -n 127.0.0.1:9876
create topic to 127.0.0.1:10911 success.
TopicConfig [topicName=TopicTest, readQueueNums=8, writeQueueNums=8, perm=RW-, topicFilterType=SINGLE_TAG, topicSysFlag=0, order=true, attributes=null]

其次就是要保证 nameserver 中的配置参数:orderMessageEnablereturnOrderTopicConfigToBroker 必须是 true,如果上述任意一个条件不满足,则是保证可用性而不是严格顺序

总结

该篇博文主要介绍 Message 结构体有哪些以及它里面属性的作用,同步与异步发送之间的区别参数:SendCallback 回调接口,指定 Queue 投递消息:MessageQueue,生产者侧确保消息能够有序地投递:MessageQueueSelector,以及在 RocketMQ 如何保证全局有序、局部有序:生产者确保单 Queue 全局有序、生产者确保多 Queue 局部有序,希望该篇博文你能够喜欢,感谢三连支持❤️

学习指南针:

Rocket 官方文档

RocketMQ GitHub

🌟🌟🌟愿你我都能够在寒冬中相互取暖,互相成长,只有不断积累、沉淀自己,后面有机会自然能破冰而行!

博文放在 RocketMQ 专栏里,欢迎订阅,会持续更新!

如果觉得博文不错,关注我 vnjohn,后续会有更多实战、源码、架构干货分享!

推荐专栏:Spring、MySQL,订阅一波不再迷路

大家的「关注❤️ + 点赞👍 + 收藏⭐」就是我创作的最大动力!谢谢大家的支持,我们下文见!

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.coloradmin.cn/o/1362837.html

如若内容造成侵权/违法违规/事实不符,请联系多彩编程网进行投诉反馈,一经查实,立即删除!

相关文章

条件竞争之文件上传

一、条件竞争介绍 条件竞争,在程序员日常的Web应用开发中&#xff0c;通常不如其他漏洞受到的关注度高。因为普遍的共识是&#xff0c;条件竞争是不可靠的&#xff0c;大多数时候只能靠代码审计来识别发现&#xff0c;而依赖现有的工具或技术很难在黑盒灰盒中识别并进行攻击。…

精进单元测试技能 —— Pytest断言的艺术!

本篇文章主要是阐述Pytest在断言方面的应用。让大家能够了解和掌握Pytest针对断言设计了多种功能以适应在不同测试场景上使用。 了解断言的基础 在Pytest中&#xff0c;断言是通过 assert 语句来实现的。简单的断言通常用于验证预期值和实际值是否相等&#xff0c;例如&#…

小游戏实战丨基于PyGame的消消乐小游戏

文章目录 写在前面PyGame消消乐注意事项系列文章写在后面 写在前面 本期内容&#xff1a;基于pygame实现喜羊羊与灰太狼版消消乐小游戏 下载地址&#xff1a;https://download.csdn.net/download/m0_68111267/88700193 实验环境 python3.11及以上pycharmpygame 安装pygame…

ENVI 各版本安装指南

ENVI下载链接 https://pan.baidu.com/s/1APpjHHSsrXMaCcJUQGmFBA?pwd0531 1.鼠标右击【ENVI 5.6(64bit&#xff09;】压缩包&#xff08;win11及以上系统需先点击“显示更多选项”&#xff09;选择【解压到 ENVI 5.6(64bit&#xff09;】。 2.打开解压后的文件夹&#xff0c…

一些想法:关于行人检测与重识别OIMLoss

本文主要是介绍我们录用于 ECCV18 的一个工作&#xff1a;Person Search via A Mask-guided Two-stream CNN Model. 这篇文章着眼于 Person Search 这个任务&#xff0c;即同时考虑行人检测&#xff08;Pedestrian Detection&#xff09;与行人重识别&#xff08;Person Re-ide…

2024PMP考试新考纲-【人员领域】近期典型真题和超详细解析(5)

今天华研荟继续为您分享PMP新考纲下的【人员People领域】近年真题&#xff0c;帮助大家举一反三&#xff0c;一次性通过2024年的PMP考试。 2024年PMP考试新考纲-【人员领域】真题解析21 题&#xff1a;项目经理正在为一个项目工作。该项目由于人员流动&#xff0c;相关方登记册…

Matlab二维绘图

低级绘图命令line 有什么点就点哪里&#xff0c;然后连起来&#xff0c;没什么细节&#xff0c;不光滑&#xff0c;所以基本不会用到。 x0:0.2*pi:2*pi; ysin(x); line(x,y);%画一条sin函数线 line([-5,5],[2,2]);%画一条水平线 line([5,5],[0,2]);%画一条竖线 高级绘图命令…

MySQL之视图外连接、内连接和子查询的使用

一、视图 1.1 含义 虚拟表&#xff0c;和普通表一样使用 1.2 操作 创建视图 create view 视图名 as 修改视图 方式一&#xff1a; create or replace view 视图名 as 【查看视图相关字段】 方式二&#xff1a; alter view 视图名 as 【查看的SQL语句】 查看视图 方式一&…

BUUCTF--actf_2019_babyheap1

这题看名字就知道是堆题&#xff0c;先看保护&#xff1a; 保护除了PIE全开&#xff0c;黑盒测试&#xff1a; 题目提供增删查&#xff0c;没有改。看看IDA中代码逻辑&#xff1a; 逻辑跟我前面做的一题极为相似&#xff0c;就不过多分析。 这是free&#xff1a; 因为题目不能…

基于SSM的驾校考试预约管理设计与实现

末尾获取源码 开发语言&#xff1a;Java Java开发工具&#xff1a;JDK1.8 后端框架&#xff1a;SSM 前端&#xff1a;采用JSP技术开发 数据库&#xff1a;MySQL5.7和Navicat管理工具结合 服务器&#xff1a;Tomcat8.5 开发软件&#xff1a;IDEA / Eclipse 是否Maven项目&#x…

求两个数之间的最小公约数

目录 前言 方法&#xff1a;求两个数之间的最小公约数 1.欧几里得算法 2.枚举法 3.公共因子积 4.更相减损术 5.Stein算法 解题&#xff1a;在链表中插入最大公约数 总结 前言 今天刷每日一题&#xff1a;2807. 在链表中插入最大公约数 - 力扣&#xff08;LeetCode&#xff09;…

数据库攻防学习之MySQL

MySQL 0x01mysql学习 MySQL 是瑞典的MySQL AB公司开发的一个可用于各种流行操作系统平台的关系数据库系统&#xff0c;它具有客户机/服务器体系结构的分布式数据库管理系统。可以免费使用使用&#xff0c;用的人数很多。 0x02环境搭建 这里演示用&#xff0c;phpstudy搭建的…

SpringBoot 如何 返回页面

背景 RestController ResponseBody Controller Controller中的方法无法返回jsp页面&#xff0c;或者html&#xff0c;配置的视图解析器 InternalResourceViewResolver不起作用&#xff0c;返回的内容就是Return 里的内容。 Mapping ResponseBody 也会出现同样的问题。 解…

2024年阿里云优惠活动清单_优惠代金券领取大全

阿里云服务器优惠活动大全包括&#xff1a;云服务器新人特惠、云小站、阿里云免费中心、学生主机优惠、云服务器精选特惠、阿里云领券中心等&#xff0c;活动上阿里云服务器ECS经济型e实例2核2G、3M固定带宽99元一年、轻量应用服务器2核2G3M带宽轻量服务器一年61元&#xff0c;…

Linux系统性能优化:七个实战经验

Linux系统的性能是指操作系统完成任务的有效性、稳定性和响应速度。Linux系统管理员可能经常会遇到系统不稳定、响应速度慢等问题&#xff0c;例如在Linux上搭建了一个web服务&#xff0c;经常出现网页无法打开、打开速度慢等现象&#xff0c;而遇到这些问题&#xff0c;就有人…

Spring Cloud Sleuth+zipkin实现链路追踪

Spring Cloud Sleuth 主要功能就是在分布式系统中提供追踪解决方案&#xff0c;并且兼容支持了 zipkin&#xff0c;你只需要在pom文件中引入相应的依赖即可。 微服务架构上通过业务来划分服务的&#xff0c;通过REST调用&#xff0c;对外暴露的一个接口&#xff0c;可能需要很…

【数据库系统概论】数据库并发控制机制——并发控制的主要技术之封锁(Locking)

系统文章目录 数据库的四个基本概念&#xff1a;数据、数据库、数据库管理系统和数据库系统 数据库系统的三级模式和二级映射 数据库系统外部的体系结构 数据模型 关系数据库中的关系操作 SQL是什么&#xff1f;它有什么特点&#xff1f; 数据定义之基本表的定义/创建、修改和…

【Mars3d】new mars3d.layer.GeoJsonLayer({不规则polygon加载label不在正中间的解决方案

问题&#xff1a; 1.new mars3d.layer.GeoJsonLayer({type: "polygon",在styleOptions里配置label的时候&#xff0c;发现这个 不规则polygon加载的时候&#xff0c;会出现label不在中心位置。 graphicLayer new mars3d.layer.GeoJsonLayer({ name: "全国省界…

网络通信过程的一些基础问题

客户端A在和服务器进行TCP/IP通信时&#xff0c;发送和接收数据使用的是同一个端口吗&#xff1f; 这个问题可以这样来思考&#xff1a;在客户端A与服务器B建立连接时&#xff0c;A需要指定一个端口a向服务器发送数据。当服务器接收到A的报文时&#xff0c;从报文头部解析出A的…

电脑开启虚拟化如何查看自己的主机主板型号

问题描述 在使用virtualbox、vmware安装虚拟机的时候&#xff0c;需要本机电脑能够支持虚拟化。 但是不同厂家的主机&#xff08;主板&#xff09;幸好并不一致&#xff0c;所以需要先了解自己的电脑主板型号 操作方法 1、win r 键打开运行窗口&#xff0c;输入cmd并确定打开…