Browse Source

mq消息序列化

master
tangmingyou 4 years ago
parent
commit
21a33b5f74
  1. 11
      im-das-api/im-das-user-api/pom.xml
  2. 15
      im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQ.java
  3. 39
      im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMessageConverter.java
  4. 27
      im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMsg.java
  5. 47
      im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/MQType.java
  6. 4
      im-das/im-das-user/pom.xml
  7. 16
      im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Receiver.java
  8. 22
      im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Sender.java
  9. 8
      im-das/im-das-user/src/main/java/net/sopod/soim/das/user/config/RabbitConfig.java

11
im-das-api/im-das-user-api/pom.xml

@ -17,6 +17,17 @@
<artifactId>mybatis-plus-annotation</artifactId> <artifactId>mybatis-plus-annotation</artifactId>
<version>3.5.0</version> <version>3.5.0</version>
</dependency> </dependency>
<dependency>
<groupId>net.sopod</groupId>
<artifactId>im-common</artifactId>
<version>${soim.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
<scope>provided</scope>
</dependency>
</dependencies> </dependencies>
</project> </project>

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

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

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

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

4
im-das/im-das-user/pom.xml

@ -88,6 +88,10 @@
<groupId>org.springframework.boot</groupId> <groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId> <artifactId>spring-boot-starter-amqp</artifactId>
</dependency> </dependency>
<dependency>
<groupId>org.msgpack</groupId>
<artifactId>jackson-dataformat-msgpack</artifactId>
</dependency>
<dependency> <dependency>
<groupId>org.springframework.boot</groupId> <groupId>org.springframework.boot</groupId>

16
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; 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.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
/** /**
* Listener * Listener
* 只能接受字符串和字节数组
* *
* @author tmy * @author tmy
* @date 2022-05-28 16:26 * @date 2022-05-28 16:26
@ -15,8 +19,16 @@ import org.springframework.stereotype.Component;
public class Receiver { public class Receiver {
@RabbitHandler @RabbitHandler
public void process(String hello) { public void process(Message message) {
System.out.println("Receiver: " + hello); 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);
} }
} }

22
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; package net.sopod.soim.das.user.amqp;
import net.sopod.soim.common.util.ImClock; 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.AmqpTemplate;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.boot.context.event.ApplicationReadyEvent; import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.ApplicationListener; import org.springframework.context.ApplicationListener;
import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import java.util.UUID;
/** /**
* Sender * Sender
* *
@ -24,11 +22,17 @@ public class Sender implements ApplicationListener<ApplicationReadyEvent> {
public void onApplicationEvent(ApplicationReadyEvent event) { public void onApplicationEvent(ApplicationReadyEvent event) {
ConfigurableApplicationContext context = event.getApplicationContext(); ConfigurableApplicationContext context = event.getApplicationContext();
AmqpTemplate amqpTemplate = context.getBean(AmqpTemplate.class); AmqpTemplate amqpTemplate = context.getBean(AmqpTemplate.class);
amqpTemplate.convertAndSend("hello", "hello:" + ImClock.millis()); // amqpTemplate.convertAndSend("hello", "hello:" + ImClock.millis());
MessageProperties messageProperties = new MessageProperties(); ImMessage imMessage = new ImMessage()
messageProperties.setMessageId(UUID.randomUUID().toString()); .setId(10086L)
Message message = new Message(new byte[] {}, messageProperties); .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...."); System.out.println("send msg....");
} }

8
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; package net.sopod.soim.das.user.config;
import org.springframework.amqp.core.Queue; import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Configuration;
@ -19,4 +20,11 @@ public class RabbitConfig {
return new Queue("hello"); return new Queue("hello");
} }
@Bean
public RabbitTemplate rabbitTemplate() {
RabbitTemplate rabbitTemplate = new RabbitTemplate();
rabbitTemplate.setMessageConverter(null);
return null;
}
} }

Loading…
Cancel
Save