Browse Source

im-router 同步服务

master
tangmingyou 4 years ago
parent
commit
20f0b80abd
  1. 53
      README.md
  2. 119
      im-common/src/main/java/net/sopod/soim/common/util/Jackson.java
  3. 18
      im-common/src/main/resources/banner.txt
  4. 8
      im-service/im-router/pom.xml
  5. 2
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java
  6. 31
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncService.java
  7. 50
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java
  8. 50
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogDataCodec.java
  9. 27
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java
  10. 33
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServerInboundHandler.java
  11. 50
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/CodecUtil.java
  12. 61
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java
  13. 34
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncLogEncoder.java
  14. 113
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java
  15. 38
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncLog.java
  16. 41
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java
  17. 20
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java
  18. 21
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java
  19. 21
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogServerHandler.java
  20. 56
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncSession.java
  21. 11
      pom.xml

53
README.md

@ -1,4 +1,48 @@
TODO
### 开发计划
- (05-08~05-11)router 一致性 hash,新增、删除、备份处理
- 数据同步删除后,保留id一段时间,有请求进行重定向/转发
- 新增:重算hash,发起数据同步,接收其他节点推送数据及数据更改日志,注册服务,其他节点删除数据
- 正常删除节点:重算 hash,数据和更改日志推送给其他节点,取消注册,关闭服务
- 宕机:备份服务接收全量同步数据和更改日志,服务探测到主服务不可用,注册到注册中心提供服务
- (05-12~05-12) 备用计划:router 层 id 号段负载均衡
- (05-13~05-14) dubbo 服务异步处理
- 功能开发:
- (05-13~05-13) 用户注册功能
- (05-14~05-14) 好友列表(在线状态:批量uid一致性hash, router查询)
- 用户查询
- 添加好友
- 好友在线状态列表
- (05-15~05-15) 消息群发(im-router 群消息路由,批量uid一致性哈希路由)
- 创群:
- 加群:发送申请,推送申请,推送回复
- 删群
- 群列表
- (05-16~05-16) 聊天记录查询
- (05-17~05-18) das 部分接口消息队列异步写
- (05-19~05-20) 集群部署:docker swarm
- 监控:log4j2完善 + prometheus + grafana + loki + arthas
- (05-21~05-21) prometheus 服务发现,exporter(服务发现) 开发
- (05-22~05-22) 容器监控,主机监控,数据库监控,日志收集
- (05-23~05-24) 应用监控
- 通讯指标监控,请求量,吞吐量,延时,请求节点分布,链路追踪
- 数据库连接池监控
- (05-25~05-25) 压测:压测开发
- (05-26~05-26) client: 控制台完善,grallvm 打包
- 后续:
- 通讯加密
- websocket 网关
- web 页面开发
- 异/同设备,多地登录
### TODO
- router -> entry 负载均衡
- router 新增节点顺时针相邻节点数据一致性哈希迁移
- router 冗余节点存储数据不提供服务
@ -22,9 +66,9 @@ TODO
- 内存存储层 router
- 固化存储层 das
功能组件:
1.0
- IOC: guice
- 通信层: netty
- 序列化: protobuf
@ -43,14 +87,13 @@ TODO
- ID生成器
2.0
- quarkus
- quarkus
集群负载均衡:
子网1,子网2,公网
子网1和2不相通
dubbo native image
https://dubbo.apache.org/zh/docs/references/graalvm/support-graalvm/
@ -58,11 +101,9 @@ https://dubbo.apache.org/zh/docs/references/graalvm/support-graalvm/
entry <--> client 通信
- serviceId --> paramClass --> serviceHandler
entry <--> logic dubbo service
client jconsle cmd
router <--> cache
das <--> db
router <--> das

119
im-common/src/main/java/net/sopod/soim/common/util/Jackson.java

