Day12 RocketMQ的简单使用
今天没学多少,接前一天
- 延迟消息
-
指定延迟级别 message.setDelayTimeLevel(3);

-
指定发送时间 message.setDeliverTimeMs(System.currentTimeMillis() + 10_000L);
-
- 批量消息
- 多条消息合并成一批发出去,减少网络IO,提升消息发送的吞吐量
- 批量消息最好不要超过1M,同一批消息的topic必须相同,且不支持延迟
- 事物消息
-
通过RocketMQ的事物机制,保证上下游数据一致性,重点是监听数据库事物,根据数据库的事物判断MQ是否投递给消费者 示例代码请下载源码查看org.apache.rocketmq.example.transaction.TransactionProducer

-
生产者将消息发送至Apache RocketMQ服务端。
-
Apache RocketMQ服务端将消息持久化成功之后,向生产者返回Ack确认消息已经发送成功,此时消息被标记为"暂不能投递",这种状态下的消息即为半事务消息。
半消息对消费者不可见,实际上是将消息放进来一个叫RMQ_SYS_TRANS_HALF_TOPIC的系统topic
-
生产者开始执行本地事务逻辑。
-
生产者根据本地事务执行结果向服务端提交二次确认结果(Commit或是Rollback),服务端收到确认结果后处理逻辑如下:
- 二次确认结果为Commit:服务端将半事务消息标记为可投递,并投递给消费者。
- 二次确认结果为Rollback:服务端将回滚事务,不会将半事务消息投递给消费者。
-
在断网或者是生产者应用重启的特殊情况下,若服务端未收到发送者提交的二次确认结果,或服务端收到的二次确认结果为Unknown未知状态,经过固定时间后,服务端将对消息生产者即生产者集群中任一生产者实例发起消息回查。
回查次数通过transactionCheckMax参数设置,默认15次,回查间隔通过transactionCheckInterval参数设置,默认60s
-
生产者收到消息回查后,需要检查对应消息的本地事务执行的最终结果。
-
生产者根据检查到的本地事务的最终状态再次提交二次确认,服务端仍按照步骤4对半事务消息进行处理。
-
-
- ACL权限机制
- topic权限 通过perm字段配置 2:禁写禁订阅,4:可订阅,不能写,6:可写可订阅
- broker权限配置
-
在broker.conf 设置aclEnable=true
-
修改plain_acl.yml(热加载,不需要重启)

-
客户端使用,引入org.apache.rocketmq.rocketmq-acl包,声明时传入AclClientRPCHook对象即可
-
- SpringBoot整合RocketMQ使用
▼xml复制代码<dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-spring-boot-starter</artifactId> <version>2.3.0</version> <!-- 推荐使用最新稳定版、注意版本兼容性 --> </dependency>
▼yml复制代码rocketmq: name-server: localhost:9876 # 必填!NameServer 地址 producer: group: my-producer-group # Producer 组名(必须唯一) send-message-timeout: 3000 # 发送超时(毫秒) retry-times-when-send-failed: 2 # 同步发送失败重试次数 # consumer: # 如果需要消费消息,再配置 consumer # listeners: # your-topic: your-consumer-group
▼java复制代码// Producer import org.apache.rocketmq.spring.core.RocketMQTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @Service public class MessageService { @Autowired private RocketMQTemplate rocketMQTemplate; public void sendMessage(String topic, String message) { // 同步发送:topic + 消息体(自动序列化为 JSON) rocketMQTemplate.convertAndSend(topic, message); // 也可以带 tag:topic:tag // rocketMQTemplate.convertAndSend("OrderTopic:pay", order); } } // Consumer import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Component; /** * 参考官方测试类 com.roy.rocketmq.SpringRocketTest * 一个RocketMQTemplate实例只能包含一个生产者,也只能往一个topic下发送消息,如果需要往另一个topic下发消息 * 就要通过@ExtRocketMQTemplateConfiguration()声明一个子实例 * 事物消息使用 @RocketMQTransactionListener,通过rocketMQTemplateBeanName指向具体子类 */ @Component @RocketMQMessageListener( topic = "TestTopic", consumerGroup = "my-consumer-group" // 必须与 producer 的 group 不同 ) public class TestTopicConsumer implements RocketMQListener<String> { @Override public void onMessage(String message) { System.out.println("Received message: " + message); // 处理业务逻辑 } }



