diff --git a/src/main/java/cn/hezhaohui/rabbitmq/amqp/MessageConsumer.java b/src/main/java/cn/hezhaohui/rabbitmq/amqp/MessageConsumer.java index 40f7fa5..678a99f 100644 --- a/src/main/java/cn/hezhaohui/rabbitmq/amqp/MessageConsumer.java +++ b/src/main/java/cn/hezhaohui/rabbitmq/amqp/MessageConsumer.java @@ -1,21 +1,29 @@ package cn.hezhaohui.rabbitmq.amqp; -import org.springframework.amqp.core.Message; +import cn.hezhaohui.rabbitmq.entity.Message; +import lombok.extern.slf4j.Slf4j; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; +import java.util.Arrays; + +import static cn.hezhaohui.rabbitmq.constant.RabbitmqConstant.QUEUE_EMAIL_NAME; +import static cn.hezhaohui.rabbitmq.constant.RabbitmqConstant.QUEUE_SMS_NAME; + /** * 消息消费者 */ @Component +@Slf4j public class MessageConsumer { - // TODO: 实现消息监听方法 - // TODO: 定义队列名称常量 - // TODO: 可能需要添加日志记录 + @RabbitListener(queues = QUEUE_EMAIL_NAME) + public void receiveEmailMessage(Message message) { + log.info("receiveEmailMessage: {}", message.getContent()); + } -// @RabbitListener(queues = "test.queue") -// public void receiveMessage(Message message) { -// // TODO: 实现消息处理逻辑 -// } + @RabbitListener(queues = QUEUE_SMS_NAME) + public void receiveSmsMessage(Message message) { + log.info("receiveSmsMessage: {}", message.getContent()); + } } diff --git a/src/main/java/cn/hezhaohui/rabbitmq/entity/Message.java b/src/main/java/cn/hezhaohui/rabbitmq/entity/Message.java index 27ee973..a976342 100644 --- a/src/main/java/cn/hezhaohui/rabbitmq/entity/Message.java +++ b/src/main/java/cn/hezhaohui/rabbitmq/entity/Message.java @@ -1,11 +1,17 @@ package cn.hezhaohui.rabbitmq.entity; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import lombok.AllArgsConstructor; import lombok.Data; +import lombok.NoArgsConstructor; /** * 消息实体类 */ @Data +@JsonIgnoreProperties(ignoreUnknown = true) +@AllArgsConstructor +@NoArgsConstructor public class Message { - final private String content; + private String content; } diff --git a/src/test/java/cn/hezhaohui/rabbitmq/RabbitMQTest.java b/src/test/java/cn/hezhaohui/rabbitmq/RabbitMQTest.java deleted file mode 100644 index f4b96df..0000000 --- a/src/test/java/cn/hezhaohui/rabbitmq/RabbitMQTest.java +++ /dev/null @@ -1,29 +0,0 @@ -package cn.hezhaohui.rabbitmq; - -import cn.hezhaohui.rabbitmq.amqp.MessageProducer; -import jakarta.annotation.Resource; -import org.junit.jupiter.api.Test; -import org.springframework.boot.test.context.SpringBootTest; - -@SpringBootTest -public class RabbitMQTest { - - @Resource - private MessageProducer messageProducer; - - @Test - public void testSendMessage() { - // TODO: 实现生产者发送消息的测试 - // 1. 创建测试消息对象 - // 2. 调用messageProducer.sendMessage()方法 - // 3. 验证消息是否成功发送 - } - - @Test - public void testReceiveMessage() { - // TODO: 实现消费者接收消息的测试 - // 1. 发送测试消息到队列 - // 2. 等待消息被消费 - // 3. 验证消息内容是否正确 - } -} diff --git a/src/test/java/cn/hezhaohui/rabbitmq/amqp/MessageProducerTest.java b/src/test/java/cn/hezhaohui/rabbitmq/amqp/MessageProducerTest.java index 1eaaa2e..fb14596 100644 --- a/src/test/java/cn/hezhaohui/rabbitmq/amqp/MessageProducerTest.java +++ b/src/test/java/cn/hezhaohui/rabbitmq/amqp/MessageProducerTest.java @@ -12,13 +12,21 @@ class MessageProducerTest { @Resource private MessageProducer messageProducer; + +/** + * 测试方法:用于发送电子邮件消息 + * 该方法使用messageProducer对象调用sendEmailMessage方法 + * 传入一个Message对象,其中包含电子邮件地址"hello@example.com" + */ @Test - void sendEmailMessage() { - messageProducer.sendEmailMessage(new Message("hello@example.com")); + void sendEmailMessage() throws InterruptedException { // 测试用例:验证电子邮件消息发送功能 + messageProducer.sendEmailMessage(new Message("hello@example.com")); // 创建并发送电子邮件消息 + Thread.sleep(1000); } @Test - void sendSmsMessage() { + void sendSmsMessage() throws InterruptedException { messageProducer.sendSmsMessage(new Message("1234567890")); + Thread.sleep(1000); } } \ No newline at end of file