-
Notifications
You must be signed in to change notification settings - Fork 615
rocketmq spring boot starter使用指南
javahongxi edited this page Dec 28, 2018
·
7 revisions
针对官方starter修改点 官方
- 支持连接多个集群(订阅) (官方一个应用只能连接一个集群)
- 顺序消息消费失败,可配重试次数 (非顺序消息默认重试16次,每次时间延后)
- 发送延时消息方法参数优化(魔法参数改为枚举)
- 优化getMessageType方法,支持 MyConsumer extends AbstractConsumer implements RocketMQListener
(官方只支持MyConsumer implements RocketMQListener) - RocketMQTemplate方法重载(加入keys)
- 暂未加入事务消息功能 (官方最新版支持)
<dependency>
<groupId>org.hongxi</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>1.0.0-SNAPSHOT</version>
</dependency>
spring:
rocketmq:
nameServer: 127.0.0.1:9876
producer:
group: boot_producer_group
@SpringBootApplication
public class ProducerApplication implements CommandLineRunner {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public static void main(String[] args){
SpringApplication.run(ProducerApplication.class, args);
}
public void run(String... args) throws Exception {
// 如下两种方式等价
rocketMQTemplate.convertAndSend("test-topic-1", "Hello, World!");
rocketMQTemplate.send("test-topic-1", MessageBuilder.withPayload("Hello, World! I'm from spring message").build());
// 第三个参数为key
rocketMQTemplate.syncSend("test-topic-1", "Hello, World! I'm from simple message", "18122811143034568830");
// topic: arkOrder,tag: paid, cacel
rocketMQTemplate.convertAndSend("arkOrder:paid", "Hello, World!");
rocketMQTemplate.convertAndSend("arkOrder:cancel", "Hello, World!");
// 消息体为自定义对象
rocketMQTemplate.convertAndSend("test-topic-2", new OrderPaidEvent("T_001", new BigDecimal("88.00")));
// 发送延迟消息
rocketMQTemplate.sendDelayed("test-topic-1", "I'm delayed message", MessageDelayLevel.TIME_1M);
// 发送即发即失消息(不关心发送结果)
rocketMQTemplate.sendOneWay("test-topic-1", MessageBuilder.withPayload("I'm one way message").build());
// 发送顺序消息
rocketMQTemplate.syncSendOrderly("test-topic-4", "I'm order message", "1234");
// 发送异步消息
rocketMQTemplate.asyncSend("test-topic-1", MessageBuilder.withPayload("I'm one way message").build(), new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
}
@Override
public void onException(Throwable e) {
}
});
System.out.println("send finished!");
}
}
@SpringBootApplication
public class ConsumerApplication{
public static void main(String[] args){
SpringApplication.run(ConsumerApplication.class, args);
}
}
@Slf4j
@Service
@RocketMQMessageListener(topic = "test-topic-1", consumerGroup = "my-consumer_test-topic-1")
public class MyConsumer implements RocketMQListener<String> {
public void onMessage(String message) {
log.info("received message: " + message);
}
}
@Slf4j
@Service
@RocketMQMessageListener(topic = "test-topic-2", consumerGroup = "my-consumer_test-topic-2")
public class MyConsumer2 implements RocketMQListener<OrderPaidEvent> {
public void onMessage(OrderPaidEvent orderPaidEvent) {
log.info("received orderPaidEvent: " + orderPaidEvent);
}
}
/**
* 指定连接某个MQ集群
*/
@Slf4j
@Service
@RocketMQMessageListener(nameServer = "127.0.0.1:9877", instanceName = "tradeCluster", topic = "test-topic-3", consumerGroup = "my-consumer_test-topic-3")
public class MyConsumer3 implements RocketMQListener<String> {
public void onMessage(String message) {
log.info("received message: " + message);
}
}
/**
* 顺序消息消费失败,默认不重试(官方是一直重试)
*/
@Slf4j
@Service
@RocketMQMessageListener(topic = "test-topic-4", consumerGroup = "my-consumer_test-topic-4",
consumeMode = ConsumeMode.ORDERLY)
public class MyConsumer4 implements RocketMQListener<String> {
public void onMessage(String message) {
log.info("received message: " + message);
int a = 1 / 0;
}
}
/**
* 配置重试次数 reconsumeTimes = 3
*/
@Slf4j
@Service
@RocketMQMessageListener(topic = "test-topic-4", consumerGroup = "my-consumer_test-topic-5",
consumeMode = ConsumeMode.ORDERLY, reconsumeTimes = 3)
public class MyConsumer5 implements RocketMQListener<MessageExt> {
public void onMessage(MessageExt messageExt) {
log.info("received message: " + messageExt);
int a = 1 / 0;
}
}
/**
* 配置重试次数 reconsumeTimes = -1 代表一直重试
*/
@Slf4j
@Service
@RocketMQMessageListener(topic = "test-topic-4", consumerGroup = "my-consumer_test-topic-6",
consumeMode = ConsumeMode.ORDERLY, reconsumeTimes = -1)
public class MyConsumer6 implements RocketMQListener<MessageExt> {
public void onMessage(MessageExt messageExt) {
log.info("received message: " + messageExt);
int a = 1 / 0;
}
}
注意:Consumer里抛出异常才会重试,所以使用者不要把Consumer里的整个代码try-catch
wiki.hongxi.org
首页
Java核心技术
- JUC JMM与线程安全
- JUC 指令重排与内存屏障
- JUC Java内存模型FAQ
- JUC 同步和Java内存模型
- JUC volatile实现原理
- JUC AQS详解
- JUC AQS理解
- JUC synchronized优化
- JUC 线程和同步
- JUC 线程状态
- JUC 线程通信
- JUC ThreadLocal介绍及原理
- JUC 死锁及避免方案
- JUC 读写锁简单实现
- JUC 信号量
- JUC 阻塞队列
- NIO Overview
- NIO Channel
- NIO Buffer
- NIO Scatter与Gather
- NIO Channel to Channel Transfers
- NIO Selector
- NIO FileChannel
- NIO SocketChannel
- NIO ServerSocketChannel
- NIO Non-blocking Server
- NIO DatagramChannel
- NIO Pipe
- NIO NIO vs. IO
- NIO DirectBuffer
- NIO zero-copy
- NIO Source Code
- NIO HTTP Protocol
- NIO epoll bug
- Reflection 基础
- Reflection 动态代理
- JVM相关
- 设计模式典型案例
Netty
RocketMQ深入研究
kafka深入研究
Pulsar深入研究
Dubbo源码导读
- Dubbo SPI
- Dubbo 自适应拓展机制
- Dubbo 服务导出
- Dubbo 服务引用
- Dubbo 服务字典
- Dubbo 服务路由
- Dubbo 集群
- Dubbo 负载均衡
- Dubbo 服务调用过程
微服务架构
Redis
Elasticsearch
其他
- Dubbo 框架设计
- Dubbo 优雅停机
- dubbo-spring-boot-starter使用指南
- rocketmq-spring-boot-starter使用指南
- Mybatis multi-database in spring-boot 2
- RocketMQ 客户端简单封装
- Otter 入门
杂谈
关于我