springboot-rocketmq-starter
starter源码github地址:https://github.com/Lickey1991/rocketmq-spring-boot-starter
笔者这里fork https://github.com/maihaoche/rocketmq-spring-boot-starter项目,做了一下小改动,一个项目支持多个producer group,如果不需要此修改,可以check此项目。
springboot-demo
springboot接入demo源码 github地址:https://github.com/Lickey1991/rocketmq-springbooot-demo
rocketmq-demo-springMVC
springMVC接入demo源码 github地址:https://github.com/Lickey1991/rocketmq-springmvc-demo
RocketMq版本 4.5
1.spring Boot 项目接入
1.1 maven 添加依赖
<dependency>
<groupId>com.lickey</groupId>
<artifactId>spring-boot-starter-rocketmq</artifactId>
<version>1.0.0-SNAPSHOT</version>
</dependency>
1.2 添加rocketMq配置
rocketMq NameSvrAddr spring.rocketmq.name-server-address: 10.17.0.171:9876;10.17.0.172:9876
# 可选, 如果无需发送消息则忽略该配置
spring.rocketmq.producer-group: ${rocketmq.groupid.demo}
# 发送超时配置毫秒数, 可选, 默认3000
spring.rocketmq.send-msg-timeout: 5000
# 追溯消息具体消费情况的开关,默认打开
#trace-enabled: false
2. Sping MVC 项目接入
2.1 maven 添加依赖
<dependencies>
<dependency>
<groupId>com.lickey</groupId>
<artifactId>spring-boot-starter-rocketmq</artifactId>
<version>1.0.0-SNAPSHOT</version>
<exclusions>
<exclusion>
<artifactId>slf4j-log4j12</artifactId>
<groupId>org.slf4j</groupId>
</exclusion>
</exclusions>
</dependency>
2.2 添加rocketMq配置
添加rocketmq.properties 文件,配置内容与springboot一致
2.3 springMvc 配置@Value读取properties
- application-context.xml配置
<!--加载 RocketMq starter autoConfiguration-->
<context:component-scan base-package="com.lickey.starter.rocketmq.config"/>
<!-- 加载配置属性文件 -->
<context:property-placeholder ignore-unresolvable="true" location="classpath:rocketmq.properties" />
2.spring-mvc.xml或者dispatcher-servlet.xml 配置
<!-- 加载配置属性文件 -->
<context:property-placeholder ignore-unresolvable="true" location="classpath:rocketmq.properties" />
因为controller配置在spring-mvc.xml/dispatcher-servlet.xml 中,所以想要在controller中使用@value("${key}")读取properties 需要在spring-mvc.xml/dispatcher-servlet.xml 也定义properties配置
RocketMq
配置中 name-sercer-address可以在 rocketMq
的控制台中查看,daily控制台中查看。 多个地址,中间使用分号【;】进行分隔
3. 创建producer
3.1 普通producer
- 创建`producer
参考demo中RocketMqProducer,如果需要在消息发送完成后,做统一处理逻辑例如记录日志等,需要重写doAfterSyncSend方法
使用@MQProducer注解创建普通producer,producerGroup为生产者组,如果此处不指定生产这组,则取配置文件中spring.rocketmq.producer-group值为默认生产者组,二者必须有个值不能为空,否则创建失败,需要继承 AbstractMQProducer 类
2.发送同步消息
参考TestServiceImpl#sendMqMessqge,调用producer的syncSend (Message message)
使用MessageBuilder构建消息体Message,通过静态of方法,指定message的topic与tag,或者指定messageBody信息。同时提供topic,tag,massageBody对相应的set方法。
setKey设置rocketMq消息中的key,消息发送后在rocketMq控制台查看detail信息如下
3.发送异步消息
参考demoTestServiceImpl#asyncSend,调用producer的asyncSend (Message message, SendCallback sendCallback)方法
sendCallBack
为发送后的回调,实现方式参考DemoSendCallback,实现SendCallBack接口,重写onSuccess和onException方法
4.发送顺序消息
调用producer的syncSendOrderly(Message message, String hashKey)方法,该方式使用rocketMq顺序消息的hashKey计算消息顺序方法
3.2创建事物producer
1.创建producer
参考demo中的RocketMqTransactionProducer,使用@MQTransactionProducer注解创建事物producer,producerGroup为生产者组,必须有值不能为空且不能与其他producer组相同,否则创建失败,需要继承 AbstractMQTransactionProducer类,实现executeLocalTransaction、checkLocalTransaction方法(checkLocalTransation方法作用本地事物执行时间过长 或者集群收到producer传过来的状态是unknow,集群通过查询事物消息TOPIC,回调check),该功能在rocketMq 3.6版本中去除。回调的实现通过ClientRemotingProcessor#processRequest。 所以checkLocalTransaction 无需写实际方法体。想要测试,在executeLocalTransaction中返回UNKNOW,processRequest处断点观察。回调触发,有兴趣可以看下相关文章
2.发送事物消息
参考demoTestServiceImpl#transactionSend,调用producer的sendMessageInTransaction(Message msg, Object arg)方法
4.创建消费者
参考demo中DemoTopicConsumer,使用@MQConsumer注解创建事物,consumerGroup
为消费者组,topic为监听的topic,重写process方法,方法参数message为反序列化后的messageBody实例,extMap key包括MessageExtConst中 MessageExt、Message.property两部分
MessageExtConst中extMap包含的key
/** 来自 MessageExt */
public static final String PROPERTY_EXT_QUEUE_ID = "QUEUE_ID";
public static final String PROPERTY_EXT_STORE_SIZE = "STORE_SIZE";
public static final String PROPERTY_EXT_QUEUE_OFFSET = "QUEUE_OFFSET";
public static final String PROPERTY_EXT_SYS_FLAG = "SYS_FLAG";
public static final String PROPERTY_EXT_BORN_TIMESTAMP = "BORN_TIMESTAMP";
public static final String PROPERTY_EXT_BORN_HOST = "BORN_HOST";
public static final String PROPERTY_EXT_STORE_TIMESTAMP = "STORE_TIMESTAMP";
public static final String PROPERTY_EXT_STORE_HOST = "STORE_HOST";
public static final String PROPERTY_EXT_MSG_ID = "MSG_ID";
public static final String PROPERTY_EXT_COMMIT_LOG_OFFSET = "COMMIT_LOG_OFFSET";
public static final String PROPERTY_EXT_RECONSUME_TIMES = "RECONSUME_TIMES";
public static final String PROPERTY_EXT_PREPARED_TRANSACTION_OFFSET = "PREPARED_TRANSACTION_OFFSET";
public static final String PROPERTY_EXT_BODY_CRC = "BODY_CRC";
/** 以下属性来自 Message.property */
public static final String PROPERTY_KEYS = "KEYS";
public static final String PROPERTY_TAGS = "TAGS";
public static final String PROPERTY_WAIT_STORE_MSG_OK = "WAIT";
public static final String PROPERTY_DELAY_TIME_LEVEL = "DELAY";
public static final String PROPERTY_RETRY_TOPIC = "RETRY_TOPIC";
public static final String PROPERTY_REAL_TOPIC = "REAL_TOPIC";
public static final String PROPERTY_REAL_QUEUE_ID = "REAL_QID";
public static final String PROPERTY_TRANSACTION_PREPARED = "TRAN_MSG";
public static final String PROPERTY_PRODUCER_GROUP = "PGROUP";
public static final String PROPERTY_MIN_OFFSET = "MIN_OFFSET";
public static final String PROPERTY_MAX_OFFSET = "MAX_OFFSET";
public static final String PROPERTY_BUYER_ID = "BUYER_ID";
public static final String PROPERTY_ORIGIN_MESSAGE_ID = "ORIGIN_MESSAGE_ID";
public static final String PROPERTY_TRANSFER_FLAG = "TRANSFER_FLAG";
public static final String PROPERTY_CORRECTION_FLAG = "CORRECTION_FLAG";
public static final String PROPERTY_MQ2_FLAG = "MQ2_FLAG";
public static final String PROPERTY_RECONSUME_TIME = "RECONSUME_TIME";
public static final String PROPERTY_MSG_REGION = "MSG_REGION";
public static final String PROPERTY_TRACE_SWITCH = "TRACE_ON";
public static final String PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX = "UNIQ_KEY";
public static final String PROPERTY_MAX_RECONSUME_TIMES = "MAX_RECONSUME_TIMES";
public static final String PROPERTY_CONSUME_START_TIMESTAMP = "CONSUME_START_TIME";
以上是接入starter Demo