《RocketMQ 官网》阅读笔记 RocketMQ 消息队列 MessageQueue 消息 Messagege
消息队列 (Message Queue)
定义
队列是 Apache RocketMQ 中用于存储和传输消息的容器。队列是 Apache RocketMQ 消息存储的最小单位。Apache RocketMQ 中的一个主题由多个队列组成。
队列具有以下优势:
- 有序存储:队列天生具有顺序性,消息按入队顺序存储。最早的消息位于队列起始位置,最新的消息位于队尾。通过偏移量(Offset)来标记队列中消息的位置和顺序。
- 流式操作语义:Apache RocketMQ 的基于队列的存储允许消费者从特定偏移量读取一条或多条消息。这有助于实现聚合读取和回溯读取等功能。这些功能在 RabbitMQ 或 ActiveMQ 中是不具备的。
模型关系
默认情况下,Apache RocketMQ 提供可靠的消息存储。所有成功发送的消息都会持久化存储在队列中。消息由生产者发送,由消费者客户端接收。每条消息均可保证至少被成功投递一次。
Apache RocketMQ 的队列模型类似于 Kafka 的分区模型。在 Apache RocketMQ 中,队列是主题的一部分。尽管消息由主题管理,但操作是在队列级别进行的。例如,当生产者向特定主题发送消息时,消息会被发送到该主题下的一个队列中。
您可以通过调整 Apache RocketMQ 中的队列数量来进行扩容或缩容。
内部属性
读写权限
定义:决定当前队列是否可以进行读写操作。
取值:由 Broker 定义。枚举值说明如下:
- 6:读写。当前队列既可以写入消息也可以读取消息。
- 4:只读。当前队列可以读取消息,但不能写入消息。
- 2:只写。当前队列可以写入消息,但不能读取消息。
- 0:读写不可用。当前队列不允许任何读写操作。
约束:读写权限与运维操作相关。建议不要频繁修改权限。
行为约束
每个主题包含一个或多个用于存储消息的队列。每个主题的队列数量与消息类型及实例所在的地域有关。队列数量不可随意更改。
版本兼容性
队列名称因 Apache RocketMQ Broker 的版本而异。差异如下:
- Broker 3.x 和 4.x 版本:队列名称由主题名称、Broker ID 和队列 ID 组成,并绑定到物理节点。
- Broker 5.x 版本:队列名称是一个由集群分配的全局唯一字符串,与物理节点解耦。
建议不要手动构建队列名称或将其绑定到其他操作中,否则在 Broker 更新时,这些队列名称可能无法解析。
使用说明
队列数量设置
您可以在创建或修改主题时指定 Apache RocketMQ 中的队列数量。建议配置较少的队列,避免添加不必要的队列。
主题中包含过多队列会导致以下问题:
- 集群元 数据量增加:Apache RocketMQ 基于队列采集指标和监控数据。过多的队列可能导致元数据量剧增。
- 客户端负载过高:Apache RocketMQ 的消息读写基于队列执行。大量的队列可能会产生空轮询请求,从而增加系统负载。
扩容队列的场景
- 物理节点负载均衡
Apache RocketMQ 中每个主题的队列可以分布在不同的服务节点上。为确保集群扩容后的流量负载均衡,建议增加队列或将之前的队列迁移到新的服务节点上。 - 与 FIFO(顺序)消息相关的性能瓶颈
在 Apache RocketMQ Broker 4.x 版本中,FIFO 消息仅在队列级别生效。因此,FIFO 消息的并发性取决于队列数量。当系统出现性能瓶颈时,建议增加队列数量。
消息 (Message)
定义
消息是 Apache RocketMQ 中数据传输的最小单元。生产者将业务数据封装成消息,并将消息发送到 Apache RocketMQ Broker。随后,Broker 会根据相关语义将消息分发给 费者。
Apache RocketMQ 中消息模型的特征为:
- 不可变性:消息一旦产生即成为一个事件。在消息生成后,其内容不会发生改变。即使消息经过传输通道,其内容也保持不变。消费者获取的消息均为只读消息。
- 持久性:默认情况下,Apache RocketMQ 会对消息进行持久化处理。接收到的消息存储在 Apache RocketMQ Broker 的存储文件中,以确保在系统故障时可以对消息进行追踪和恢复。
模型关系
消息类型
- 普通消息:普通消息。普通消息不需要特殊语义,且与其他普通消息不相关联。
- FIFO 消息:顺序消息。Apache RocketMQ 使用消息组来确定一组特定消息的顺序。消息将按照发送顺序进行投递。
- 延迟消息:延迟消息。您可以指定延迟时间,使消息在延迟时间过后才对消费者可见,而不是在生产时立即投递。
- 事务消息:事务消息。Apache RocketMQ 支持分布式事务消息,确保 数据库更新和消息调用之间的事务一致性。
消息队列
消息所属的队列。由 Broker 指定并填充。
消息位点 (Message Offset)
当前消息在队列中的存储位置。由 Broker 指定并填充。有效值:0 到 Long.Max。
消息 ID
消息的唯一标识符。每条消息的 ID 在集群中全局唯一。由生产者客户端自动生成。消息 ID 是一个由 32 个字符组成的字符串,包含数字和大写字母。
消息 Key(可选)
消息的索引键列表。您可以配置不同的 Key 来区分消息并快速查找消息。由生产者客户端定义。
消息 Tag(可选)
用于过滤消息的标签。 消费者可以通过标签过滤消息,仅接收包含指定标签的消息。由生产者客户端定义。每条消息只能指定一个标签。
定时时间(可选)
在定时场景中,消息触发延迟投递时所使用的毫秒级时间戳。由消息生产者定义。最长持续时间为 40 天。
消息发送时间
消息发送时,生产者客户端的本地毫秒级时间戳。客户端时间可能与 Broker 时间不同。在这种情况下,消息发送时间以客户端时间为准。
消息存储时间
消息存储时,Apache RocketMQ Broker 的本地毫秒级时间戳。对于延迟消息和事务消息,消息保留时间是指消息生效时展示给消费者的 Broker 时间。客户端时间可能与 Broker 时间不同。在这种情况下,消息保留时间以 Broker 时间为准。
重试次数
消息消费失败后,Apache RocketMQ Broker 重新投递消息的次数。每次重试后,最大重试次数加一。由 Broker 标记。第一次消费消息时,重试次数为零。第一次消费失败时,重试次数为一。
消息的自定义属性
自定义属性
生产者可以指定的扩展信息。由生产者基于字符串键值对进行指定。
消息负载 (Message Payload)
业务消息的实际数据。由生产者序列化并以二进制字节流进行传输。
行为约束
消息大小不能超过上限。如果消息大小超过相应的上限,消息将发送失败。默认限制:最大消息大小:4 MB。
使用说明
不建议单条消息进行超大负载传输。
Apache RocketMQ 是一种用于传输业务事件数据的消息中间件。如果消息过大,网络传输层可能会过载。这会影响错误重试和限流机制。建议限制单条消息事件的数据大小。如果生产环境确实需要传输大数据,建议根据固定大小对消息进行拆分,或采用文件存储方式。
消息的不可变性
在 Apache RocketMQ Broker 5.x 版本中,消息无法被修改, 消费者获取的消息均为只读消息。3.x 和 4.x 版本未实施有关不可变性的强制约束。建议若需传输消息,请重新初始化消息。
- 正确示例
Messagem=Consumer.receive();Messagem2=MessageBuilder.buildFrom(m);Producer.send(m2);- 错误示例:
Messagem=Consumer.receive();m.update();Producer.send(m);