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; + } + }