diff --git a/labx-10/labx-10-sc-stream-rabbitmq-consumer-ack/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/Demo01Consumer.java b/labx-10/labx-10-sc-stream-rabbitmq-consumer-ack/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/Demo01Consumer.java index 2009a2f8..42a776b1 100644 --- a/labx-10/labx-10-sc-stream-rabbitmq-consumer-ack/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/Demo01Consumer.java +++ b/labx-10/labx-10-sc-stream-rabbitmq-consumer-ack/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/Demo01Consumer.java @@ -27,7 +27,7 @@ public class Demo01Consumer { logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message); // 提交消费进度 // if (message.getId() % 2 == 1) { - if (index.incrementAndGet() % 2 == 1) { + if (index.incrementAndGet() == 1) { // ack 确认消息 // 第二个参数 multiple ,用于批量确认消息,为了减少网络流量,手动确认可以被批处。 // 1. 当 multiple 为 true 时,则可以一次性确认 deliveryTag 小于等于传入值的所有消息 diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/pom.xml b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/pom.xml new file mode 100644 index 00000000..c68b3e40 --- /dev/null +++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/pom.xml @@ -0,0 +1,58 @@ + + + + labx-10 + cn.iocoder.springboot.labs + 1.0-SNAPSHOT + + 4.0.0 + + labx-10-sc-stream-rabbitmq-producer-confirm + + + 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-stream-rabbit + + + + diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/ProducerApplication.java b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/ProducerApplication.java new file mode 100644 index 00000000..c128d2f4 --- /dev/null +++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/ProducerApplication.java @@ -0,0 +1,16 @@ +package cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo; + +import cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.message.MySource; +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.cloud.stream.annotation.EnableBinding; + +@SpringBootApplication +@EnableBinding(MySource.class) +public class ProducerApplication { + + public static void main(String[] args) { + SpringApplication.run(ProducerApplication.class, args); + } + +} diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java new file mode 100644 index 00000000..032b356b --- /dev/null +++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java @@ -0,0 +1,67 @@ +package cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.controller; + +import cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.message.Demo01Message; +import cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.message.MySource; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.integration.annotation.ServiceActivator; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.MessageBuilder; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; + +import java.util.Random; + +@RestController +@RequestMapping("/demo01") +public class Demo01Controller { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @Autowired + private MySource mySource; + + @GetMapping("/send") + public boolean send() { + // 创建 Message + Demo01Message message = new Demo01Message() + .setId(new Random().nextInt()); + // 创建 Spring Message 对象 + Message springMessage = MessageBuilder.withPayload(message) + .build(); + // 发送消息 + return mySource.demo01Output().send(springMessage); + } + +// @StreamListener(IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME) // errorChannel +// public void globalHandleError(ErrorMessage errorMessage) { +// logger.error("[globalHandleError][payload:{}]", errorMessage.getPayload().getMessage()); +// logger.error("[globalHandleError][originalMessage:{}]", errorMessage.getOriginalMessage()); +// logger.error("[globalHandleError][headers:{}]", errorMessage.getHeaders()); +// } + +// @StreamListener("ooxx") // errorChannel +// public void ooxx(ErrorMessage errorMessage) { +// logger.error("[globalHandleError][payload:{}]", errorMessage.getPayload().getMessage()); +// logger.error("[globalHandleError][originalMessage:{}]", errorMessage.getOriginalMessage()); +// logger.error("[globalHandleError][headers:{}]", errorMessage.getHeaders()); +// } + +//// @ServiceActivator(inputChannel = "demo-producer-application.ooxx") +// @ServiceActivator(inputChannel = "demo-producer-application.ooxx") +//// @StreamListener("demo-producer-application.ooxx") // errorChannel +// public void handleError(Message errorMessage) { +//// logger.error("[handleError][payload:{}]", errorMessage.getPayload().getMessage()); +//// logger.error("[handleError][originalMessage:{}]", errorMessage.getOriginalMessage()); +//// logger.error("[handleError][headers:{}]", errorMessage.getHeaders()); +// System.out.println(); +// } + + @ServiceActivator(inputChannel = "publisher-confirm") + public void onPublisherConfirm(Message message) { + logger.debug("on publisher confirm"); + } + +} diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/Demo01Message.java b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/Demo01Message.java new file mode 100644 index 00000000..7f015623 --- /dev/null +++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/Demo01Message.java @@ -0,0 +1,29 @@ +package cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.message; + +/** + * 示例 01 的 Message 消息 + */ +public class Demo01Message { + + /** + * 编号 + */ + private Integer id; + + public Demo01Message setId(Integer id) { + this.id = id; + return this; + } + + public Integer getId() { + return id; + } + + @Override + public String toString() { + return "Demo01Message{" + + "id=" + id + + '}'; + } + +} diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/MySource.java b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/MySource.java new file mode 100644 index 00000000..64d06834 --- /dev/null +++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/MySource.java @@ -0,0 +1,11 @@ +package cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.message; + +import org.springframework.cloud.stream.annotation.Output; +import org.springframework.messaging.MessageChannel; + +public interface MySource { + + @Output("demo01-output") + MessageChannel demo01Output(); + +} diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/resources/application.yml b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/resources/application.yml new file mode 100644 index 00000000..20d859ac --- /dev/null +++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-confirm/src/main/resources/application.yml @@ -0,0 +1,39 @@ +spring: + application: + name: demo-producer-application + cloud: + # Spring Cloud Stream 配置项,对应 BindingServiceProperties 类 + stream: + # Binder 配置项,对应 BinderProperties Map + binders: + rabbit001: + type: rabbit # 设置 Binder 的类型 + environment: # 设置 Binder 的环境配置 + # 如果是 RabbitMQ 类型的时候,则对应的是 RabbitProperties 类 + spring: + rabbitmq: + host: 127.0.0.1 # RabbitMQ 服务的地址 + port: 5672 # RabbitMQ 服务的端口 + username: guest # RabbitMQ 服务的账号 + password: guest # RabbitMQ 服务的密码 + publisherConfirms: true + publisherReturns: true + publisher-confirm-type: simple # 设置 Confirm 类型为 SIMPLE 。 + # Binding 配置项,对应 BindingProperties Map + bindings: + demo01-output: + destination: DEMO-TOPIC-01 # 目的地。这里使用 RabbitMQ Exchange + content-type: application/json # 内容格式。这里使用 JSON + binder: rabbit001 # 设置使用的 Binder 名字 + producer: + errorChannelEnabled: true + # RabbitMQ 自定义 Binding 配置项,对应 RabbitBindingProperties Map + rabbit: + bindings: + demo01-output: + # RabbitMQ Producer 配置项,对应 RabbitProducerProperties 类 + producer: + confirmAckChannel: publisher-confirm + +server: + port: 18080 diff --git a/labx-10/pom.xml b/labx-10/pom.xml index 100dc550..4cd102a6 100644 --- a/labx-10/pom.xml +++ b/labx-10/pom.xml @@ -41,6 +41,9 @@ labx-10-sc-stream-rabbitmq-consumer-ack + + labx-10-sc-stream-rabbitmq-producer-confirm +