@ -8,6 +8,8 @@ import com.fasterxml.jackson.core.json.JsonReadFeature;
import com.fasterxml.jackson.core.json.PackageVersion;
import com.fasterxml.jackson.databind.*;
import com.fasterxml.jackson.databind.module.SimpleModule;
import com.fasterxml.jackson.databind.type.CollectionLikeType;
import com.fasterxml.jackson.databind.type.MapType;
import com.fasterxml.jackson.datatype.jsr310.deser.LocalDateDeserializer;
import com.fasterxml.jackson.datatype.jsr310.deser.LocalDateTimeDeserializer;
import com.fasterxml.jackson.datatype.jsr310.deser.LocalTimeDeserializer;
@ -24,9 +26,7 @@ import java.time.LocalDateTime;
import java.time.LocalTime;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
import java.util.Locale;
import java.util.Map;
import java.util.TimeZone;
import java.util.*;
import java.util.function.Consumer;
import java.util.function.Supplier;
@ -46,6 +46,7 @@ public class Jackson {
private static volatile Jackson JSON_INSTANCE;
private static volatile Jackson YAML_INSTANCE;
private static volatile Jackson XML_INSTANCE;
private static volatile Jackson MSGPACK_INSTANCE;
private final ObjectMapper objectMapper;
@ -176,6 +177,19 @@ public class Jackson {
return XML_INSTANCE;
}
/**
* 序列化的时候不写入字段名字会按字段顺序写入值
* 如果在bean中要增加新字段请务必保证新字段加在字段序的最后
* 对象新增字段放在中间位置会导致序列化失败
*/
public static Jackson msgpack() {
getFactoryInstance("org.msgpack.jackson.dataformat.MessagePackFactory",
() -> MSGPACK_INSTANCE == null,
jackson -> MSGPACK_INSTANCE = jackson
);
return MSGPACK_INSTANCE;
}
public boolean canSerialize(Class<?> clazz) {
return objectMapper.canSerialize(clazz);
}
@ -199,8 +213,107 @@ public class Jackson {
}
}
public <T> byte[] serializeBytes(T value) {
try {
return objectMapper.writeValueAsBytes(value);
} catch (JsonProcessingException e) {
e.printStackTrace();
return null;
}
}
public <T> T deserializeBytes(byte[] bytes, Class<T> valueType) {
try {
return objectMapper.readValue(bytes, valueType);
} catch (IOException e) {
e.printStackTrace();
return null;
}
}
public <T> T toPojo(Map<String, Object> fromValue, Class<T> toValueType) {
return objectMapper.convertValue(fromValue, toValueType);
}
public <K, V> Map<K, V> readMap(String content, Class<K> keyClass, Class<V> valueClass) {
if (content == null || content.length() == 0) {
return Collections.emptyMap();
} else {
try {
return objectMapper.readValue(content, getMapType(keyClass, valueClass));
} catch (IOException e) {
e.printStackTrace();
return null;
}
}
}
public <K, V> Map<K, V> readMap(byte[] content, Class<K> keyClass, Class<V> valueClass) {
if (content == null || content.length == 0) {
return Collections.emptyMap();
} else {
try {
return objectMapper.readValue(content, getMapType(keyClass, valueClass));
} catch (IOException e) {
e.printStackTrace();
return null;
}
}
}
private MapType getMapType(Class<?> keyClass, Class<?> valueClass) {
return objectMapper.getTypeFactory().constructMapType(Map.class, keyClass, valueClass);
}
public List<Map<String, Object>> readListMap(String content) {
return readListMap(content, Object.class);
}
public <V> List<Map<String, V>> readListMap(String content, Class<V> valueType) {
if (content == null || content.length() == 0) {
return Collections.emptyList();
} else {
try {
return objectMapper.readValue(content, objectMapper.getTypeFactory()
.constructCollectionLikeType(List.class, getMapType(String.class, valueType)));
} catch (IOException e) {
e.printStackTrace();
return null;
}
}
}
public <T> List<T> readList(byte[] content, Class<T> elementClass) {
if (content == null || content.length == 0) {
return Collections.emptyList();
} else {
try {
return objectMapper.readValue(content, getListType(elementClass));
} catch (IOException e) {
e.printStackTrace();
return null;
}
}
}
public <T> List<T> readList(String content, Class<T> elementClass) {
if (content == null || content.length() == 0) {
return Collections.emptyList();
} else {
try {
return objectMapper.readValue(content, getListType(elementClass));
} catch (IOException e) {
e.printStackTrace();
return null;
}
}
}
private CollectionLikeType getListType(Class<?> elementClass) {
objectMapper.getTypeFactory().constructCollectionLikeType(List.class, getMapType(String.class, Object.class));
return objectMapper.getTypeFactory().constructCollectionLikeType(List.class, elementClass);
}
}

18
im-common/src/main/resources/banner.txt

