From 21a33b5f7401e5dd39d5da6ffdc598f9a22446a0 Mon Sep 17 00:00:00 2001
From: tangmingyou <234767776@qq.com>
Date: Mon, 30 May 2022 18:09:02 +0800
Subject: [PATCH] =?UTF-8?q?mq=E6=B6=88=E6=81=AF=E5=BA=8F=E5=88=97=E5=8C=96?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
im-das-api/im-das-user-api/pom.xml | 11 +++++
.../sopod/soim/das/user/api/mq/ChatMQ.java | 15 ++++++
.../das/user/api/mq/ChatMessageConverter.java | 39 +++++++++++++++
.../sopod/soim/das/user/api/mq/ChatMsg.java | 27 +++++++++++
.../sopod/soim/das/user/api/mq/MQType.java | 47 +++++++++++++++++++
im-das/im-das-user/pom.xml | 4 ++
.../sopod/soim/das/user/amqp/Receiver.java | 16 ++++++-
.../net/sopod/soim/das/user/amqp/Sender.java | 22 +++++----
.../soim/das/user/config/RabbitConfig.java | 8 ++++
9 files changed, 178 insertions(+), 11 deletions(-)
create mode 100644 im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQ.java
create mode 100644 im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMessageConverter.java
create mode 100644 im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMsg.java
create mode 100644 im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/MQType.java
diff --git a/im-das-api/im-das-user-api/pom.xml b/im-das-api/im-das-user-api/pom.xml
index 4cd41e5..ab81e1d 100644
--- a/im-das-api/im-das-user-api/pom.xml
+++ b/im-das-api/im-das-user-api/pom.xml
@@ -17,6 +17,17 @@
mybatis-plus-annotation
3.5.0
+
+ net.sopod
+ im-common
+ ${soim.version}
+ provided
+
+
+ org.springframework.boot
+ spring-boot-starter-amqp
+ provided
+
\ No newline at end of file
diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQ.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQ.java
new file mode 100644
index 0000000..4c4d390
--- /dev/null
+++ b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQ.java
@@ -0,0 +1,15 @@
+package net.sopod.soim.das.user.api.mq;
+
+/**
+ * ChatMQ
+ *
+ * @author tmy
+ * @date 2022-05-30 17:43
+ */
+public class ChatMQ {
+
+ private MQType mqType;
+
+ private byte[] body;
+
+}
diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMessageConverter.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMessageConverter.java
new file mode 100644
index 0000000..86c0f44
--- /dev/null
+++ b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMessageConverter.java
@@ -0,0 +1,39 @@
+package net.sopod.soim.das.user.api.mq;
+
+import net.sopod.soim.common.constant.Consts;
+import net.sopod.soim.common.util.Jackson;
+import net.sopod.soim.das.user.api.model.entity.ImMessage;
+import org.springframework.amqp.core.Message;
+import org.springframework.amqp.core.MessageProperties;
+import org.springframework.amqp.support.converter.AbstractMessageConverter;
+import org.springframework.amqp.support.converter.MessageConversionException;
+
+/**
+ * ChatMessageConverter
+ *
+ * @author tmy
+ * @date 2022-05-30 18:07
+ */
+public class ChatMessageConverter extends AbstractMessageConverter {
+
+ @Override
+ protected Message createMessage(Object object, MessageProperties messageProperties) {
+ if (object instanceof Message) {
+ return (Message) object;
+ }
+ byte[] bytes = Jackson.msgpack().serializeBytes(object);
+ if (bytes.length > 100 * Consts.KB) {
+ // TODO 压缩
+ }
+
+ return null;
+ }
+
+ @Override
+ public Object fromMessage(Message message) throws MessageConversionException {
+ byte[] body = message.getBody();
+ Jackson.msgpack().deserializeBytes(body, ImMessage.class);
+ return null;
+ }
+
+}
diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMsg.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMsg.java
new file mode 100644
index 0000000..8818df2
--- /dev/null
+++ b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMsg.java
@@ -0,0 +1,27 @@
+package net.sopod.soim.das.user.api.mq;
+
+import net.sopod.soim.common.util.Jackson;
+import net.sopod.soim.common.util.StringUtil;
+import net.sopod.soim.das.user.api.model.entity.ImMessage;
+import org.springframework.amqp.core.Message;
+import org.springframework.amqp.core.MessageProperties;
+
+/**
+ * ChatMsg
+ *
+ * @author tmy
+ * @date 2022-05-30 15:03
+ */
+public class ChatMsg extends Message {
+
+ public ChatMsg(ImMessage imMessage) {
+ super(Jackson.msgpack().serializeBytes(imMessage), getMessageProperties(imMessage));
+ }
+
+ private static MessageProperties getMessageProperties(ImMessage imMessage) {
+ MessageProperties messageProperties = new MessageProperties();
+ messageProperties.setMessageId(StringUtil.toString(imMessage.getId()));
+ return messageProperties;
+ }
+
+}
diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/MQType.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/MQType.java
new file mode 100644
index 0000000..041e047
--- /dev/null
+++ b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/MQType.java
@@ -0,0 +1,47 @@
+package net.sopod.soim.das.user.api.mq;
+
+import net.sopod.soim.das.user.api.model.entity.ImGroupMessage;
+import net.sopod.soim.das.user.api.model.entity.ImMessage;
+
+/**
+ * MQType
+ *
+ * @author tmy
+ * @date 2022-05-30 17:43
+ */
+public enum MQType {
+
+ /** 单聊消息 */
+ IM_MESSAGE("IM_MESSAGE", ImMessage.class),
+
+ /** 群聊消息 */
+ IM_GROUP_MESSAGE("IM_GROUP_MESSAGE", ImGroupMessage.class);
+
+ private final Class> beanClazz;
+
+ private final String queueName;
+
+ MQType(String queueName, Class> beanClazz) {
+ this.queueName = queueName;
+ this.beanClazz = beanClazz;
+ }
+
+ public Class> getBeanClazz() {
+ return beanClazz;
+ }
+
+ public String getQueueName() {
+ return queueName;
+ }
+
+ public static MQType getMQType(Class> beanClazz) {
+ MQType[] values = values();
+ for (MQType type : values) {
+ if (type.beanClazz == beanClazz) {
+ return type;
+ }
+ }
+ return null;
+ }
+
+}
diff --git a/im-das/im-das-user/pom.xml b/im-das/im-das-user/pom.xml
index b2b935e..96c68ec 100644
--- a/im-das/im-das-user/pom.xml
+++ b/im-das/im-das-user/pom.xml
@@ -88,6 +88,10 @@
org.springframework.boot
spring-boot-starter-amqp
+
+ org.msgpack
+ jackson-dataformat-msgpack
+
org.springframework.boot
diff --git a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Receiver.java b/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Receiver.java
index 17cf173..8ae2c93 100644
--- a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Receiver.java
+++ b/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Receiver.java
@@ -1,11 +1,15 @@
package net.sopod.soim.das.user.amqp;
+import net.sopod.soim.common.util.Jackson;
+import net.sopod.soim.das.user.api.model.entity.ImMessage;
+import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
/**
* Listener
+ * 只能接受字符串和字节数组
*
* @author tmy
* @date 2022-05-28 16:26
@@ -15,8 +19,16 @@ import org.springframework.stereotype.Component;
public class Receiver {
@RabbitHandler
- public void process(String hello) {
- System.out.println("Receiver: " + hello);
+ public void process(Message message) {
+ byte[] body = message.getBody();
+ ImMessage imMessage = Jackson.msgpack().deserializeBytes(body, ImMessage.class);
+ System.out.println("Receiver: " + imMessage);
+ }
+
+ @RabbitHandler
+ public void process(byte[] payload) {
+ ImMessage imMessage = Jackson.msgpack().deserializeBytes(payload, ImMessage.class);
+ System.out.println("receive2: " + imMessage);
}
}
diff --git a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Sender.java b/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Sender.java
index c5b48b0..6cb0076 100644
--- a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Sender.java
+++ b/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Sender.java
@@ -1,16 +1,14 @@
package net.sopod.soim.das.user.amqp;
import net.sopod.soim.common.util.ImClock;
+import net.sopod.soim.das.user.api.model.entity.ImMessage;
+import net.sopod.soim.das.user.api.mq.ChatMsg;
import org.springframework.amqp.core.AmqpTemplate;
-import org.springframework.amqp.core.Message;
-import org.springframework.amqp.core.MessageProperties;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.ApplicationListener;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.stereotype.Component;
-import java.util.UUID;
-
/**
* Sender
*
@@ -24,11 +22,17 @@ public class Sender implements ApplicationListener {
public void onApplicationEvent(ApplicationReadyEvent event) {
ConfigurableApplicationContext context = event.getApplicationContext();
AmqpTemplate amqpTemplate = context.getBean(AmqpTemplate.class);
- amqpTemplate.convertAndSend("hello", "hello:" + ImClock.millis());
- MessageProperties messageProperties = new MessageProperties();
- messageProperties.setMessageId(UUID.randomUUID().toString());
- Message message = new Message(new byte[] {}, messageProperties);
-
+ // amqpTemplate.convertAndSend("hello", "hello:" + ImClock.millis());
+ ImMessage imMessage = new ImMessage()
+ .setId(10086L)
+ .setSender(10808L)
+ .setReceiver(180002L)
+ .setFriendId(101010L)
+ .setContent("何时")
+ .setCreateTime(ImClock.date());
+ ChatMsg chatMsg = new ChatMsg(imMessage);
+ amqpTemplate.convertAndSend("hello", 123);
+ amqpTemplate.send("hello", chatMsg);
System.out.println("send msg....");
}
diff --git a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/config/RabbitConfig.java b/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/config/RabbitConfig.java
index 67f920d..2ed242b 100644
--- a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/config/RabbitConfig.java
+++ b/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/config/RabbitConfig.java
@@ -1,6 +1,7 @@
package net.sopod.soim.das.user.config;
import org.springframework.amqp.core.Queue;
+import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@@ -19,4 +20,11 @@ public class RabbitConfig {
return new Queue("hello");
}
+ @Bean
+ public RabbitTemplate rabbitTemplate() {
+ RabbitTemplate rabbitTemplate = new RabbitTemplate();
+ rabbitTemplate.setMessageConverter(null);
+ return null;
+ }
+
}