增加 spring boot rocketmq 示例

This commit is contained in:
YunaiV
2019-12-05 01:44:47 +08:00
parent 6b00101687
commit eb5efcf923
4 changed files with 99 additions and 4 deletions

View File

@@ -0,0 +1,25 @@
package cn.iocoder.springboot.lab31.rocketmqdemo.consumer;
import cn.iocoder.springboot.lab31.rocketmqdemo.message.Demo05Message;
import cn.iocoder.springboot.lab31.rocketmqdemo.message.Demo07Message;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
@Component
@RocketMQMessageListener(
topic = Demo07Message.TOPIC,
consumerGroup = "demo07-consumer-group-" + Demo05Message.TOPIC
)
public class Demo07Consumer implements RocketMQListener<Demo07Message> {
private Logger logger = LoggerFactory.getLogger(getClass());
@Override
public void onMessage(Demo07Message message) {
logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message);
}
}

View File

@@ -1,4 +1,4 @@
package cn.iocoder.springboot.lab31.rocketmqdemo;
package cn.iocoder.springboot.lab31.rocketmqdemo.core;
import org.apache.rocketmq.spring.annotation.ExtRocketMQTemplateConfiguration;
import org.apache.rocketmq.spring.core.RocketMQTemplate;

View File

@@ -1,18 +1,54 @@
package cn.iocoder.springboot.lab31.rocketmqdemo.producer;
import org.apache.rocketmq.client.producer.SendResult;
import cn.iocoder.springboot.lab31.rocketmqdemo.message.Demo07Message;
import org.apache.rocketmq.client.producer.TransactionSendResult;
import org.apache.rocketmq.spring.annotation.RocketMQTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionState;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
@Component
public class Demo07Producer {
private static final String TX_PRODUCER_GROUP = "demo07-producer-group";
@Autowired
private RocketMQTemplate rocketMQTemplate;
public SendResult syncSendOrderly(Integer id) {
return null;
public TransactionSendResult sendMessageInTransaction(Integer id) {
// 创建 Demo07Message 消息
Message message = MessageBuilder.withPayload(new Demo07Message().setId(id))
.build();
// 发送事务消息
return rocketMQTemplate.sendMessageInTransaction(TX_PRODUCER_GROUP, Demo07Message.TOPIC, message,
id);
}
@RocketMQTransactionListener(txProducerGroup = TX_PRODUCER_GROUP)
public class TransactionListenerImpl implements RocketMQLocalTransactionListener {
private Logger logger = LoggerFactory.getLogger(getClass());
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// ... local transaction process, return rollback, commit or unknown
logger.info("[executeLocalTransaction][执行本地事务,消息:{} arg{}]", msg, arg);
return RocketMQLocalTransactionState.UNKNOWN;
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
// ... check transaction status and return rollback, commit or unknown
logger.info("[checkLocalTransaction][回查消息:{}]", msg);
return RocketMQLocalTransactionState.COMMIT;
}
}
}

View File

@@ -0,0 +1,34 @@
package cn.iocoder.springboot.lab31.rocketmqdemo.producer;
import cn.iocoder.springboot.lab31.rocketmqdemo.Application;
import org.apache.rocketmq.client.producer.SendResult;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.junit4.SpringRunner;
import java.util.concurrent.CountDownLatch;
@RunWith(SpringRunner.class)
@SpringBootTest(classes = Application.class)
public class Demo07ProducerTest {
private Logger logger = LoggerFactory.getLogger(getClass());
@Autowired
private Demo07Producer producer;
@Test
public void testSendMessageInTransaction() throws InterruptedException {
int id = (int) (System.currentTimeMillis() / 1000);
SendResult result = producer.sendMessageInTransaction(id);
logger.info("[testSendMessageInTransaction][发送编号:[{}] 发送结果:[{}]]", id, result);
// 阻塞等待,保证消费
new CountDownLatch(1).await();
}
}