@ -0,0 +1,18 @@
⣿⣿⣿⣿⣿⣿⣿⣿⣿⠻⢿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿
⣿⣿⣿⣿⣿⣿⣿⣿⣿⠀⠀⠹⢿⣿⣿⣿⣿⣿⡿⠋⣿⣿⣿⣿⣿
⣿⣿⣿⣿⣿⣿⣿⣿⣿⠀⠀⠀⠈⠻⣿⣿⡿⠏⠀⠀⣿⣿⣿⣿⣿
⣿⣿⣿⣿⣿⣿⣿⣿⣿⠀⠀⠀⠀⠀⠙⠋⠀⠀⠀⢀⣿⣿⣿⣿⣿
⣿⣿⣷⡈⠉⠉⠉⠉⠉⠀⣀⡤⠴⠶⠶⠶⠤⣄⡀⠸⠿⠿⠛⢛⣿
⣿⣿⣿⣿⡄⠀⠀⠀⣰⠚⠁⠀⠀⠀⠀⠀⠀⠀⠳⣄⠀⠀⢠⣾⣿
⣿⣿⣿⣿⣷⠀⠀⡜⠁⠀⠀⠀⠀⣀⣀⣀⣀⣀⣀⠘⡆⢠⣿⣿⣿
⣿⣿⠿⠛⠉⠀⢰⠇⠀⠰⣾⡯⠭⠭⣭⠭⡭⠭⠭⠭⢿⡀⠙⠿⣿
⣿⣤⣀⠀⠀⠀⢸⠀⢠⣧⣤⣤⣤⠤⢼⣾⠥⣤⡤⠤⠬⡇⣠⣴⣾
⣿⣿⣿⣿⡦⠀⢸⠀⠈⢧⡀⠁⠀⠀⡠⠛⣄⠈⠀⢀⣠⠇⣿⣿⣿
⣿⣿⠿⠋⠀⠀⣸⡄⠀⢤⣉⠓⠚⠋⠁⡀⢸⡩⣯⡭⢿⡀⠙⠿⣿
⣿⣷⣤⣄⡀⢼⠋⠀⠀⠀⠉⠉⠁⠀⠀⠳⣼⠇⠀⠀⣼⢹⣾⣿⣿
⣿⣿⣿⣿⡟⠈⠳⣤⠀⠀⡀⠀⠀⣀⣀⣀⣀⣀⡀⠀⣿⡻⣿⣿⣿
⣿⣿⣿⣿⣾⣿⣿⠞⣆⠸⡇⠚⠉⠁⠀⣿⣿⡟⠉⢁⣿⣿⣿⣿⣿
⣿⣿⣿⣿⣿⣿⣿⣶⣿⣆⠈⠁⠀⠰⡖⠙⠛⠃⣠⣿⣿⣿⣿⣿⣿
⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⡦⣄⣀⣀⣠⣴⣾⣿⣿⣿⣿⣿⣿⣿
⣿⣿⣿⣿⣿⣿⡿⠿⠟⠛⡛⢷⣤⣀⣠⢾⣛⠛⠿⢿⣿⣿⣿⣿⣿
⣿⣿⣿⣿⣿⠃⠀⠀⠀⢰⠃⢸⠀⠀⠀⠘⡟⣆⠀⠀⢹⣿⣿⣿⣿

8
im-service/im-router/pom.xml

@ -85,6 +85,14 @@
<groupId>cglib</groupId>
<artifactId>cglib</artifactId>
</dependency>
<dependency>
<groupId>org.msgpack</groupId>
<artifactId>jackson-dataformat-msgpack</artifactId>
</dependency>
<dependency>
<groupId>org.xerial.snappy</groupId>
<artifactId>snappy-java</artifactId>
</dependency>
</dependencies>
</project>

2
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java

@ -1,6 +1,6 @@
package net.sopod.soim.router.datasync;
import net.sopod.soim.router.datasync.server.SyncLog;
import net.sopod.soim.router.datasync.server.data.SyncLog;
import java.util.Queue;
import java.util.concurrent.ConcurrentHashMap;

31
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncService.java

@ -0,0 +1,31 @@
package net.sopod.soim.router.datasync;
import net.sopod.soim.common.util.Jackson;
import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.cache.RouterUserStorage;
import java.util.Map;
/**
* SyncService
*
* @author tmy
* @date 2022-05-08 09:42
*/
public class SyncService {
public void fullSync() {
Map<Long, RouterUser> routerUserMap = RouterUserStorage.getInstance().getRouterUserMap();
for (Map.Entry<Long, RouterUser> entry : routerUserMap.entrySet()) {
Long key = entry.getKey();
RouterUser value = entry.getValue();
Jackson.json().serialize(value);
}
}
public static void main(String[] args) {
}
}

50
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java

