From 0f1dcfa82fb2e14396b7ad33c29c3ef8b604db78 Mon Sep 17 00:00:00 2001 From: YunaiV <> Date: Sun, 15 Dec 2019 16:46:54 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=20activemq=20=E7=A4=BA?= =?UTF-8?q?=E4=BE=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- lab-32/lab-32-activemq-demo-delay/pom.xml | 30 +++++++++++++ .../lab32/activemqdemo/Application.java | 13 ++++++ .../activemqdemo/consumer/Demo02Consumer.java | 19 ++++++++ .../activemqdemo/message/Demo02Message.java | 30 +++++++++++++ .../activemqdemo/producer/Demo02Producer.java | 32 ++++++++++++++ .../src/main/resources/application.yaml | 8 ++++ .../lab32/activemqdemo/package-info.java | 1 + .../producer/Demo02ProducerTest.java | 44 +++++++++++++++++++ .../target/classes/application.yaml | 8 ++++ lab-32/pom.xml | 1 + 10 files changed, 186 insertions(+) create mode 100644 lab-32/lab-32-activemq-demo-delay/pom.xml create mode 100644 lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/Application.java create mode 100644 lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/consumer/Demo02Consumer.java create mode 100644 lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/message/Demo02Message.java create mode 100644 lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo02Producer.java create mode 100644 lab-32/lab-32-activemq-demo-delay/src/main/resources/application.yaml create mode 100644 lab-32/lab-32-activemq-demo-delay/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/package-info.java create mode 100644 lab-32/lab-32-activemq-demo-delay/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo02ProducerTest.java create mode 100644 lab-32/lab-32-activemq-demo-delay/target/classes/application.yaml diff --git a/lab-32/lab-32-activemq-demo-delay/pom.xml b/lab-32/lab-32-activemq-demo-delay/pom.xml new file mode 100644 index 00000000..899b4c9b --- /dev/null +++ b/lab-32/lab-32-activemq-demo-delay/pom.xml @@ -0,0 +1,30 @@ + + + + org.springframework.boot + spring-boot-starter-parent + 2.2.1.RELEASE + + + 4.0.0 + + lab-32-activemq-demo-delay + + + + + org.springframework.boot + spring-boot-starter-activemq + + + + + org.springframework.boot + spring-boot-starter-test + test + + + + diff --git a/lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/Application.java b/lab-32/lab-32-activemq-demo-delay/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-delay/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-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/consumer/Demo02Consumer.java b/lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/consumer/Demo02Consumer.java new file mode 100644 index 00000000..d8e5e86c --- /dev/null +++ b/lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/consumer/Demo02Consumer.java @@ -0,0 +1,19 @@ +package cn.iocoder.springboot.lab32.activemqdemo.consumer; + +import cn.iocoder.springboot.lab32.activemqdemo.message.Demo02Message; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.jms.annotation.JmsListener; +import org.springframework.stereotype.Component; + +@Component +public class Demo02Consumer { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @JmsListener(destination = Demo02Message.QUEUE) + public void onMessage(Demo02Message message) { + logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message); + } + +} diff --git a/lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/message/Demo02Message.java b/lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/message/Demo02Message.java new file mode 100644 index 00000000..8637e71e --- /dev/null +++ b/lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/message/Demo02Message.java @@ -0,0 +1,30 @@ +package cn.iocoder.springboot.lab32.activemqdemo.message; + +import java.io.Serializable; + +public class Demo02Message implements Serializable { + + public static final String QUEUE = "QUEUE_DEMO_02"; + + /** + * 编号 + */ + private Integer id; + + public Demo02Message setId(Integer id) { + this.id = id; + return this; + } + + public Integer getId() { + return id; + } + + @Override + public String toString() { + return "Demo01Message{" + + "id=" + id + + '}'; + } + +} diff --git a/lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo02Producer.java b/lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo02Producer.java new file mode 100644 index 00000000..09fc28d7 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-delay/src/main/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo02Producer.java @@ -0,0 +1,32 @@ +package cn.iocoder.springboot.lab32.activemqdemo.producer; + +import cn.iocoder.springboot.lab32.activemqdemo.message.Demo02Message; +import org.apache.activemq.ScheduledMessage; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.jms.core.JmsMessagingTemplate; +import org.springframework.stereotype.Component; + +import java.util.HashMap; +import java.util.Map; + +@Component +public class Demo02Producer { + + @Autowired + private JmsMessagingTemplate jmsTemplate; + + public void syncSend(Integer id, Integer delay) { + // 创建 ClusteringMessage 消息 + Demo02Message message = new Demo02Message(); + message.setId(id); + // 创建 Header + Map headers = null; + if (delay != null && delay > 0) { + headers = new HashMap<>(); + headers.put(ScheduledMessage.AMQ_SCHEDULED_DELAY, delay); + } + // 同步发送消息 + jmsTemplate.convertAndSend(Demo02Message.QUEUE, message, headers); + } + +} diff --git a/lab-32/lab-32-activemq-demo-delay/src/main/resources/application.yaml b/lab-32/lab-32-activemq-demo-delay/src/main/resources/application.yaml new file mode 100644 index 00000000..8179b589 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-delay/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-delay/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/package-info.java b/lab-32/lab-32-activemq-demo-delay/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-delay/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-delay/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo02ProducerTest.java b/lab-32/lab-32-activemq-demo-delay/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo02ProducerTest.java new file mode 100644 index 00000000..40c7311e --- /dev/null +++ b/lab-32/lab-32-activemq-demo-delay/src/test/java/cn/iocoder/springboot/lab32/activemqdemo/producer/Demo02ProducerTest.java @@ -0,0 +1,44 @@ +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 Demo02ProducerTest { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @Autowired + private Demo02Producer producer; + + @Test + public void testSyncSend01() throws InterruptedException { + // 不设置消息的过期时间 + this.testSyncSendDelay(null); + } + + @Test + public void testSyncSend02() throws InterruptedException { + // 设置发送消息的过期时间为 5000 毫秒 + this.testSyncSendDelay(5000); + } + + private void testSyncSendDelay(Integer delay) throws InterruptedException { + int id = (int) (System.currentTimeMillis() / 1000); + producer.syncSend(id, delay); + logger.info("[testSyncSendDelay][发送编号:[{}] 发送成功]", id); + + // 阻塞等待,保证消费 + new CountDownLatch(1).await(); + } + +} diff --git a/lab-32/lab-32-activemq-demo-delay/target/classes/application.yaml b/lab-32/lab-32-activemq-demo-delay/target/classes/application.yaml new file mode 100644 index 00000000..8179b589 --- /dev/null +++ b/lab-32/lab-32-activemq-demo-delay/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/pom.xml b/lab-32/pom.xml index bcb70c44..37308ba2 100644 --- a/lab-32/pom.xml +++ b/lab-32/pom.xml @@ -15,6 +15,7 @@ lab-32-activemq-native lab-32-activemq-demo lab-32-activemq-demo-message-model + lab-32-activemq-demo-delay