diff --git a/labx-13/labx-13-sc-sleuth-mq-activemq/pom.xml b/labx-13/labx-13-sc-sleuth-mq-activemq/pom.xml new file mode 100644 index 00000000..83027555 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-mq-activemq/pom.xml @@ -0,0 +1,64 @@ + + + + labx-13 + cn.iocoder.springboot.labs + 1.0-SNAPSHOT + + 4.0.0 + + labx-13-sc-sleuth-mq-activemq + + + 1.8 + 1.8 + 2.2.4.RELEASE + Hoxton.SR1 + + + + + + + org.springframework.boot + spring-boot-starter-parent + ${spring.boot.version} + pom + import + + + org.springframework.cloud + spring-cloud-dependencies + ${spring.cloud.version} + pom + import + + + + + + + + org.springframework.boot + spring-boot-starter-web + + + + + org.springframework.cloud + spring-cloud-starter-zipkin + + + + + org.springframework.boot + spring-boot-starter-activemq + + + + diff --git a/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/ActiveMQApplication.java b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/ActiveMQApplication.java new file mode 100644 index 00000000..1e5a2725 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/ActiveMQApplication.java @@ -0,0 +1,13 @@ +package cn.iocoder.springboot.labx13.activemqdemo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class ActiveMQApplication { + + public static void main(String[] args) { + SpringApplication.run(ActiveMQApplication.class, args); + } + +} diff --git a/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/consumer/DemoConsumer.java b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/consumer/DemoConsumer.java new file mode 100644 index 00000000..ed4bfb9b --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/consumer/DemoConsumer.java @@ -0,0 +1,19 @@ +package cn.iocoder.springboot.labx13.activemqdemo.consumer; + +import cn.iocoder.springboot.labx13.activemqdemo.message.DemoMessage; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.jms.annotation.JmsListener; +import org.springframework.stereotype.Component; + +@Component +public class DemoConsumer { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @JmsListener(destination = DemoMessage.QUEUE) + public void onMessage(DemoMessage message) { + logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message); + } + +} diff --git a/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/controller/DemoController.java b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/controller/DemoController.java new file mode 100644 index 00000000..d20de264 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/controller/DemoController.java @@ -0,0 +1,26 @@ +package cn.iocoder.springboot.labx13.activemqdemo.controller; + +import cn.iocoder.springboot.labx13.activemqdemo.producer.DemoProducer; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; + +@RestController +@RequestMapping("/demo") +public class DemoController { + + @Autowired + private DemoProducer producer; + + @GetMapping("/activemq") + public String echo() { + this.sendMessage(1); + return "activemq"; + } + + public void sendMessage(Integer id) { + producer.syncSend(id); + } + +} diff --git a/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/message/DemoMessage.java b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/message/DemoMessage.java new file mode 100644 index 00000000..8754fdb1 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/message/DemoMessage.java @@ -0,0 +1,30 @@ +package cn.iocoder.springboot.labx13.activemqdemo.message; + +import java.io.Serializable; + +public class DemoMessage implements Serializable { + + public static final String QUEUE = "QUEUE_DEMO_"; + + /** + * 编号 + */ + private Integer id; + + public DemoMessage setId(Integer id) { + this.id = id; + return this; + } + + public Integer getId() { + return id; + } + + @Override + public String toString() { + return "DemoMessage{" + + "id=" + id + + '}'; + } + +} diff --git a/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/producer/DemoProducer.java b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/producer/DemoProducer.java new file mode 100644 index 00000000..ab2f3d32 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/java/cn/iocoder/springboot/labx13/activemqdemo/producer/DemoProducer.java @@ -0,0 +1,22 @@ +package cn.iocoder.springboot.labx13.activemqdemo.producer; + +import cn.iocoder.springboot.labx13.activemqdemo.message.DemoMessage; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.jms.core.JmsMessagingTemplate; +import org.springframework.stereotype.Component; + +@Component +public class DemoProducer { + + @Autowired + private JmsMessagingTemplate jmsTemplate; + + public void syncSend(Integer id) { + // 创建 DemoMessage 消息 + DemoMessage message = new DemoMessage(); + message.setId(id); + // 同步发送消息 + jmsTemplate.convertAndSend(DemoMessage.QUEUE, message); + } + +} diff --git a/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/resources/application.yaml b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/resources/application.yaml new file mode 100644 index 00000000..7776f853 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-mq-activemq/src/main/resources/application.yaml @@ -0,0 +1,23 @@ +spring: + application: + name: demo-application-activemq + + # ActiveMQ 配置项,对应 ActiveMQProperties 配置类 + activemq: + broker-url: tcp://127.0.0.1:61616 # Activemq Broker 的地址 + user: admin # 账号 + password: admin # 密码 + packages: + trust-all: true # 可信任的反序列化包 + + # Zipkin 配置项,对应 ZipkinProperties 类 + zipkin: + base-url: http://127.0.0.1:9411 # Zipkin 服务的地址 + + # Spring Cloud Sleuth 配置项 + sleuth: + messaging: + # Spring Cloud Sleuth 针对 JMS 组件的配置项 + jms: + enabled: true # 是否开启 + remote-service-name: jms # 远程服务名,默认为 jms diff --git a/labx-13/labx-13-sc-sleuth-mq-kafka/labx-13-sc-sleuth-mq-kafka-producer/src/main/resources/application.yml b/labx-13/labx-13-sc-sleuth-mq-kafka/labx-13-sc-sleuth-mq-kafka-producer/src/main/resources/application.yml index ad0fcea5..c8717725 100644 --- a/labx-13/labx-13-sc-sleuth-mq-kafka/labx-13-sc-sleuth-mq-kafka-producer/src/main/resources/application.yml +++ b/labx-13/labx-13-sc-sleuth-mq-kafka/labx-13-sc-sleuth-mq-kafka-producer/src/main/resources/application.yml @@ -8,15 +8,20 @@ spring: # binders: # Binding 配置项,对应 BindingProperties Map bindings: - demo01-input: + demo01-output: destination: DEMO-TOPIC-01 # 目的地。这里使用 Kafka Topic content-type: application/json # 内容格式。这里使用 JSON - group: demo01-consumer-group # 消费者分组 # Spring Cloud Stream Kafka 配置项 kafka: # Kafka Binder 配置项,对应 KafkaBinderConfigurationProperties 类 binder: brokers: 127.0.0.1:9092 # 指定 Kafka Broker 地址,可以设置多个,以逗号分隔 + # Kafka 自定义 Binding 配置项,对应 KafkaBindingProperties Map + bindings: + demo01-output: + # Kafka Producer 配置项,对应 KafkaProducerProperties 类 + producer: + sync: true # 是否同步发送消息,默认为 false 异步。 # Zipkin 配置项,对应 ZipkinProperties 类 zipkin: diff --git a/labx-13/labx-13-sc-sleuth-mq-kafka/labx-13-sc-stream-mq-kafka-consumer/src/main/resources/application.yml b/labx-13/labx-13-sc-sleuth-mq-kafka/labx-13-sc-stream-mq-kafka-consumer/src/main/resources/application.yml index 4d001e8d..0459a71f 100644 --- a/labx-13/labx-13-sc-sleuth-mq-kafka/labx-13-sc-stream-mq-kafka-consumer/src/main/resources/application.yml +++ b/labx-13/labx-13-sc-sleuth-mq-kafka/labx-13-sc-stream-mq-kafka-consumer/src/main/resources/application.yml @@ -25,7 +25,7 @@ spring: # Spring Cloud Sleuth 配置项 sleuth: messaging: - # Spring Cloud Sleuth 针对 kafka 组件的配置项,例如说 SpringMVC + # Spring Cloud Sleuth 针对 kafka 组件的配置项kafka kafka: enabled: true # 是否开启 remote-service-name: kafka # 远程服务名,默认为 kafka diff --git a/labx-13/pom.xml b/labx-13/pom.xml index 981d234e..50ca3d14 100644 --- a/labx-13/pom.xml +++ b/labx-13/pom.xml @@ -26,6 +26,7 @@ labx-13-sc-sleuth-mq-rabbitmq labx-13-sc-sleuth-mq-kafka + labx-13-sc-sleuth-mq-activemq