@ -1,9 +1,17 @@
package net.sopod.soim.router.datasync.server;
import io.netty.bootstrap.Bootstrap;
import io.netty.channel.Channel;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.util.AttributeKey;
import net.sopod.soim.router.datasync.server.codec.SyncCmdCodec;
import net.sopod.soim.router.datasync.server.codec.SyncLogEncoder;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import net.sopod.soim.router.datasync.server.data.SyncLog;
import net.sopod.soim.router.datasync.server.handler.SyncCmdClientHandler;
import net.sopod.soim.router.datasync.server.handler.SyncLogClientHandler;
/**
* SyncClient
@ -13,18 +21,52 @@ import io.netty.channel.socket.SocketChannel;
*/
public class SyncClient {
private NioEventLoopGroup group;
private final Bootstrap bootstrap;
private Channel clientChannel;
public SyncClient() {
NioEventLoopGroup group = new NioEventLoopGroup(2);
new Bootstrap()
.group(group)
this.bootstrap = new Bootstrap()
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel channel) throws Exception {
channel.pipeline()
.addLast(new SyncLogDataCodec());
.addLast(new SyncCmdCodec())
.addLast(new SyncLogEncoder())
.addLast(new SyncCmdClientHandler())
.addLast(new SyncLogClientHandler());
}
});
}
public void connect(String host, int port) throws InterruptedException {
this.group = new NioEventLoopGroup(2);
this.clientChannel = bootstrap
.group(group)
.connect(host, port)
.await().channel();
}
public void write(SyncCmd syncCmd) {
clientChannel.writeAndFlush(syncCmd);
}
public void write(SyncLog syncLog) {
clientChannel.write(syncLog);
}
public <T> void setAttr(AttributeKey<T> key, T val) {
clientChannel.attr(key).set(val);
}
public <T> void getAttr(AttributeKey<T> key) {
clientChannel.attr(key).get();
}
public void close() {
this.group.shutdownGracefully();
}
}

50
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogDataCodec.java

@ -1,50 +0,0 @@
package net.sopod.soim.router.datasync.server;
import io.netty.buffer.ByteBuf;
import io.netty.channel.*;
import io.netty.handler.codec.ByteToMessageDecoder;
import io.netty.handler.codec.MessageToByteEncoder;
import java.util.List;
/**
* SyncDataInboundHandler
*
* @author tmy
* @date 2022-05-05 10:20
*/
public class SyncLogDataCodec
extends CombinedChannelDuplexHandler<SyncLogDataCodec.SyncDataDecoder,
SyncLogDataCodec.SyncDataEncoder> {
public SyncLogDataCodec() {
super(new SyncDataDecoder(), new SyncDataEncoder());
}
/**
* 解码器
*/
public static class SyncDataDecoder extends ByteToMessageDecoder {
@Override
protected void decode(ChannelHandlerContext channelHandlerContext, ByteBuf byteBuf, List<Object> list) throws Exception {
SyncLog syncLog = SyncLog.read(byteBuf);
list.add(syncLog);
}
}
/**
* 编码器
*/
public static class SyncDataEncoder extends MessageToByteEncoder<SyncLog> {
@Override
protected void encode(ChannelHandlerContext channelHandlerContext, SyncLog syncLog, ByteBuf byteBuf) throws Exception {
byte[] bytes = syncLog.toBytes();
byteBuf.writeBytes(bytes);
}
}
}

27
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java

@ -10,6 +10,10 @@ import io.netty.channel.socket.nio.NioServerSocketChannel;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
import net.sopod.soim.common.util.netty.FastThreadLocalThreadFactory;
import net.sopod.soim.router.datasync.server.codec.SyncCmdCodec;
import net.sopod.soim.router.datasync.server.codec.SyncLogEncoder;
import net.sopod.soim.router.datasync.server.handler.SyncCmdServerHandler;
import net.sopod.soim.router.datasync.server.handler.SyncLogServerHandler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -25,16 +29,14 @@ public class SyncServer {
private static final Logger logger = LoggerFactory.getLogger(SyncServer.class);
private final int port;
private NioEventLoopGroup boss;
private NioEventLoopGroup worker;
public SyncServer(int port) {
this.port = port;
public SyncServer() {
}
public void start(Consumer<Throwable> onFail) throws InterruptedException {
public void start(int port, Consumer<Throwable> onFail) {
this.boss = new NioEventLoopGroup(1, new FastThreadLocalThreadFactory("sync-server-boss-%d", Thread.NORM_PRIORITY));
this.worker = new NioEventLoopGroup(2, new FastThreadLocalThreadFactory("sync-server-worker-%d", Thread.NORM_PRIORITY));
ServerBootstrap serverBoot = new ServerBootstrap()
@ -49,8 +51,10 @@ public class SyncServer {
: logger.isErrorEnabled() ? LogLevel.ERROR : LogLevel.INFO;
ChannelPipeline pipeline = channel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel))
.addLast(new SyncLogDataCodec())
.addLast(new SyncServerInboundHandler());
.addLast(new SyncCmdCodec())
.addLast(new SyncLogEncoder())
.addLast(new SyncCmdServerHandler())
.addLast(new SyncLogServerHandler());
}
});
serverBoot.bind(port).addListener((ChannelFutureListener) future -> {
@ -73,4 +77,13 @@ public class SyncServer {
}
}
public static void main(String[] args) {
new SyncServer()
.start(9999, err -> {
logger.error("服务启动失败");
});
}
}

