diff --git a/lab-32/lab-32-activemq-demo-consume-retry/pom.xml b/lab-32/lab-32-activemq-demo-consume-retry/pom.xml new file mode 100644 index 00000000..b7b8e6c9 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-consume-retry/pom.xml @@ -0,0 +1,30 @@ + + + + org.springframework.boot + spring-boot-starter-parent + 2.2.1.RELEASE + + + 4.0.0 + + lab-32-activemq-demo-consume-retry + + + + + org.springframework.boot + spring-boot-starter-activemq + + + + + org.springframework.boot + spring-boot-starter-test + test + + + + diff --git a/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/Application.java b/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/Application.java new file mode 100644 index 00000000..09cc5bc3 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/Application.java @@ -0,0 +1,13 @@ +package cn.iocoder.springboot.lab32.activemqdemo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class Application { + + public static void main(String[] args) { + SpringApplication.run(Application.class, args); + } + +} diff --git a/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/config/ActiveMQConnectionFactoryCustomizerImpl.java b/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/config/ActiveMQConnectionFactoryCustomizerImpl.java new file mode 100644 index 00000000..7a689bd0 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/config/ActiveMQConnectionFactoryCustomizerImpl.java @@ -0,0 +1,15 @@ +package cn.iocoder.springboot.lab32.activemqdemo.config; + +import org.apache.activemq.ActiveMQConnectionFactory; +import org.springframework.boot.autoconfigure.jms.activemq.ActiveMQConnectionFactoryCustomizer; +import org.springframework.context.annotation.Configuration; + +@Configuration +public class ActiveMQConnectionFactoryCustomizerImpl implements ActiveMQConnectionFactoryCustomizer { + + @Override + public void customize(ActiveMQConnectionFactory factory) { + System.out.println(); + } + +} diff --git a/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/consumer/Demo01Consumer.java b/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/consumer/Demo01Consumer.java new file mode 100644 index 00000000..6aaf02fb --- /dev/null +++ b/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/consumer/Demo01Consumer.java @@ -0,0 +1,21 @@ +package cn.iocoder.springboot.lab32.activemqdemo.consumer; + +import cn.iocoder.springboot.lab32.activemqdemo.message.Demo05Message; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.jms.annotation.JmsListener; +import org.springframework.stereotype.Component; + +@Component +public class Demo01Consumer { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @JmsListener(destination = Demo05Message.QUEUE) + public void onMessage(Demo05Message message) { + logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message); + // 注意,此处抛出一个 RuntimeException 异常,模拟消费失败 + throw new RuntimeException("我就是故意抛出一个异常"); + } + +} diff --git a/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/message/Demo05Message.java b/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/message/Demo05Message.java new file mode 100644 index 00000000..20365976 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/message/Demo05Message.java @@ -0,0 +1,30 @@ +package cn.iocoder.springboot.lab32.activemqdemo.message; + +import java.io.Serializable; + +public class Demo05Message implements Serializable { + + public static final String QUEUE = "QUEUE_DEMO_05"; + + /** + * 编号 + */ + private Integer id; + + public Demo05Message setId(Integer id) { + this.id = id; + return this; + } + + public Integer getId() { + return id; + } + + @Override + public String toString() { + return "Demo05Message{" + + "id=" + id + + '}'; + } + +} diff --git a/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo01Producer.java b/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo01Producer.java new file mode 100644 index 00000000..b5b8f45d --- /dev/null +++ b/lab-32/lab-32-activemq-demo-consume-retry/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo01Producer.java @@ -0,0 +1,22 @@ +package cn.iocoder.springboot.lab32.activemqdemo.producer; + +import cn.iocoder.springboot.lab32.activemqdemo.message.Demo05Message; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.jms.core.JmsMessagingTemplate; +import org.springframework.stereotype.Component; + +@Component +public class Demo01Producer { + + @Autowired + private JmsMessagingTemplate jmsTemplate; + + public void syncSend(Integer id) { + // 创建 ClusteringMessage 消息 + Demo05Message message = new Demo05Message(); + message.setId(id); + // 同步发送消息 + jmsTemplate.convertAndSend(Demo05Message.QUEUE, message); + } + +} diff --git a/lab-32/lab-32-activemq-demo-consume-retry/src/main/resources/application.yaml b/lab-32/lab-32-activemq-demo-consume-retry/src/main/resources/application.yaml new file mode 100644 index 00000000..8179b589 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-consume-retry/src/main/resources/application.yaml @@ -0,0 +1,8 @@ +spring: + # ActiveMQ 配置项,对应 ActiveMQProperties 配置类 + activemq: + broker-url: tcp://127.0.0.1:61616 # RabbitMQ Broker 的地址 + user: admin # 账号 + password: admin # 密码 + packages: + trust-all: true # 可信任的反序列化包 diff --git a/lab-32/lab-32-activemq-demo-consume-retry/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/package-info.java b/lab-32/lab-32-activemq-demo-consume-retry/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/package-info.java new file mode 100644 index 00000000..6878c298 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-consume-retry/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/package-info.java @@ -0,0 +1 @@ +package cn.iocoder.springboot.lab32.activemqdemo; diff --git a/lab-32/lab-32-activemq-demo-consume-retry/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo01ProducerTest.java b/lab-32/lab-32-activemq-demo-consume-retry/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo01ProducerTest.java new file mode 100644 index 00000000..0e4560c3 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-consume-retry/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo01ProducerTest.java @@ -0,0 +1,34 @@ +package cn.iocoder.springboot.lab32.activemqdemo.producer; + +import cn.iocoder.springboot.lab32.activemqdemo.Application; +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 Demo01ProducerTest { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @Autowired + private Demo01Producer producer; + + @Test + public void testSyncSend() throws InterruptedException { + // 发送消息 + int id = (int) (System.currentTimeMillis() / 1000); + producer.syncSend(id); + logger.info("[testSyncSend][发送编号:[{}] 发送成功]", id); + + // 阻塞等待,保证消费 + new CountDownLatch(1).await(); + } + +} diff --git a/lab-32/lab-32-activemq-demo-consume-retry/target/classes/application.yaml b/lab-32/lab-32-activemq-demo-consume-retry/target/classes/application.yaml new file mode 100644 index 00000000..8179b589 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-consume-retry/target/classes/application.yaml @@ -0,0 +1,8 @@ +spring: + # ActiveMQ 配置项,对应 ActiveMQProperties 配置类 + activemq: + broker-url: tcp://127.0.0.1:61616 # RabbitMQ Broker 的地址 + user: admin # 账号 + password: admin # 密码 + packages: + trust-all: true # 可信任的反序列化包 diff --git a/lab-32/lab-32-activemq-demo-message-model/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/ClusteringProducer.java b/lab-32/lab-32-activemq-demo-message-model/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/ClusteringProducer.java index bad7a71e..044d0454 100644 --- a/lab-32/lab-32-activemq-demo-message-model/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/ClusteringProducer.java +++ b/lab-32/lab-32-activemq-demo-message-model/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/ClusteringProducer.java @@ -11,14 +11,14 @@ import javax.annotation.Resource; public class ClusteringProducer { @Resource(name = ActiveMQConfig.CLUSTERING_JMS_TEMPLATE_BEAN_NAME) - private JmsMessagingTemplate rabbitTemplate; + private JmsMessagingTemplate jmsTemplate; public void syncSend(Integer id) { // 创建 ClusteringMessage 消息 ClusteringMessage message = new ClusteringMessage(); message.setId(id); // 同步发送消息 - rabbitTemplate.convertAndSend(ClusteringMessage.QUEUE, message); + jmsTemplate.convertAndSend(ClusteringMessage.QUEUE, message); } } diff --git a/lab-32/pom.xml b/lab-32/pom.xml index 5be9db26..6f2b08f4 100644 --- a/lab-32/pom.xml +++ b/lab-32/pom.xml @@ -18,6 +18,7 @@ lab-32-activemq-demo-delay lab-32-activemq-demo-concurrency lab-32-activemq-demo-orderly + lab-32-activemq-demo-consume-retry