<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
</dependency>
rocketmq:
name-server: 127.0.0.1:9876
# 纯消费者不需要以下配置
producer:
group: test-group
@Autowired
private final RocketMQTemplate rocketMQTemplate;
// 默认使用同步发送, 但拿不到回执, 源码见下文org.apache.rocketmq.spring.core.RocketMQTemplate.doSent
rocketMQTemplate.convertAndSend("test-topic", entity);
rocketMQTemplate.send("test-topic", MessageBuilder.withPayload(entity).build());
// 带tag
rocketMQTemplate.convertAndSend("test-topic:tag1", entity);
rocketMQTemplate.send("test-topic:tag2", MessageBuilder.withPayload(entity).build());
rocketMQTemplate.sendOneWay("test-topic", MessageBuilder.withPayload("oneway message").build());
SendResult result = rocketMQTemplate.syncSend("test-topic", entity, timeout, delayLevel);
SendResult result = rocketMQTemplate.syncSendOrderly("test-topic", "order message", "hashkey", timeout);
rocketMQTemplate.asyncSend("test-topic", MessageBuilder.withPayload(entity).build(), new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {}
@Override
public void onException(Throwable e) {}
}, timeout);
@Component
@RocketMQMessageListener(topic = "test-topic", consumerGroup = "test-consumer", consumeMode = ConsumeMode.ORDERLY)
public class DemoListener implements RocketMQListener<UserEntity> {
...
@Override
public void onMessage(MyEntity entity) {
logger.info("---->收到了消息了!");
logger.info("---->" + entity.toString());
}
}
org.apache.rocketmq.spring.core.RocketMQTemplate.doSent
@Override
protected void doSend(String destination, Message<?> message) {
SendResult sendResult = syncSend(destination, message);
log.debug("send message to `{}` finished. result:{}", destination, sendResult);
}