33
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServerInboundHandler.java

@ -1,33 +0,0 @@
package net.sopod.soim.router.datasync.server;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
/**
* SyncLogInboundHandler
*
* @author tmy
* @date 2022-05-05 14:57
*/
public class SyncServerInboundHandler extends ChannelInboundHandlerAdapter {
public SyncServerInboundHandler() {
}
@Override
public void channelActive(ChannelHandlerContext ctx) throws Exception {
super.channelActive(ctx);
}
@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception {
super.channelInactive(ctx);
}
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
super.exceptionCaught(ctx, cause);
}
}

50
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/CodecUtil.java

@ -0,0 +1,50 @@
package net.sopod.soim.router.datasync.server.codec;
import net.sopod.soim.common.util.Jackson;
import org.xerial.snappy.Snappy;
import java.io.IOException;
/**
* CodecUtil
*
* @author tmy
* @date 2022-05-08 16:54
*/
public class CodecUtil {
public static byte[] encode(Object data) {
return Jackson.msgpack().serializeBytes(data);
}
public static byte[] compress(byte[] bytes) {
try {
return Snappy.compress(bytes);
} catch (IOException e) {
throw new RuntimeException(e);
}
}
public static byte[] codecAndCompress(Object data) {
byte[] bytes = encode(data);
return compress(bytes);
}
public static <T> T decode(byte[] bytes, Class<T> type) {
return Jackson.msgpack().deserializeBytes(bytes, type);
}
public static byte[] uncompress(byte[] bytes) {
try {
return Snappy.uncompress(bytes);
} catch (IOException e) {
throw new RuntimeException(e);
}
}
public static <T> T uncompressAndDecode(byte[] bytes, Class<T> type) {
byte[] uncompress = uncompress(bytes);
return decode(uncompress, type);
}
}

61
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java

@ -0,0 +1,61 @@
package net.sopod.soim.router.datasync.server.codec;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.CombinedChannelDuplexHandler;
import io.netty.handler.codec.ByteToMessageDecoder;
import io.netty.handler.codec.MessageToByteEncoder;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import net.sopod.soim.router.datasync.server.data.SyncLog;
import java.util.List;
/**
* SyncChannelHandler
* SyncCmd 编解码器
*
* @author tmy
* @date 2022-05-07 21:37
*/
public class SyncCmdCodec extends CombinedChannelDuplexHandler<SyncCmdCodec.SyncCmdDecoder,
SyncCmdCodec.SyncCmdEncoder> {
public SyncCmdCodec() {
super(new SyncCmdDecoder(), new SyncCmdEncoder());
}
public static class SyncCmdDecoder extends ByteToMessageDecoder {
@Override
protected void decode(ChannelHandlerContext channelHandlerContext, ByteBuf byteBuf, List<Object> list) throws Exception {
SyncCmd syncCmd = SyncCmd.read(byteBuf);
if (!SyncCmd.SYNC_LOG.equals(syncCmd.getCmdType())) {
list.add(syncCmd);
return;
}
// 数据同步请求命令
// 同步数据请求体长度
int syncLogBytesLen = byteBuf.readInt();
// 读取、解压、解析字节数据
byte[] syncLogCompressBytes = new byte[syncLogBytesLen];
byteBuf.readBytes(syncLogCompressBytes);
byte[] syncLogBytes = CodecUtil.uncompress(syncLogCompressBytes);
ByteBuf syncLogByteBuf = Unpooled.wrappedBuffer(syncLogBytes);
try {
SyncLog syncLog = SyncLog.read(syncLogByteBuf);
list.add(syncLog);
} finally {
syncLogByteBuf.release();
}
}
}
public static class SyncCmdEncoder extends MessageToByteEncoder<SyncCmd> {
@Override
protected void encode(ChannelHandlerContext channelHandlerContext, SyncCmd syncCmd, ByteBuf byteBuf) throws Exception {
byte[] bytes = syncCmd.toByte();
byteBuf.writeBytes(bytes);
}
}
}

34
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncLogEncoder.java

