From 078f0e0b408f86e4de2cc06e0231ab1ba34d816d Mon Sep 17 00:00:00 2001
From: tangmingyou <234767776@qq.com>
Date: Wed, 11 May 2022 17:58:54 +0800
Subject: [PATCH] =?UTF-8?q?im-router=E6=95=B0=E6=8D=AE=E6=9B=B4=E6=94=B9?=
=?UTF-8?q?=E6=97=A5=E5=BF=97=E5=90=8C=E6=AD=A5?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
.../soim/client/session/SoImSession.java | 2 +
im-common/pom.xml | 6 ++
.../soim/common/util/HashAlgorithms.java | 19 ----
.../net/sopod/soim/common/util/Jackson.java | 73 ++++++++++++-
.../net/sopod/soim/common/util/Reflects.java | 29 +++++
.../common/util/cache/LinkedMapLRUCache.java | 35 ++++++
.../common/util/netty/Varint32FrameCodec.java | 28 +++++
.../soim/entry/server/ImEntryInitializer.java | 2 +
.../router/api/route/ConsistentHashTest.java | 4 -
.../api/route/UidConsistentHashSelector.java | 24 ++++-
.../soim/router/config/AppContextHolder.java | 32 +++++-
.../router/config/ImRouterAppOnReady.java | 45 ++++----
.../router/datasync/DataChangeTrigger.java | 100 ++++++++++++++++--
.../router/datasync/DataSyncProxyFactory.java | 80 +++++++-------
.../{server => }/SyncLogByHashService.java | 59 +++++++----
.../datasync/SyncLogMigrateService.java | 48 +++++++--
.../router/datasync/server/SyncClient.java | 4 +
.../router/datasync/server/SyncServer.java | 4 +
.../datasync/server/codec/SyncCmdCodec.java | 1 +
.../router/datasync/server/data/SyncCmd.java | 2 +-
.../router/datasync/server/data/SyncLog.java | 15 +--
.../server/handler/SyncCmdClientHandler.java | 11 ++
.../server/handler/SyncCmdServerHandler.java | 2 +-
.../server/handler/SyncLogClientHandler.java | 81 ++++++++++++--
24 files changed, 551 insertions(+), 155 deletions(-)
create mode 100644 im-common/src/main/java/net/sopod/soim/common/util/cache/LinkedMapLRUCache.java
create mode 100644 im-common/src/main/java/net/sopod/soim/common/util/netty/Varint32FrameCodec.java
rename im-service/im-router/src/main/java/net/sopod/soim/router/datasync/{server => }/SyncLogByHashService.java (85%)
diff --git a/im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java b/im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java
index 6051eab..af66dac 100644
--- a/im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java
+++ b/im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java
@@ -10,6 +10,7 @@ import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;
import net.sopod.soim.client.logger.Logger;
+import net.sopod.soim.common.util.netty.Varint32FrameCodec;
import net.sopod.soim.core.net.ImMessageCodec;
import net.sopod.soim.data.msg.auth.Auth;
import net.sopod.soim.data.msg.chat.Chat;
@@ -55,6 +56,7 @@ public class SoImSession {
@Override
protected void initChannel(SocketChannel ch) throws Exception {
ch.pipeline()
+ .addLast(new Varint32FrameCodec())
.addLast(new ImMessageCodec())
.addLast(messageDispatcher);
}
diff --git a/im-common/pom.xml b/im-common/pom.xml
index 03ef45a..b4deffc 100644
--- a/im-common/pom.xml
+++ b/im-common/pom.xml
@@ -36,6 +36,12 @@
${netty.version}
provided
+
+ io.netty
+ netty-codec
+ ${netty.version}
+ provided
+
\ No newline at end of file
diff --git a/im-common/src/main/java/net/sopod/soim/common/util/HashAlgorithms.java b/im-common/src/main/java/net/sopod/soim/common/util/HashAlgorithms.java
index c2760cc..4cc7f3e 100644
--- a/im-common/src/main/java/net/sopod/soim/common/util/HashAlgorithms.java
+++ b/im-common/src/main/java/net/sopod/soim/common/util/HashAlgorithms.java
@@ -35,25 +35,6 @@ public class HashAlgorithms {
return hashCode & 0xffffffffL;
}
- public static long md5Hash(String value, int number) {
- MessageDigest md5;
- try {
- md5 = MessageDigest.getInstance("MD5");
- } catch (NoSuchAlgorithmException e) {
- throw new RuntimeException("MD5 not supported", e);
- }
- md5.reset();
- byte[] keyBytes = value.getBytes(StandardCharsets.UTF_8);
- md5.update(keyBytes);
- byte[] digest = md5.digest();
- // dubbo md5 hash algorithms
- return (((long) (digest[3 + number * 4] & 0xFF) << 24)
- | ((long) (digest[2 + number * 4] & 0xFF) << 16)
- | ((long) (digest[1 + number * 4] & 0xFF) << 8)
- | (digest[number * 4] & 0xFF))
- & 0xFFFFFFFFL;
- }
-
/**
* CRC系列算法本身并非查表,可是,查表是它的一种最快的实现方式。以下是CRC32的实现
diff --git a/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java b/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java
index b1fb6ea..acd57d1 100644
--- a/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java
+++ b/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java
@@ -10,6 +10,7 @@ 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.databind.type.TypeFactory;
import com.fasterxml.jackson.datatype.jsr310.deser.LocalDateDeserializer;
import com.fasterxml.jackson.datatype.jsr310.deser.LocalDateTimeDeserializer;
import com.fasterxml.jackson.datatype.jsr310.deser.LocalTimeDeserializer;
@@ -20,6 +21,8 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
+import java.lang.reflect.ParameterizedType;
+import java.lang.reflect.Type;
import java.text.SimpleDateFormat;
import java.time.LocalDate;
import java.time.LocalDateTime;
@@ -312,8 +315,76 @@ public class Jackson {
}
private CollectionLikeType getListType(Class> elementClass) {
- objectMapper.getTypeFactory().constructCollectionLikeType(List.class, getMapType(String.class, Object.class));
return objectMapper.getTypeFactory().constructCollectionLikeType(List.class, elementClass);
}
+ public ObjectMapper getObjectMapper() {
+ return objectMapper;
+ }
+
+ public TypeFactory getTypeFactory() {
+ return objectMapper.getTypeFactory();
+ }
+
+ /**
+ * @param content 序列化内容
+ * @param type 方法参数类型等可能包含泛型类型
+ * @return 解析结果
+ */
+ public Object readWithType(String content, Type type) {
+ if (content == null || content.length() == 0) {
+ return null;
+ } else {
+ JavaType javaType = getJavaType(type);
+ try {
+ return objectMapper.readValue(content, javaType);
+ } catch (JsonProcessingException e) {
+ e.printStackTrace();
+ return null;
+ }
+ }
+ }
+
+ public Object readWithType(byte[] bytes, Type type) {
+ if (bytes == null || bytes.length == 0) {
+ return null;
+ } else {
+ JavaType javaType = getJavaType(type);
+ try {
+ return objectMapper.readValue(bytes, javaType);
+ } catch (IOException e) {
+ e.printStackTrace();
+ return null;
+ }
+ }
+ }
+
+ /**
+ * 解析任意层数泛型对象
+ * @param type 泛型类型
+ * @return jackson JavaType
+ */
+ private JavaType getJavaType(Type type) {
+ //判断是否带有泛型
+ if (type instanceof ParameterizedType) {
+ Type[] actualTypeArguments = ((ParameterizedType) type).getActualTypeArguments();
+ System.out.println(Arrays.asList(actualTypeArguments));
+ //获取泛型类型
+ Class> rowClass = (Class>) ((ParameterizedType) type).getRawType();
+
+ JavaType[] javaTypes = new JavaType[actualTypeArguments.length];
+
+ for (int i = 0; i < actualTypeArguments.length; i++) {
+ //泛型也可能带有泛型,递归获取
+ javaTypes[i] = getJavaType(actualTypeArguments[i]);
+ }
+ return objectMapper.getTypeFactory().constructParametricType(rowClass, javaTypes);
+ } else {
+ //简单类型直接用该类构建JavaType
+ Class> cla = (Class>) type;
+ // TypeFactory.defaultInstance()
+ return objectMapper.getTypeFactory().constructParametricType(cla, new JavaType[0]);
+ }
+ }
+
}
diff --git a/im-common/src/main/java/net/sopod/soim/common/util/Reflects.java b/im-common/src/main/java/net/sopod/soim/common/util/Reflects.java
index 092536c..a27af60 100644
--- a/im-common/src/main/java/net/sopod/soim/common/util/Reflects.java
+++ b/im-common/src/main/java/net/sopod/soim/common/util/Reflects.java
@@ -1,5 +1,7 @@
package net.sopod.soim.common.util;
+import java.lang.reflect.Method;
+import java.lang.reflect.Modifier;
import java.lang.reflect.Type;
import java.util.ArrayList;
import java.util.Arrays;
@@ -56,4 +58,31 @@ public class Reflects {
return Arrays.asList(genericNames);
}
+ public static List getNonstaticMethods(Class> clazz, String name) {
+ return getNonstaticMethods(clazz, name, -1);
+ }
+
+ /**
+ * @param clazz 类
+ * @param name 方法名称
+ * @param paramSize 大于0,匹配参数个数
+ * @return 匹配到的方法集合
+ */
+ public static List getNonstaticMethods(Class> clazz, String name, int paramSize) {
+ List methodList = new ArrayList<>();
+ Method[] methods = clazz.getMethods();
+ for (Method method : methods) {
+ if (method.getName().equals(name)
+ && !Modifier.isStatic(method.getModifiers())) {
+ if (paramSize < 0) {
+ methodList.add(method);
+ } else if (method.getParameterCount() == paramSize) {
+ // 参数个数匹配
+ methodList.add(method);
+ }
+ }
+ }
+ return methodList;
+ }
+
}
diff --git a/im-common/src/main/java/net/sopod/soim/common/util/cache/LinkedMapLRUCache.java b/im-common/src/main/java/net/sopod/soim/common/util/cache/LinkedMapLRUCache.java
new file mode 100644
index 0000000..758bdec
--- /dev/null
+++ b/im-common/src/main/java/net/sopod/soim/common/util/cache/LinkedMapLRUCache.java
@@ -0,0 +1,35 @@
+package net.sopod.soim.common.util.cache;
+
+import java.util.LinkedHashMap;
+import java.util.Map;
+
+/**
+ * LinkedMapLRUCache
+ *
+ * @author tmy
+ * @date 2022-05-11 17:46
+ */
+public class LinkedMapLRUCache extends LinkedHashMap {
+
+ /**
+ * 缓存容量
+ */
+ private final int maxEntries;
+
+ public LinkedMapLRUCache(int maxEntries) {
+ this.maxEntries = maxEntries;
+ }
+
+ /**
+ * 通过重写removeEldestEntry方法,加入一定的条件,满足条件返回true。
+ *
+ * @param eldest 大链表头节点
+ * @return true,表示允许移除头节点;false,表示不允许移除头节点
+ */
+ @Override
+ protected boolean removeEldestEntry(Map.Entry eldest) {
+ //如果节点数量大于LRU缓存容量,那么返回true
+ return size() > maxEntries;
+ }
+
+}
diff --git a/im-common/src/main/java/net/sopod/soim/common/util/netty/Varint32FrameCodec.java b/im-common/src/main/java/net/sopod/soim/common/util/netty/Varint32FrameCodec.java
new file mode 100644
index 0000000..ada48aa
--- /dev/null
+++ b/im-common/src/main/java/net/sopod/soim/common/util/netty/Varint32FrameCodec.java
@@ -0,0 +1,28 @@
+package net.sopod.soim.common.util.netty;
+
+import io.netty.channel.CombinedChannelDuplexHandler;
+import io.netty.handler.codec.protobuf.ProtobufVarint32FrameDecoder;
+import io.netty.handler.codec.protobuf.ProtobufVarint32LengthFieldPrepender;
+
+/**
+ * 帧编解码器
+ * 接收字节可能会分段到达,添加帧编解码器,修改应用程序代码以不调用 readVarInt 和 writeVarInt
+ * 接收到的 ByteBuf 的大小将始终与发送的 ByteBuf 的大小匹配,因为解码器负责聚合任何片段
+ * https://netty.io/4.0/xref/io/netty/example/worldclock/WorldClockServerInitializer.html
+ *
+ * {@link ProtobufVarint32FrameDecoder} 它确保管道中的下一个处理程序一次总是接收准确的一帧字节
+ * {@link ProtobufVarint32LengthFieldPrepender} 负责将帧长度 header 添加到所有传出消息
+ *
+ * @author tmy
+ * @date 2022-05-11 17:06
+ */
+public class Varint32FrameCodec extends CombinedChannelDuplexHandler<
+ ProtobufVarint32FrameDecoder,
+ ProtobufVarint32LengthFieldPrepender> {
+
+ public Varint32FrameCodec() {
+ // FrameDecoder, LengthFieldPrepender
+ super(new ProtobufVarint32FrameDecoder(), new ProtobufVarint32LengthFieldPrepender());
+ }
+
+}
diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java b/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java
index b3b70f1..d8df009 100644
--- a/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java
+++ b/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java
@@ -5,6 +5,7 @@ import io.netty.channel.ChannelPipeline;
import io.netty.channel.socket.SocketChannel;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
+import net.sopod.soim.common.util.netty.Varint32FrameCodec;
import net.sopod.soim.core.net.ImMessageCodec;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -25,6 +26,7 @@ public class ImEntryInitializer extends ChannelInitializer {
LogLevel logLevel = logger.isDebugEnabled() ? LogLevel.DEBUG : LogLevel.INFO;
ChannelPipeline pipeline = socketChannel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel))
+ .addLast(new Varint32FrameCodec())
.addLast(new ImMessageCodec())
.addLast(new InboundImMessageHandler());
}
diff --git a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ConsistentHashTest.java b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ConsistentHashTest.java
index 061639f..9a512f4 100644
--- a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ConsistentHashTest.java
+++ b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ConsistentHashTest.java
@@ -51,10 +51,6 @@ public class ConsistentHashTest {
return HashAlgorithms.md5Hash(value);
}
- public static long hash(String value, int number) {
- return HashAlgorithms.md5Hash(value, number);
- }
-
private static void testConsistentHash() {
String[] nodes = {
"192.168.31.156:3032",
diff --git a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java
index 4f08442..be4ce17 100644
--- a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java
+++ b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java
@@ -4,6 +4,9 @@ import net.sopod.soim.common.util.HashAlgorithms;
import org.apache.commons.lang3.tuple.ImmutablePair;
import org.apache.commons.lang3.tuple.Pair;
+import java.nio.charset.StandardCharsets;
+import java.security.MessageDigest;
+import java.security.NoSuchAlgorithmException;
import java.util.*;
/**
@@ -91,7 +94,26 @@ public class UidConsistentHashSelector {
}
private static long hash(String value, int number) {
- return HashAlgorithms.md5Hash(value, number);
+ return md5Hash(value, number);
+ }
+
+ private static long md5Hash(String value, int number) {
+ MessageDigest md5;
+ try {
+ md5 = MessageDigest.getInstance("MD5");
+ } catch (NoSuchAlgorithmException e) {
+ throw new RuntimeException("MD5 not supported", e);
+ }
+ md5.reset();
+ byte[] keyBytes = value.getBytes(StandardCharsets.UTF_8);
+ md5.update(keyBytes);
+ byte[] digest = md5.digest();
+ // dubbo md5 hash algorithms
+ return (((long) (digest[3 + number * 4] & 0xFF) << 24)
+ | ((long) (digest[2 + number * 4] & 0xFF) << 16)
+ | ((long) (digest[1 + number * 4] & 0xFF) << 8)
+ | (digest[number * 4] & 0xFF))
+ & 0xFFFFFFFFL;
}
public int getIdentityHashCode() {
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java b/im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java
index 8853c4d..9c3c4ae 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java
@@ -5,13 +5,19 @@ import com.alibaba.nacos.api.exception.NacosException;
import com.alibaba.nacos.api.naming.NamingService;
import com.alibaba.nacos.api.naming.pojo.Instance;
import net.sopod.soim.common.constant.AppConstant;
+import net.sopod.soim.common.constant.DubboConstant;
+import net.sopod.soim.common.util.Collects;
import net.sopod.soim.common.util.HashAlgorithms;
import net.sopod.soim.common.util.StringUtil;
import org.apache.dubbo.common.URL;
+import org.apache.dubbo.registry.Registry;
+import org.apache.dubbo.registry.support.RegistryManager;
+import org.apache.dubbo.rpc.model.ApplicationModel;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.ApplicationContext;
+import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Properties;
@@ -116,8 +122,7 @@ public class AppContextHolder {
}
}
try {
- List allInstances = namingService.getAllInstances(AppConstant.APP_IM_ROUTER_NAME);
- return allInstances;
+ return namingService.getAllInstances(AppConstant.APP_IM_ROUTER_NAME);
} catch (NacosException e) {
throw new IllegalStateException(AppConstant.APP_IM_ROUTER_NAME + "集群信息获取失败:", e);
} finally {
@@ -131,4 +136,27 @@ public class AppContextHolder {
}
}
+ /**
+ * 注册 im-router 的API接口服务
+ */
+ public static void doRegistry() {
+ RegistryManager registryManager = ApplicationModel.defaultModel().getBeanFactory()
+ .getBean(RegistryManager.class);
+ Collection registries = registryManager.getRegistries();
+ List registryInvokerUrls = AppContextHolder.getRegistryInvokerUrls();
+ if (Collects.isNotEmpty(registries)
+ && Collects.isNotEmpty(registryInvokerUrls)) {
+ for (Registry registry : registries) {
+ for (URL invokerUrl : registryInvokerUrls) {
+ // 添加 im-router 服务id参数,生成新的 url
+ URL url = invokerUrl.addParameter(
+ DubboConstant.IM_ROUTER_ID_KEY,
+ AppContextHolder.IM_ROUTER_ID
+ );
+ registry.register(url);
+ }
+ }
+ }
+ }
+
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java b/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java
index 70eddaf..90e87c8 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java
@@ -27,6 +27,10 @@ import org.springframework.context.annotation.Configuration;
import org.springframework.core.Ordered;
import java.util.*;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors;
/**
@@ -58,13 +62,15 @@ public class ImRouterAppOnReady implements ApplicationListener {
+ for (int i = 0; i < 10; i++) {
+ RouterUser routerUser = RouterUserStorage.getInstance().get(uidCounter.getAndIncrement());
+ if (routerUser == null) {
+ return;
+ }
+ routerUser.setAccount("changeAccount:" + uidCounter.get());
+ }
+ logger.info("change log....: {}", uidCounter.get());
+ }, 30, 10, TimeUnit.SECONDS);
}
/**
@@ -86,6 +104,7 @@ public class ImRouterAppOnReady implements ApplicationListener registries = registryManager.getRegistries();
- List registryInvokerUrls = AppContextHolder.getRegistryInvokerUrls();
- if (Collects.isNotEmpty(registries)
- && Collects.isNotEmpty(registryInvokerUrls)) {
- for (Registry registry : registries) {
- for (URL invokerUrl : registryInvokerUrls) {
- // 添加 im-router 服务id参数,生成新的 url
- URL url = invokerUrl.addParameter(
- DubboConstant.IM_ROUTER_ID_KEY,
- AppContextHolder.IM_ROUTER_ID
- );
- registry.register(url);
- }
- }
- }
- }
@Override
public int getOrder() {
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java
index 1fe5326..098fb22 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java
@@ -1,8 +1,11 @@
package net.sopod.soim.router.datasync;
+import net.sopod.soim.router.config.ImRouterAppOnReady;
import net.sopod.soim.router.datasync.server.data.SyncLog;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
-import java.util.Queue;
+import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicInteger;
@@ -16,6 +19,8 @@ import java.util.concurrent.atomic.AtomicInteger;
*/
public class DataChangeTrigger {
+ private static final Logger logger = LoggerFactory.getLogger(DataChangeTrigger.class);
+
private static DataChangeTrigger INSTANCE;
// TODO subscribes 订阅 selector.select(dataKey);
@@ -31,41 +36,79 @@ public class DataChangeTrigger {
return INSTANCE;
}
- private final ConcurrentHashMap seqCounterMap = new ConcurrentHashMap<>(128);
+ /**
+ *
+ * TODO 优化
+ *
+ */
+ private final ConcurrentHashMap seqCounterMap = new ConcurrentHashMap<>(1024);
- private final Queue logQueue = new ConcurrentLinkedQueue<>();
+ /**
+ * 监听者列表
+ */
+ private final List subscribes = new ArrayList<>();
+
+ private List getSubscribes(String dataKey) {
+ if (subscribes.isEmpty()) {
+ return Collections.emptyList();
+ }
+ List acceptSubs = new ArrayList<>(5);
+ for (SyncLogSubscribe subscribe : subscribes) {
+ if (subscribe.accept(dataKey)) {
+ acceptSubs.add(subscribe);
+ }
+ }
+ return acceptSubs;
+ }
public void onUpdate(SyncTypes.SyncType syncType, String dataKey, String method, Object[] args) {
+ List acceptSubs;
+ if ((acceptSubs = getSubscribes(dataKey)).isEmpty()) {
+ return;
+ }
// 序列化 args,避免后续更改
SyncLog.UpdateLog updateLog = SyncLog.updateLog(getSeq(dataKey), syncType)
.setDataKey(dataKey)
.setMethod(method)
.setArgs(args);
- publishLog(updateLog);
+ publishLog(updateLog, acceptSubs);
+
}
public void onAdd(SyncTypes.SyncType syncType, T data) {
+ List acceptSubs;
String dataKey = syncType.getDataKey(data);
+ if ((acceptSubs = getSubscribes(dataKey)).isEmpty()) {
+ return;
+ }
AtomicInteger seqCounter = getSeqCounter(dataKey);
// 序列化 data,避免后续更改
SyncLog.AddLog addLog = SyncLog.addLog(seqCounter.getAndIncrement(), syncType)
.addData(data);
- publishLog(addLog);
+ publishLog(addLog, acceptSubs);
}
/**
* 数据删除日志
*/
public void onRemove(SyncTypes.SyncType syncType, String dataKey) {
+ List acceptSubs;
+ if ((acceptSubs = getSubscribes(dataKey)).isEmpty()) {
+ return;
+ }
SyncLog.RemoveLog removeLog = SyncLog.removeLog(getSeq(dataKey), syncType)
.setDataKey(dataKey);
- publishLog(removeLog);
+ publishLog(removeLog, acceptSubs);
}
- private void publishLog(SyncLog log) {
- // TODO 判断无订阅者跳过,新增节点同步按 dataKey 单独订阅每一个数据更新
- logQueue.add(log);
- System.out.println("publish log: "+log);
+ private void publishLog(SyncLog log, List acceptSubs) {
+ for (SyncLogSubscribe acceptSub : acceptSubs) {
+ try {
+ acceptSub.onSyncLog(log);
+ } catch (Exception e) {
+ logger.error("监听者 {} 执行错误", acceptSub, e);
+ }
+ }
}
private int getSeq(String dataKey) {
@@ -76,6 +119,13 @@ public class DataChangeTrigger {
return seqCounterMap.computeIfAbsent(dataKey, key -> new AtomicInteger());
}
+ /**
+ * 添加更改日志监听者
+ */
+ public void addSubscribe(SyncLogSubscribe syncLogSubscribe) {
+ subscribes.add(syncLogSubscribe);
+ }
+
/**
* 删除暂存数据字段等...
*/
@@ -83,4 +133,34 @@ public class DataChangeTrigger {
}
+ public static interface SyncLogSubscribe {
+
+ public abstract boolean accept(String dataKey);
+
+ public abstract void onSyncLog(SyncLog syncLog);
+
+ }
+
+ public static abstract class DataKeySyncLogSubscribe implements SyncLogSubscribe {
+
+ private final Set subscribeDataKeys;
+
+ public DataKeySyncLogSubscribe() {
+ this.subscribeDataKeys = new HashSet<>(1024);
+ }
+
+ @Override
+ public boolean accept(String dataKey) {
+ return subscribeDataKeys.contains(dataKey);
+ }
+
+ public void addSubscribeDataKey(String dataKey) {
+ subscribeDataKeys.add(dataKey);
+ }
+
+ @Override
+ public abstract void onSyncLog(SyncLog syncLog);
+
+ }
+
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
index a2c4988..5e0a02d 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
@@ -145,57 +145,49 @@ public class DataSyncProxyFactory {
// 记录数据更新操作(方法和参数)
DataChangeTrigger.instance().onUpdate(syncType, syncType.getDataKey((T) instance), methodName, args);
-
- System.out.println("intercept invoke:" + methodName);
+ // System.out.println("intercept invoke:" + methodName);
return methodProxy.invokeSuper(instance, args);
}
}
- public static void main(String[] args) {
-// RouterUser user1 = new RouterUser();
-// user1.setUid(10086L);
-// user1.setAccount("日月光");
-//
-// RouterUser routerUser = newProxyInstance(SyncTypes.ROUTER_USER, user1);
-// System.out.println(routerUser);
-// routerUser.setOnlineTime(ImClock.millis());
-// System.out.println(routerUser);
-
- ConcurrentHashMap map = new ConcurrentHashMap<>();
- map.put("10081", new RouterUser().setAccount("阿基过天玺"));
- map.put("10082", new RouterUser().setAccount("家国"));
- map.put("10083", new RouterUser().setAccount("天下"));
- Collection values = map.values();
- map.put("10084", new RouterUser().setAccount("晚风"));
- new Thread(() -> {
- for (int i = 0; i < 100; i++) {
- map.put("100" + i, new RouterUser().setAccount("灯" + i));
- if (i % 10 == 0) {
- try {
- Thread.sleep(100);
- } catch (InterruptedException e) {
- e.printStackTrace();
- }
- }
- }
- }).start();
- int i = 0;
-
- for (RouterUser value : values) {
- System.out.println(value);
- i++;
- if (i % 3 == 0) {
- try {
- Thread.sleep(100);
- } catch (InterruptedException e) {
- e.printStackTrace();
- }
- }
+ public static class A {
+ public static String name() {
+ return "1";
}
- System.out.println(values.getClass());
- System.out.println(values);
+ public void age() {}
+ }
+
+ public static class B extends A {
+ public void some() {}
+ public static void what() {}
+ }
+ public static void main(String[] args) {
+// Enhancer enhancer = new Enhancer();
+// enhancer.setSuperclass(B.class);
+// enhancer.setCallback(new MethodInterceptor() {
+// @Override
+// public Object intercept(Object o, Method method, Object[] objects, MethodProxy methodProxy) throws Throwable {
+// System.out.println("intercept:" + method.getName());
+// return methodProxy.invokeSuper(o, objects);
+// }
+// });
+// B b = (B)enhancer.create();
+// b.some();
+// b.age();
+// b.name();
+// b.what();
+// b.equals(b);
+// B.name();
+
+ for (Method m : B.class.getDeclaredMethods()) {
+ System.out.println(m.getName());
+ }
+ System.out.println("=============================");
+ for (Method m : B.class.getMethods()) {
+ System.out.println(m.getName());
+ }
}
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashService.java
similarity index 85%
rename from im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java
rename to im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashService.java
index b6783d4..778aecd 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashService.java
@@ -1,13 +1,11 @@
-package net.sopod.soim.router.datasync.server;
+package net.sopod.soim.router.datasync;
import io.netty.channel.Channel;
import io.netty.util.AttributeKey;
import net.sopod.soim.common.util.Collects;
import net.sopod.soim.router.api.route.UidConsistentHashSelector;
import net.sopod.soim.router.config.AppContextHolder;
-import net.sopod.soim.router.datasync.DataSync;
-import net.sopod.soim.router.datasync.DataSyncStorage;
-import net.sopod.soim.router.datasync.SyncTypes;
+import net.sopod.soim.router.datasync.server.data.SyncCmd;
import net.sopod.soim.router.datasync.server.data.SyncLog;
import org.apache.commons.lang3.tuple.Pair;
import org.slf4j.Logger;
@@ -18,13 +16,12 @@ import java.util.*;
import java.util.concurrent.atomic.AtomicInteger;
/**
- * SyncLogPushService
+ * SyncLogByHashService
*
* @author tmy
* @date 2022-05-10 00:30
*/
-
-public class SyncLogByHashService {
+public class SyncLogByHashService extends DataChangeTrigger.DataKeySyncLogSubscribe {
private static final Logger logger = LoggerFactory.getLogger(SyncLogByHashService.class);
@@ -37,7 +34,7 @@ public class SyncLogByHashService {
private final UidConsistentHashSelector selector;
- Set syncedDataKeys = new HashSet<>();
+ // Set syncedDataKeys = new HashSet<>();
public SyncLogByHashService(Channel clientChannel, String newNodeAddr) {
this.clientChannel = new WeakReference<>(clientChannel);
@@ -47,18 +44,9 @@ public class SyncLogByHashService {
twoNodes.put(newNodeAddr, newNodeAddr);
twoNodes.put(AppContextHolder.getAppAddr(), AppContextHolder.getAppAddr());
selector = new UidConsistentHashSelector<>(twoNodes, twoNodes.hashCode());
- }
- public static void main(String[] args) {
- Map> twoNodes = new HashMap<>();
- twoNodes.put("192.168.101.69:3032", Pair.of("192.168.101.69:3032", new AtomicInteger()));
- twoNodes.put("192.168.101.69:3031", Pair.of("192.168.101.69:3031", new AtomicInteger()));
- UidConsistentHashSelector> selector = new UidConsistentHashSelector<>(twoNodes, twoNodes.hashCode());
- for (int i = 10000; i < 11000; i++) {
- Pair pair = selector.select(i + "");
- pair.getRight().incrementAndGet();
- }
- System.out.println(twoNodes);
+ // 添加数据变化监听
+ DataChangeTrigger.instance().addSubscribe(this);
}
private Map, SyncTypes.SyncType extends DataSync>> storages;
@@ -102,8 +90,10 @@ public class SyncLogByHashService {
// TODO 数据加锁
addLog.addData(data);
// TODO DataChangeTrigger.instance().subscribe(syncType, dataKey)
+
// 注册更改日志
- syncedDataKeys.add(dataKey);
+ super.addSubscribeDataKey(dataKey);
+
// TODO 解锁
count ++; totalCount ++;
@@ -159,6 +149,35 @@ public class SyncLogByHashService {
private void syncFinishSuccess() {
logger.info("sync push finish and success!");
+ Channel channel = clientChannel.get();
+ if (channel != null && channel.isActive()) {
+ SyncCmd syncEndCmd = new SyncCmd();
+ syncEndCmd.setCmdType(SyncCmd.SYNC_END);
+ channel.writeAndFlush(syncEndCmd);
+ }
+ }
+
+ /**
+ * 修改数据日志
+ */
+ @Override
+ public void onSyncLog(SyncLog syncLog) {
+ Channel channel = clientChannel.get();
+ if (channel != null && channel.isActive()) {
+ channel.writeAndFlush(syncLog);
+ }
+ }
+
+ public static void main(String[] args) {
+ Map> twoNodes = new HashMap<>();
+ twoNodes.put("192.168.101.69:3032", Pair.of("192.168.101.69:3032", new AtomicInteger()));
+ twoNodes.put("192.168.101.69:3031", Pair.of("192.168.101.69:3031", new AtomicInteger()));
+ UidConsistentHashSelector> selector = new UidConsistentHashSelector<>(twoNodes, twoNodes.hashCode());
+ for (int i = 10000; i < 11000; i++) {
+ Pair pair = selector.select(i + "");
+ pair.getRight().incrementAndGet();
+ }
+ System.out.println(twoNodes);
}
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java
index 9a07e74..83ab332 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java
@@ -12,10 +12,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
-import java.util.Iterator;
-import java.util.LinkedList;
-import java.util.List;
-import java.util.Map;
+import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
/**
@@ -33,12 +30,16 @@ public class SyncLogMigrateService {
private volatile List> migrateHosts;
private LinkedList> curMigrateHosts;
+ private List migrateClients = new ArrayList<>();
+
+ private SyncClient curClient;
+
/**
* 同步数据
* @param migrateHosts 同步数据节点
*/
public void migrateSyncLog(List> migrateHosts) {
- Preconditions.checkState(migrateHosts != null, "当前正在进行数据同步");
+ Preconditions.checkState(this.migrateHosts == null, "当前正在进行数据同步");
this.migrateHosts = migrateHosts;
this.curMigrateHosts = new LinkedList<>(migrateHosts);
logger.info("开始连接{}节点同步服务...", AppConstant.APP_IM_ROUTER_NAME);
@@ -50,14 +51,22 @@ public class SyncLogMigrateService {
* @return 是否还有下一个同步数据节点
*/
public synchronized boolean syncNextHost() {
+ if (this.curClient != null) {
+ // TODO 优化多节点同时同步数据
+ throw new IllegalStateException("当前有客户端正在全量同步数据");
+ }
if (this.curMigrateHosts.isEmpty()) {
return false;
}
Pair nextHost = this.curMigrateHosts.removeFirst();
+ if (nextHost == null) {
+ this.allHostSyncFinish();
+ return true;
+ }
try {
- SyncClient client = new SyncClient();
- client.connect(nextHost.getLeft(), nextHost.getRight());
- client.syncLogByHash(AppContextHolder.getAppAddr());
+ this.curClient = new SyncClient();
+ this.curClient.connect(nextHost.getLeft(), nextHost.getRight());
+ this.curClient.syncLogByHash(AppContextHolder.getAppAddr());
} catch (InterruptedException e) {
logger.error("节点{}:{}连接失败, 跳过!", nextHost.getLeft(), nextHost.getRight(), e);
// 同步下一个节点
@@ -66,6 +75,29 @@ public class SyncLogMigrateService {
return true;
}
+ public synchronized void curHostSyncEnd() {
+ if (this.curClient == null) {
+ logger.warn("当前无正在全量同步数据的客户端");
+ return;
+ }
+ this.migrateClients.add(this.curClient);
+ this.curClient = null;
+ this.syncNextHost();
+ }
+
+ /**
+ * 所有节点全量同步完成
+ * 注册服务
+ * 接受后续日志数据
+ * n秒无日志数据后,发送消息可移除数据
+ * 断开连接
+ */
+ private void allHostSyncFinish() {
+ logger.info("所有节点数据同步完成:执行注册服务....");
+ AppContextHolder.doRegistry();
+ logger.info("注册服务成功");
+ }
+
@Deprecated
private void fullSync() {
Map routerUserMap = RouterUserStorage.getInstance().getRouterUserMap();
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java
index a02b248..f9e1fd1 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java
@@ -9,6 +9,7 @@ import io.netty.channel.socket.nio.NioSocketChannel;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
import io.netty.util.AttributeKey;
+import net.sopod.soim.common.util.netty.Varint32FrameCodec;
import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.datasync.SyncTypes;
import net.sopod.soim.router.datasync.server.codec.SyncCmdCodec;
@@ -38,8 +39,11 @@ public class SyncClient {
.handler(new ChannelInitializer() {
@Override
protected void initChannel(SocketChannel channel) throws Exception {
+
channel.pipeline()
.addLast(new LoggingHandler(LogLevel.INFO))
+ // 接收字节可能会分段到达,添加帧编解码器
+ .addLast(new Varint32FrameCodec())
.addLast(new SyncCmdCodec())
.addLast(new SyncLogEncoder())
.addLast(new SyncCmdClientHandler())
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java
index 5a9a238..0a2b776 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java
@@ -10,6 +10,7 @@ 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.common.util.netty.Varint32FrameCodec;
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;
@@ -49,8 +50,11 @@ public class SyncServer {
: logger.isInfoEnabled() ? LogLevel.INFO
: logger.isWarnEnabled() ? LogLevel.WARN
: logger.isErrorEnabled() ? LogLevel.ERROR : LogLevel.INFO;
+
ChannelPipeline pipeline = channel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel))
+ // 接收字节可能会分段到达,添加帧编解码器
+ .addLast(new Varint32FrameCodec())
.addLast(new SyncCmdCodec())
.addLast(new SyncLogEncoder())
.addLast(new SyncCmdServerHandler())
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java
index ad70035..4de369e 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java
@@ -1,6 +1,7 @@
package net.sopod.soim.router.datasync.server.codec;
import io.netty.buffer.ByteBuf;
+import io.netty.buffer.PooledByteBufAllocator;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.CombinedChannelDuplexHandler;
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java
index 68a1c07..f23b9a3 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java
@@ -42,7 +42,7 @@ public class SyncCmd {
/**
* 同步结束命令
*/
- public static final Integer SYNC_END = 8;
+ public static final int SYNC_END = 8;
/**
* SyncLog 推送命令
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncLog.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncLog.java
index 7ff9e8b..4650b92 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncLog.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncLog.java
@@ -107,7 +107,6 @@ public class SyncLog implements Serializable {
private byte[] toBytes0() {
byte[] dataKeyBytes = dataKey == null ? new byte[0] : dataKey.getBytes();
- // byte[] clazzBytes = clazz == null ? new byte[0] : clazz.getBytes();
byte[] methodBytes = method == null ? new byte[0] : method.getBytes();
int argSize = args == null ? 0 : args.length;
@@ -144,7 +143,7 @@ public class SyncLog implements Serializable {
+ 4 // logSeq 序列号
+ 4 // dataKey 数据主键标识字节长度
+ dataKeyBytes.length // dataKey 数据字节
- + 4 // clazz 字节长度
+ // + 4 // clazz 字节长度
// + clazzBytes.length // clazz字节
+ 4 // method 字节长度
+ methodBytes.length // method 字节
@@ -157,12 +156,12 @@ public class SyncLog implements Serializable {
// Bytes.short2bytes(MAGIC, bytes, offset);
// offset += 2;
- Bytes.int2bytes(bytes.length - 6, bytes, offset);
+ Bytes.int2bytes(bytes.length - 4, bytes, offset);
offset += 4;
bytes[offset] = (byte) syncType;
offset += 1;
- bytes[offset] = (byte)operateType;
+ bytes[offset] = (byte) operateType;
offset += 1;
Bytes.int2bytes(logSeq, bytes, offset);
@@ -237,14 +236,6 @@ public class SyncLog implements Serializable {
int argSize = buf.readInt();
log.args = new String[argSize];
if (argSize > 0) {
-// Method method = getClassMethod(log.clazz, log.method);
-// if (method == null) {
-// throw new IllegalStateException("类" + log.clazz + "方法" + log.method + "未找到");
-// }
-// Class>[] paramTypes = method.getParameterTypes();
-// if (paramTypes.length != argSize) {
-// throw new IllegalStateException("类" + log.clazz + "方法" + log.method + "指定参数" + argSize + "个,查到参数" + paramTypes.length + "个");
-// }
// 复用 bytes
byte[] argBytes = new byte[0];
for (int i = 0; i < argSize; i++) {
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java
index fb97f87..f4ff225 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java
@@ -2,6 +2,8 @@ package net.sopod.soim.router.datasync.server.handler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
+import net.sopod.soim.router.config.AppContextHolder;
+import net.sopod.soim.router.datasync.SyncLogMigrateService;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -25,6 +27,9 @@ public class SyncCmdClientHandler extends SimpleChannelInboundHandler {
case SyncCmd.PONG:
this.handlePong(ctx, syncCmd);
break;
+ case SyncCmd.SYNC_END:
+ this.handleSyncEnd(ctx, syncCmd);
+ break;
}
}
@@ -39,4 +44,10 @@ public class SyncCmdClientHandler extends SimpleChannelInboundHandler {
System.out.println("pong: " + ctx.channel());
}
+ private void handleSyncEnd(ChannelHandlerContext ctx, SyncCmd syncCmd) {
+ logger.info("sync end: {}", ctx.channel());
+ SyncLogMigrateService migrateService = AppContextHolder.getBean(SyncLogMigrateService.class);
+ migrateService.curHostSyncEnd();
+ }
+
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java
index 8a5a792..b8d0937 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java
@@ -2,7 +2,7 @@ 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.SyncLogByHashService;
+import net.sopod.soim.router.datasync.SyncLogByHashService;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java
index d7a4a73..e90cd11 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java
@@ -2,8 +2,12 @@ package net.sopod.soim.router.datasync.server.handler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
+import net.sopod.soim.common.util.Collects;
+import net.sopod.soim.common.util.Jackson;
+import net.sopod.soim.common.util.Reflects;
+import net.sopod.soim.common.util.cache.LinkedMapLRUCache;
+import net.sopod.soim.router.cache.RouterUserStorage;
import net.sopod.soim.router.datasync.DataSync;
-import net.sopod.soim.router.datasync.DataSyncStorage;
import net.sopod.soim.router.datasync.SyncTypes;
import net.sopod.soim.router.datasync.server.codec.CodecUtil;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
@@ -11,7 +15,10 @@ import net.sopod.soim.router.datasync.server.data.SyncLog;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import java.util.List;
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
+import java.lang.reflect.Type;
+import java.util.*;
/**
* SyncLogHandler
@@ -26,17 +33,75 @@ public class SyncLogClientHandler extends SimpleChannelInboundHandler {
@Override
protected void channelRead0(ChannelHandlerContext ctx, SyncLog syncLog) {
- List bytesList = syncLog.getSerializeDataCollect();
+ switch (syncLog.getOperateType()) {
+ case SyncLog.OPT_ADD:
+ this.handleAddSyncLog(ctx, syncLog);
+ break;
+ case SyncLog.OPT_REMOVE:
+ this.handleRemoveSyncLog(ctx, syncLog);
+ break;
+ case SyncLog.OPT_UPDATE:
+ this.handleUpdateSyncLog(ctx, syncLog);
+ break;
+ }
+ }
+
+ private void handleAddSyncLog(ChannelHandlerContext ctx, SyncLog syncLog) {
SyncTypes.SyncType syncType = SyncTypes.getSyncType(syncLog.getSyncType());
- for (byte[] bytes : bytesList) {
- DataSync instance = CodecUtil.decode(bytes, syncType.dataType());
- syncType.addData(instance);
- System.out.println(instance);
+ List bytesList = syncLog.getSerializeDataCollect();
+ if (Collects.isNotEmpty(bytesList)) {
+ for (byte[] bytes : bytesList) {
+ DataSync instance = CodecUtil.decode(bytes, syncType.dataType());
+ syncType.addData(instance);
+ logger.info("addData: {}", instance);
+ }
}
// 同步完成响应
SyncCmd syncCmd = new SyncCmd().setCmdType(SyncCmd.SYNC_BY_HASH_ACK);
ctx.writeAndFlush(syncCmd);
- logger.info("storage: {}", DataSyncStorage.getStorages());
+ logger.info("users size: {}", RouterUserStorage.getInstance().getRouterUserMap().size());
+ }
+
+ private void handleRemoveSyncLog(ChannelHandlerContext ctx, SyncLog syncLog) {
+ SyncTypes.SyncType syncType = SyncTypes.getSyncType(syncLog.getSyncType());
+ syncType.removeData(syncLog.getDataKey());
+ }
+
+ private static final LinkedMapLRUCache syncTypeMethodCache = new LinkedMapLRUCache<>(16);
+
+ private void handleUpdateSyncLog(ChannelHandlerContext ctx, SyncLog syncLog) {
+ SyncTypes.SyncType syncType = SyncTypes.getSyncType(syncLog.getSyncType());
+ String methodName = syncLog.getMethod();
+ String[] args = syncLog.getArgs();
+ Class dataType = syncType.dataType();
+ String cacheKey = dataType.getName() + "." + methodName + "(" + args.length + ")";
+ Method method = syncTypeMethodCache.get(cacheKey);
+ if (method == null) {
+ synchronized(syncTypeMethodCache) {
+ if (null == (method = syncTypeMethodCache.get(cacheKey))) {
+ List methods = Reflects.getNonstaticMethods(dataType, methodName, args.length);
+ if (Collects.isEmpty(methods)) {
+ logger.warn("{} method {} not found!", dataType, methodName);
+ return;
+ }
+ method = methods.get(0);
+ syncTypeMethodCache.put(cacheKey, method);
+ }
+ }
+ }
+ Type[] types = method.getGenericParameterTypes();
+ Object[] params = new Object[types.length];
+ for (int i = 0; i < types.length; i++) {
+ params[i] = Jackson.json().readWithType(args[i], types[i]);
+ }
+ DataSync data = syncType.getData(syncLog.getDataKey());
+ try {
+ method.setAccessible(true);
+ method.invoke(data, params);
+ logger.info("invoke:{}, args:{}", method.getName(), params);
+ } catch (IllegalAccessException | InvocationTargetException e) {
+ e.printStackTrace();
+ }
}
}