diff --git a/im-common/src/main/java/net/sopod/soim/common/util/Converter.java b/im-common/src/main/java/net/sopod/soim/common/util/Converter.java
new file mode 100644
index 0000000..5c6f36d
--- /dev/null
+++ b/im-common/src/main/java/net/sopod/soim/common/util/Converter.java
@@ -0,0 +1,25 @@
+package net.sopod.soim.common.util;
+
+/**
+ * Converter
+ *
+ * @author tmy
+ * @date 2022-05-30 22:23
+ */
+public class Converter {
+
+ public static String toString(Object data) {
+ return data == null ? null : data.toString();
+ }
+
+ /**
+ * @exception NumberFormatException if the string cannot be parsed as an integer.
+ */
+ public static int toInt(Object value) {
+ if (value instanceof Number) {
+ return ((Number)value).intValue();
+ }
+ return Integer.parseInt(String.valueOf(value));
+ }
+
+}
diff --git a/im-das-api/im-das-user-api/pom.xml b/im-das-api/im-das-user-api/pom.xml
index ab81e1d..0e9b224 100644
--- a/im-das-api/im-das-user-api/pom.xml
+++ b/im-das-api/im-das-user-api/pom.xml
@@ -28,6 +28,10 @@
spring-boot-starter-amqp
provided
+
+ org.xerial.snappy
+ snappy-java
+
\ 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/config/ChatRabbitMQAutoConfiguration.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/ChatRabbitMQAutoConfiguration.java
new file mode 100644
index 0000000..3644a58
--- /dev/null
+++ b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/ChatRabbitMQAutoConfiguration.java
@@ -0,0 +1,56 @@
+package net.sopod.soim.das.user.api.config;
+
+import net.sopod.soim.das.user.api.mq.ChatMQMessageConverter;
+import net.sopod.soim.das.user.api.mq.ChatQueueType;
+import org.springframework.amqp.core.Queue;
+import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
+import org.springframework.amqp.rabbit.core.RabbitTemplate;
+import org.springframework.amqp.support.converter.MessageConverter;
+import org.springframework.beans.factory.support.BeanDefinitionRegistry;
+import org.springframework.beans.factory.support.RootBeanDefinition;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.context.annotation.Import;
+import org.springframework.context.annotation.ImportBeanDefinitionRegistrar;
+import org.springframework.core.type.AnnotationMetadata;
+
+/**
+ * ChatRabbitMQAutoConfiguration
+ * RabbitMQ的配置类,用来配置队列、交换器、路由等高级信息
+ *
+ * @author tmy
+ * @date 2022-05-30 23:14
+ */
+@Import(ChatRabbitMQAutoConfiguration.class)
+@Configuration
+public class ChatRabbitMQAutoConfiguration implements ImportBeanDefinitionRegistrar {
+
+ /**
+ * direct直连模式, 按照routingkey分发到指定队列
+ */
+ @Override
+ public void registerBeanDefinitions(AnnotationMetadata importingClassMetadata, BeanDefinitionRegistry registry) {
+ // 消息队列Queue绑定
+ for (ChatQueueType chatQueueType : ChatQueueType.values()) {
+ Queue queue = new Queue(chatQueueType.getQueueName(), true);
+ RootBeanDefinition queueDef = new RootBeanDefinition();
+ queueDef.setInstanceSupplier(() -> queue);
+ queueDef.setBeanClass(Queue.class);
+ registry.registerBeanDefinition(queue.getName(), queueDef);
+ }
+ }
+
+ @Bean
+ public MessageConverter rabbitMessageConverter() {
+ return new ChatMQMessageConverter();
+ }
+
+ @Bean
+ public RabbitTemplate rabbitTemplate(CachingConnectionFactory connectionFactory,
+ MessageConverter messageConverter) {
+ RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
+ rabbitTemplate.setMessageConverter(messageConverter);
+ return rabbitTemplate;
+ }
+
+}
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
deleted file mode 100644
index 4c4d390..0000000
--- a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQ.java
+++ /dev/null
@@ -1,15 +0,0 @@
-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/ChatMQMessageConverter.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQMessageConverter.java
new file mode 100644
index 0000000..67a3ee9
--- /dev/null
+++ b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQMessageConverter.java
@@ -0,0 +1,76 @@
+package net.sopod.soim.das.user.api.mq;
+
+import net.sopod.soim.common.constant.Consts;
+import net.sopod.soim.common.util.Converter;
+import net.sopod.soim.common.util.ImClock;
+import net.sopod.soim.common.util.Jackson;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+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;
+import org.xerial.snappy.Snappy;
+
+import java.io.IOException;
+
+/**
+ * ChatMessageConverter
+ * 聊天消息到 MQ 生产和消费序列化 {@link ChatQueueType}
+ *
+ * @author tmy
+ * @date 2022-05-30 18:07
+ */
+public class ChatMQMessageConverter extends AbstractMessageConverter {
+
+ private static final Logger logger = LoggerFactory.getLogger(ChatMQMessageConverter.class);
+
+ public static final String MQ_MSG_TYPE = "MQ_MSG_TYPE";
+
+ public static final String MQ_MSG_IS_SNAPPY = "MQ_MSG_IS_SNAPPY";
+
+ @Override
+ protected Message createMessage(Object msg, MessageProperties messageProperties) {
+ if (msg instanceof Message) {
+ return (Message) msg;
+ }
+ ChatQueueType chatQueueType = ChatQueueType.getMQType(msg.getClass(),
+ type -> new MessageConversionException(String.format("不支持%s消息类型,在MQType添加", type.getName())));
+ messageProperties.setHeader(MQ_MSG_TYPE, chatQueueType.ordinal());
+ messageProperties.setMessageId(chatQueueType.getMsgId(msg));
+ // 消息序列化成字节
+ byte[] bytes = Jackson.msgpack().serializeBytes(msg);
+ if (bytes.length > Consts.KB) {
+ // snappy 压缩
+ try {
+ long start = ImClock.millis();
+ byte[] compressed = Snappy.compress(bytes);
+ logger.info("snappy compress msg bytes {} to {}, use {} ms.", bytes.length, compressed.length, ImClock.millis() - start);
+ bytes = compressed;
+ messageProperties.setHeader(MQ_MSG_IS_SNAPPY, true);
+ } catch (IOException e) {
+ logger.error("snappy compress error: bytes len {}", bytes.length, e);
+ }
+ }
+ return new Message(bytes, messageProperties);
+ }
+
+ @Override
+ public Object fromMessage(Message message) throws MessageConversionException {
+ MessageProperties messageProperties = message.getMessageProperties();
+ Object ordinal = messageProperties.getHeaders().get(MQ_MSG_TYPE);
+ Object isSnappy = messageProperties.getHeaders().get(MQ_MSG_IS_SNAPPY);
+ ChatQueueType chatQueueType = ChatQueueType.getMQType(Converter.toInt(ordinal),
+ ori -> new MessageConversionException(String.format("没有ordinal为%d的消息类型,在MQType添加", ori)));
+ byte[] body = message.getBody();
+ if (Boolean.TRUE.equals(isSnappy)) {
+ try {
+ body = Snappy.uncompress(body);
+ } catch (IOException e) {
+ throw new MessageConversionException(String.format("消息设置为snappy压缩但解压失败,消息长度 %d bytes", body.length), e);
+ }
+ }
+ return Jackson.msgpack().deserializeBytes(body, chatQueueType.getMsgType());
+ }
+
+}
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
deleted file mode 100644
index 86c0f44..0000000
--- a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMessageConverter.java
+++ /dev/null
@@ -1,39 +0,0 @@
-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
deleted file mode 100644
index 8818df2..0000000
--- a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMsg.java
+++ /dev/null
@@ -1,27 +0,0 @@
-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/ChatQueue.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatQueue.java
new file mode 100644
index 0000000..d74a2d9
--- /dev/null
+++ b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatQueue.java
@@ -0,0 +1,14 @@
+package net.sopod.soim.das.user.api.mq;
+
+/**
+ * ChatQueue
+ *
+ * @author tmy
+ * @date 2022-05-31 00:05
+ */
+public class ChatQueue {
+
+ public static final String IM_MESSAGE_PERSISTENT = "IM_MESSAGE_PERSISTENT";
+
+ public static final String IM_GROUP_MESSAGE_PERSISTENT = "IM_GROUP_MESSAGE_PERSISTENT";
+}
diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatQueueType.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatQueueType.java
new file mode 100644
index 0000000..0454a24
--- /dev/null
+++ b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatQueueType.java
@@ -0,0 +1,73 @@
+package net.sopod.soim.das.user.api.mq;
+
+import net.sopod.soim.common.util.Converter;
+import net.sopod.soim.common.util.StringUtil;
+import net.sopod.soim.das.user.api.model.entity.ImGroupMessage;
+import net.sopod.soim.das.user.api.model.entity.ImMessage;
+
+import java.util.function.Function;
+
+/**
+ * MQType
+ *
+ * @author tmy
+ * @date 2022-05-30 17:43
+ */
+public enum ChatQueueType {
+
+ /** 单聊消息 */
+ IM_MESSAGE(ChatQueue.IM_MESSAGE_PERSISTENT,
+ ImMessage.class,
+ msg -> Converter.toString(((ImMessage)msg).getId())),
+
+ /** 群聊消息 */
+ IM_GROUP_MESSAGE(ChatQueue.IM_GROUP_MESSAGE_PERSISTENT,
+ ImGroupMessage.class,
+ msg -> Converter.toString(((ImGroupMessage)msg).getId()));
+
+ private final Class> msgType;
+
+ private final String queueName;
+
+ private final Function