@ -0,0 +1,34 @@
package net.sopod.soim.router.datasync.server.codec;
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.MessageToByteEncoder;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import net.sopod.soim.router.datasync.server.data.SyncLog;
/**
* SyncDataInboundHandler
* SyncLog 编码器
*
* @author tmy
* @date 2022-05-05 10:20
*/
public class SyncLogEncoder extends MessageToByteEncoder<SyncLog> {
@Override
protected void encode(ChannelHandlerContext channelHandlerContext, SyncLog syncLog, ByteBuf byteBuf) throws Exception {
// 同步数据指令
SyncCmd syncCmd = new SyncCmd()
.setCmdType(SyncCmd.SYNC_LOG);
byte[] syncCmdBytes = syncCmd.toByte();
byteBuf.writeBytes(syncCmdBytes);
// 同步数据内容
byte[] bytes = syncLog.toBytes();
// snappy 压缩一下
byte[] compressLogBytes = CodecUtil.compress(bytes);
byteBuf.writeInt(compressLogBytes.length);
byteBuf.writeBytes(compressLogBytes);
}
}

113
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java

@ -0,0 +1,113 @@
package net.sopod.soim.router.datasync.server.data;
import io.netty.buffer.ByteBuf;
import lombok.Data;
import lombok.experimental.Accessors;
import org.apache.dubbo.common.io.Bytes;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* SyncCmd
*
* @author tmy
* @date 2022-05-07 21:41
*/
@Data
@Accessors(chain = true)
public class SyncCmd {
private static final short MAGIC = 0x7a21;
public static final int PING = 1;
public static final int PONG = 2;
/**
* 同步命令
*/
public static final Integer SYNC = 3;
/**
* 同步结束命令
*/
public static final Integer SYNC_END = 4;
/**
* SyncLog 推送命令
*/
public static final Integer SYNC_LOG = 10;
private Integer cmdType;
private String param1;
private String param2;
private String param3;
public byte[] toByte() {
byte[] param1Bytes = param1 == null ? new byte[0] : param1.getBytes(StandardCharsets.UTF_8);
byte[] param2Bytes = param2 == null ? new byte[0] : param2.getBytes(StandardCharsets.UTF_8);
byte[] param3Bytes = param3 == null ? new byte[0] : param3.getBytes(StandardCharsets.UTF_8);
byte[] bytes = new byte[2 + 4 + 4 + param1Bytes.length + 4 + param2Bytes.length + 4 + param3Bytes.length];
int offset = 0;
Bytes.short2bytes(MAGIC, bytes, offset);
offset += 2;
Bytes.int2bytes(cmdType, bytes, offset);
offset += 4;
Bytes.int2bytes(param1Bytes.length, bytes, offset);
offset += 4;
System.arraycopy(param1Bytes, 0, bytes, offset, param1Bytes.length);
offset += param1Bytes.length;
Bytes.int2bytes(param2Bytes.length, bytes, offset);
offset += 4;
System.arraycopy(param2Bytes, 0, bytes, offset, param2Bytes.length);
offset += param2Bytes.length;
Bytes.int2bytes(param3Bytes.length, bytes, offset);
offset += 4;
System.arraycopy(param3Bytes, 0, bytes, offset, param3Bytes.length);
offset += param3Bytes.length;
return bytes;
}
public static SyncCmd read(ByteBuf buf) {
short magic = buf.readShort();
if (magic != MAGIC) {
throw new IllegalStateException("unknown bytes magic error");
}
SyncCmd cmd = new SyncCmd();
cmd.cmdType = buf.readInt();
int param1Len = buf.readInt();
byte[] paramBytes = new byte[param1Len];
if (param1Len > 0) {
buf.readBytes(paramBytes);
cmd.param1 = new String(paramBytes, StandardCharsets.UTF_8);
}
int param2Len = buf.readInt();
if (param2Len > 0) {
paramBytes = paramBytes.length >= param2Len ? paramBytes : new byte[param2Len];
buf.readBytes(paramBytes, 0, param2Len);
cmd.param2 = new String(paramBytes, 0, param2Len, StandardCharsets.UTF_8);
}
int param3Len = buf.readInt();
if (param3Len > 0) {
paramBytes = paramBytes.length >= param3Len ? paramBytes : new byte[param3Len];
buf.readBytes(paramBytes, 0, param3Len);
cmd.param3 = new String(paramBytes, 0, param3Len, StandardCharsets.UTF_8);
}
return cmd;
}
public static void main(String[] args) {
AtomicBoolean locked = new AtomicBoolean(false);
System.out.println(locked.compareAndSet(false, false));
System.out.println(locked.compareAndExchange(false, true));
}
}

38
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLog.java → im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncLog.java

@ -1,4 +1,4 @@
package net.sopod.soim.router.datasync.server;
package net.sopod.soim.router.datasync.server.data;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
@ -12,8 +12,10 @@ import net.sopod.soim.router.datasync.SyncTypes;
import org.apache.dubbo.common.io.Bytes;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.xerial.snappy.Snappy;
import javax.annotation.Nullable;
import java.io.IOException;
import java.io.Serializable;
import java.lang.reflect.Method;
import java.nio.charset.StandardCharsets;
@ -37,7 +39,7 @@ public class SyncLog implements Serializable {
private static final long serialVersionUID = -8925525709590423526L;
private static final short MAGIC = 0x7a21;
// private static final short MAGIC = 0x7a21;
public static final int OPT_ADD = 1;
public static final int OPT_REMOVE = 2;
@ -87,7 +89,7 @@ public class SyncLog implements Serializable {
protected String[] args;
/** ================ 新增数据:序列化后的数据(避免修改) ===================== */
protected List<String> serializeDataCollect;
protected List<byte[]> serializeDataCollect;
/** 数据可能同步给多个订阅者,缓存一下 */
private transient byte[] toBytesCache;
@ -125,8 +127,9 @@ public class SyncLog implements Serializable {
int dataCollectByteLen = 0;
if (dataCollectSize > 0) {
int i = 0;
for (String dataCollect : serializeDataCollect) {
byteData[i] = dataCollect.getBytes(StandardCharsets.UTF_8);
for (byte[] dataCollect : serializeDataCollect) {
// byteData[i] = dataCollect.getBytes(StandardCharsets.UTF_8);
byteData[i] = dataCollect;
dataCollectByteLen += 4;
dataCollectByteLen += byteData[i].length;
i++;
@ -134,8 +137,8 @@ public class SyncLog implements Serializable {
}
// 总长 + 同步数据类型(byte) +
byte[] bytes = new byte[
2 // 魔术
+ 4 // 请求体总长度
// 2 + // 魔术
4 // 请求体总长度
+ 1 // 同步数据类型 SyncType
+ 1 // 新增/删除/更新
+ 4 // logSeq 序列号
@ -151,8 +154,8 @@ public class SyncLog implements Serializable {
+ dataCollectByteLen // 序列化数据字节(len,dataBytes,len,dataBytes...)
];
int offset = 0;
Bytes.short2bytes(MAGIC, bytes, offset);
offset += 2;
// Bytes.short2bytes(MAGIC, bytes, offset);
// offset += 2;
Bytes.int2bytes(bytes.length - 6, bytes, offset);
offset += 4;
@ -204,10 +207,6 @@ public class SyncLog implements Serializable {
}
public static SyncLog read(ByteBuf buf) {
short magic = buf.readShort();
if (magic != MAGIC) {
throw new IllegalStateException("unknown bytes magic error");
}
// 后续bytes长度
int dataLen = buf.readInt();
int syncDataType = buf.readByte();
@ -258,18 +257,17 @@ public class SyncLog implements Serializable {
int dataSize = buf.readInt();
log.serializeDataCollect = dataSize > 0 ? new ArrayList<>(dataSize) : Collections.emptyList();
if (dataSize > 0) {
byte[] dataByte = new byte[0];
for (int i = 0; i < dataSize; i++) {
int dataByteLen = buf.readInt();
dataByte = dataByte.length >= dataByteLen ? dataByte : new byte[dataByteLen];
byte[] dataByte = new byte[dataByteLen];
buf.readBytes(dataByte, 0, dataByteLen);
log.serializeDataCollect.add(new String(dataByte, 0, dataByteLen, StandardCharsets.UTF_8));
log.serializeDataCollect.add(dataByte);
}
}
return log;
}
public static void main(String[] args) {
public static void main(String[] args) throws IOException {
RouterUser user1 = new RouterUser()
.setUid(10001L)
.setAccount("前线")
@ -288,6 +286,7 @@ public class SyncLog implements Serializable {
// System.out.println(json.getBytes(StandardCharsets.UTF_8).length);
byte[] bytes = addLog.toBytes();
byte[] compress = Snappy.compress(bytes);
ByteBuf buf = Unpooled.wrappedBuffer(bytes);
SyncLog newLog = SyncLog.read(buf);
buf.release();
@ -328,7 +327,7 @@ public class SyncLog implements Serializable {
this.logSeq = logSeq;
}
public AddLog<T> setSerializeDataCollect(List<String> serializeDataCollect) {
public AddLog<T> setSerializeDataCollect(List<byte[]> serializeDataCollect) {
this.serializeDataCollect = serializeDataCollect;
return this;
}
@ -337,7 +336,8 @@ public class SyncLog implements Serializable {
if (this.serializeDataCollect == null) {
this.serializeDataCollect = new ArrayList<>();
}
this.serializeDataCollect.add(Jackson.json().serialize(data));
byte[] dataBytes = Jackson.msgpack().serializeBytes(data);
this.serializeDataCollect.add(dataBytes);
return this;
}
}

41
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java

@ -0,0 +1,41 @@
package net.sopod.soim.router.datasync.server.handler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* SyncHandler
*
* @author tmy
* @date 2022-05-07 21:43
*/
public class SyncCmdClientHandler extends SimpleChannelInboundHandler<SyncCmd> {
private static final Logger logger = LoggerFactory.getLogger(SyncCmdClientHandler.class);
@Override
protected void channelRead0(ChannelHandlerContext ctx, SyncCmd syncCmd) {
switch (syncCmd.getCmdType()) {
case SyncCmd.PING:
this.handlePing(ctx, syncCmd);
break;
case SyncCmd.PONG:
this.handlePong(ctx, syncCmd);
break;
}
}
private void handlePing(ChannelHandlerContext ctx, SyncCmd syncCmd) {
logger.info("ping: {}", ctx.channel());
SyncCmd pong = new SyncCmd().setCmdType(SyncCmd.PONG);
ctx.writeAndFlush(pong);
}
private void handlePong(ChannelHandlerContext ctx, SyncCmd syncCmd) {
logger.info("pong: {}", ctx.channel());
}
}

20
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java

@ -0,0 +1,20 @@
package net.sopod.soim.router.datasync.server.handler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
/**
* SyncCmdServerHandler
*
* @author tmy
* @date 2022-05-08 18:48
*/
public class SyncCmdServerHandler extends SimpleChannelInboundHandler<SyncCmd> {
@Override
protected void channelRead0(ChannelHandlerContext channelHandlerContext, SyncCmd syncCmd) throws Exception {
}
}

21
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java

@ -0,0 +1,21 @@
package net.sopod.soim.router.datasync.server.handler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import net.sopod.soim.router.datasync.server.data.SyncLog;
/**
* SyncLogHandler
* 其他节点会推送同步数据到当前新增的节点这里接收数据
*
* @author tmy
* @date 2022-05-08 17:20
*/
public class SyncLogClientHandler extends SimpleChannelInboundHandler<SyncLog> {
@Override
protected void channelRead0(ChannelHandlerContext channelHandlerContext, SyncLog syncLog) throws Exception {
}
}

21
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogServerHandler.java

@ -0,0 +1,21 @@
package net.sopod.soim.router.datasync.server.handler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import net.sopod.soim.router.datasync.server.data.SyncLog;
/**
* SyncLogServerHandler
* 即将删除的节点会主动推送同步数据这里进行接收
*
* @author tmy
* @date 2022-05-08 18:48
*/
public class SyncLogServerHandler extends SimpleChannelInboundHandler<SyncLog> {
@Override
protected void channelRead0(ChannelHandlerContext channelHandlerContext, SyncLog syncLog) throws Exception {
}
}

56
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncSession.java

@ -0,0 +1,56 @@
package net.sopod.soim.router.datasync.server.session;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import net.sopod.soim.router.datasync.server.SyncServer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* SyncSession
* 作为客户端或服务端主动操作连接
*
* @author tmy
* @date 2022-05-08 17:34
*/
public class SyncSession {
private static final Logger logger = LoggerFactory.getLogger(SyncSession.class);
public static SyncSession client() {
return null;
}
public static SyncSession server() {
return null;
}
public static class SyncSessionServer {
private final SyncServer syncServer;
public SyncSessionServer(int port) {
this.syncServer = new SyncServer();
this.syncServer.start(9999, err -> {
logger.error("sync server start fail!", err);
});
}
public void writeTo() {
}
}
public static class SyncSessionClient {
public SyncSessionClient() {
}
public void write(SyncCmd syncCmd) {
}
}
}

11
pom.xml

@ -41,6 +41,7 @@
<curator5_version>4.2.0</curator5_version>
<shardingsphere.version>5.1.0</shardingsphere.version>
<nacos.version>2.0.3</nacos.version>
<msgpack.version>0.9.1</msgpack.version>
</properties>
<dependencyManagement>
@ -254,6 +255,16 @@
<artifactId>cglib</artifactId>
<version>3.3.0</version>
</dependency>
<dependency>
<groupId>org.msgpack</groupId>
<artifactId>jackson-dataformat-msgpack</artifactId>
<version>${msgpack.version}</version>
</dependency>
<dependency>
<groupId>org.xerial.snappy</groupId>
<artifactId>snappy-java</artifactId>
<version>1.1.8.4</version>
</dependency>
</dependencies>
</dependencyManagement>

Loading…
Cancel
Save