diff --git a/im-common/pom.xml b/im-common/pom.xml
index 2cb898a..03ef45a 100644
--- a/im-common/pom.xml
+++ b/im-common/pom.xml
@@ -26,6 +26,10 @@
com.google.guava
guava
+
+ org.apache.commons
+ commons-lang3
+
io.netty
netty-common
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 58c093e..061639f 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
@@ -93,7 +93,8 @@ public class ConsistentHashTest {
public static void main(String[] args) {
// {192.168.1.103=32, 192.168.1.100=26, 192.168.1.101=19, 192.168.1.102=23}
- testConsistentHash();
+ // testConsistentHash();
+ testTreeMap();
}
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 833b08e..4f08442 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
@@ -1,9 +1,10 @@
package net.sopod.soim.router.api.route;
import net.sopod.soim.common.util.HashAlgorithms;
+import org.apache.commons.lang3.tuple.ImmutablePair;
+import org.apache.commons.lang3.tuple.Pair;
-import java.util.Map;
-import java.util.TreeMap;
+import java.util.*;
/**
* UidConsistentHashSelector
@@ -18,7 +19,7 @@ public class UidConsistentHashSelector {
*/
private static final int VIRTUAL_NODE_SIZE = 120;
- private final TreeMap virtualNodeMap;
+ private final TreeMap> virtualNodeMap;
private final int identityHashCode;
@@ -33,7 +34,17 @@ public class UidConsistentHashSelector {
for (int i = 0, len = VIRTUAL_NODE_SIZE / 4; i < len; i++) {
for (int h = 0; h < 4; h++) {
long hash = hash(serverAddr + i, h);
- this.virtualNodeMap.put(hash, value);
+ Pair lastNode;
+ if (null == (lastNode = this.virtualNodeMap.get(hash))) {
+ this.virtualNodeMap.put(hash, ImmutablePair.of(serverAddr, value));
+ } else {
+ // hash 冲突取排序小的节点
+ List twoNode = Arrays.asList(lastNode.getLeft(), serverAddr);
+ Collections.sort(twoNode);
+ if (twoNode.get(0).equals(serverAddr)) {
+ this.virtualNodeMap.put(hash, ImmutablePair.of(serverAddr, value));
+ }
+ }
}
}
}
@@ -44,11 +55,39 @@ public class UidConsistentHashSelector {
throw new IllegalStateException("im-router consistent hash route, ctx uid can not be null!");
}
long hash = hash(uid, 0);
- Map.Entry entry = virtualNodeMap.ceilingEntry(hash);
+ Map.Entry> entry = virtualNodeMap.ceilingEntry(hash);
if (entry == null) {
entry = virtualNodeMap.firstEntry();
}
- return entry.getValue();
+ return entry.getValue().getRight();
+ }
+
+ /**
+ * 计算加入新节点需要迁移数据的节点
+ * @param newNode 新节点地址
+ */
+ public Set selectMigrateNodes(String newNode) {
+ Set nodes = new HashSet<>();
+ for (int i = 0, len = VIRTUAL_NODE_SIZE / 4; i < len; i++) {
+ for (int h = 0; h < 4; h++) {
+ long hash = hash(newNode + i, h);
+ Map.Entry> entry = this.virtualNodeMap.ceilingEntry(hash);
+ if (entry == null) {
+ entry = this.virtualNodeMap.firstEntry();
+ }
+ // hash 冲突,取排序最小的一个(老节点)
+ if (entry.getKey().equals(hash)) {
+ List twoNode = Arrays.asList(entry.getValue().getLeft(), newNode);
+ Collections.sort(twoNode);
+ // 选中不是该节点跳过
+ if (!twoNode.get(0).equals(newNode)) {
+ continue;
+ }
+ }
+ nodes.add(entry.getValue().getRight());
+ }
+ }
+ return nodes;
}
private static long hash(String value, int number) {
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppContextHolder.java b/im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java
similarity index 52%
rename from im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppContextHolder.java
rename to im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java
index 22450d2..0524950 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppContextHolder.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java
@@ -1,8 +1,10 @@
package net.sopod.soim.router.config;
+import net.sopod.soim.common.constant.AppConstant;
import net.sopod.soim.common.util.HashAlgorithms;
import net.sopod.soim.common.util.StringUtil;
import org.apache.dubbo.common.URL;
+import org.springframework.context.ApplicationContext;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
@@ -16,14 +18,18 @@ import java.util.concurrent.CopyOnWriteArrayList;
* @author tmy
* @date 2022-05-04 09:21
*/
-public class ImRouterAppContextHolder {
+public class AppContextHolder {
+
+ private static ApplicationContext applicationContext;
/**
* 要注册的 provider 服务的 url 列表
*/
private static final List registryInvokerUrls = new CopyOnWriteArrayList<>();
- private static String appServiceAddr;
+ private static String appAddr;
+ private static String appHost;
+ private static int appPort;
public static final String IM_ROUTER_ID;
@@ -31,6 +37,14 @@ public class ImRouterAppContextHolder {
IM_ROUTER_ID = String.valueOf(HashAlgorithms.md5Hash(StringUtil.randomUUID()));
}
+ public static void setApplicationContext(ApplicationContext applicationContext) {
+ AppContextHolder.applicationContext = applicationContext;
+ }
+
+ public static T getBean(Class type) {
+ return applicationContext.getBean(type);
+ }
+
public static void addRegistryInvokerUrl(URL registryInvokerUrl) {
registryInvokerUrls.add(registryInvokerUrl);
}
@@ -39,12 +53,22 @@ public class ImRouterAppContextHolder {
return registryInvokerUrls;
}
- public static void setAppServiceAddr(String appServiceAddr) {
- ImRouterAppContextHolder.appServiceAddr = appServiceAddr;
+ public static void setAppServiceAddr(String host, int port) {
+ AppContextHolder.appAddr = host + ":" + port;
+ AppContextHolder.appHost = host;
+ AppContextHolder.appPort = port;
+ }
+
+ public static String getAppAddr() {
+ return appAddr;
+ }
+
+ public static String getAppHost() {
+ return appHost;
}
- public static String getAppServiceAddr() {
- return appServiceAddr;
+ public static int getAppPort() {
+ return appPort;
}
}
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 b83a151..59741c0 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
@@ -8,11 +8,14 @@ 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.router.api.route.UidConsistentHashSelector;
+import net.sopod.soim.router.datasync.SyncLogMigrateService;
+import net.sopod.soim.router.datasync.server.session.SyncServerSession;
+import org.apache.commons.lang3.tuple.ImmutablePair;
+import org.apache.commons.lang3.tuple.Pair;
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.apache.dubbo.rpc.proxy.AbstractProxyInvoker;
import org.apache.dubbo.spring.boot.context.event.AwaitingNonWebApplicationListener;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -22,6 +25,7 @@ import org.springframework.context.annotation.Configuration;
import org.springframework.core.Ordered;
import java.util.*;
+import java.util.stream.Collectors;
/**
* AppcationInitialed
@@ -35,13 +39,31 @@ public class ImRouterAppOnReady implements ApplicationListener registries = registryManager.getRegistries();
- List registryInvokerUrls = ImRouterAppContextHolder.getRegistryInvokerUrls();
+ List registryInvokerUrls = AppContextHolder.getRegistryInvokerUrls();
if (Collects.isNotEmpty(registries)
&& Collects.isNotEmpty(registryInvokerUrls)) {
for (Registry registry : registries) {
@@ -59,7 +81,7 @@ public class ImRouterAppOnReady implements ApplicationListener clusterImEntryInstance = getClusterImEntryInstance();
if (Collects.isEmpty(clusterImEntryInstance)) {
logger.info("当前集群无{}节点,直接启动", AppConstant.APP_IM_ROUTER_NAME);
- return;
+ return true;
}
- for (Instance instance : clusterImEntryInstance) {
- String host = instance.getIp();
- // 同步服务器端口偏移量 1000
- int port = instance.getPort();
+ logger.info("当前集群{}节点: {} of {}",
+ AppConstant.APP_IM_ROUTER_NAME,
+ clusterImEntryInstance.size(),
+ clusterImEntryInstance.stream().map(Instance::toInetAddr).collect(Collectors.toList()));
+
+ // 构建一致性 hash 环,计算需要同步数据的节点
+ Map addrInstanceMap = Collects.collect2Map(clusterImEntryInstance,
+ Instance::toInetAddr, // 服务地址, 如: 192.168.56.1:3031
+ new HashMap<>(Collects.mapCapacity(clusterImEntryInstance.size()))
+ );
+ UidConsistentHashSelector selector = new UidConsistentHashSelector<>(addrInstanceMap, addrInstanceMap.hashCode());
+ Set migrateNodes = selector.selectMigrateNodes(AppContextHolder.getAppAddr());
+ logger.info("需迁移数据{}节点: {} of {}",
+ AppConstant.APP_IM_ROUTER_NAME,
+ migrateNodes.size(),
+ migrateNodes.stream().map(Instance::toInetAddr).collect(Collectors.toList()));
+ if (Collects.isEmpty(migrateNodes)) {
+ return true;
}
+ // 发起客户端连接,开始同步数据
+ List> migrateHosts = migrateNodes.stream()
+ .map(instance -> ImmutablePair.of(instance.getIp(), instance.getPort() + SYNC_SERVER_PORT_OFFSET))
+ .collect(Collectors.toList());
+ // 开始同步数据
+ syncLogMigrateService.migrateSyncLog(migrateHosts);
+ return false;
}
private List getClusterImEntryInstance() {
@@ -99,6 +145,14 @@ public class ImRouterAppOnReady implements ApplicationListener exporter) throws RpcException {
URL invokerUrl = exporter.getInvoker().getUrl();
if (!InjvmProtocol.NAME.equals(invokerUrl.getProtocol())) {
- ImRouterAppContextHolder.addRegistryInvokerUrl(invokerUrl);
- if (ImRouterAppContextHolder.getAppServiceAddr() == null) {
- ImRouterAppContextHolder.setAppServiceAddr(invokerUrl.getAddress());
- logger.info("im-router registry serverAddr: {}", invokerUrl.getAddress());
+ AppContextHolder.addRegistryInvokerUrl(invokerUrl);
+ if (AppContextHolder.getAppAddr() == null) {
+ // 保存服务注册地址
+ AppContextHolder.setAppServiceAddr(invokerUrl.getHost(), invokerUrl.getPort());
+ logger.info("im-router registry serverAddr: {}:{}", invokerUrl.getHost(), invokerUrl.getPort());
}
}
}
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
new file mode 100644
index 0000000..9a07e74
--- /dev/null
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java
@@ -0,0 +1,90 @@
+package net.sopod.soim.router.datasync;
+
+import com.google.common.base.Preconditions;
+import net.sopod.soim.common.constant.AppConstant;
+import net.sopod.soim.common.util.Jackson;
+import net.sopod.soim.router.cache.RouterUser;
+import net.sopod.soim.router.cache.RouterUserStorage;
+import net.sopod.soim.router.config.AppContextHolder;
+import net.sopod.soim.router.datasync.server.SyncClient;
+import org.apache.commons.lang3.tuple.Pair;
+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.concurrent.ConcurrentHashMap;
+
+/**
+ * SyncLogMigrateService
+ * 新节点集群数据迁移
+ *
+ * @author tmy
+ * @date 2022-05-08 09:42
+ */
+@Service
+public class SyncLogMigrateService {
+
+ private static final Logger logger = LoggerFactory.getLogger(SyncLogMigrateService.class);
+
+ private volatile List> migrateHosts;
+ private LinkedList> curMigrateHosts;
+
+ /**
+ * 同步数据
+ * @param migrateHosts 同步数据节点
+ */
+ public void migrateSyncLog(List> migrateHosts) {
+ Preconditions.checkState(migrateHosts != null, "当前正在进行数据同步");
+ this.migrateHosts = migrateHosts;
+ this.curMigrateHosts = new LinkedList<>(migrateHosts);
+ logger.info("开始连接{}节点同步服务...", AppConstant.APP_IM_ROUTER_NAME);
+ this.syncNextHost();
+ }
+
+ /**
+ *
+ * @return 是否还有下一个同步数据节点
+ */
+ public synchronized boolean syncNextHost() {
+ if (this.curMigrateHosts.isEmpty()) {
+ return false;
+ }
+ Pair nextHost = this.curMigrateHosts.removeFirst();
+ try {
+ SyncClient client = new SyncClient();
+ client.connect(nextHost.getLeft(), nextHost.getRight());
+ client.syncLogByHash(AppContextHolder.getAppAddr());
+ } catch (InterruptedException e) {
+ logger.error("节点{}:{}连接失败, 跳过!", nextHost.getLeft(), nextHost.getRight(), e);
+ // 同步下一个节点
+ return this.syncNextHost();
+ }
+ return true;
+ }
+
+ @Deprecated
+ private void fullSync() {
+ Map routerUserMap = RouterUserStorage.getInstance().getRouterUserMap();
+ for (Map.Entry entry : routerUserMap.entrySet()) {
+ Long key = entry.getKey();
+ RouterUser value = entry.getValue();
+ Jackson.json().serialize(value);
+ }
+ }
+
+ public static void main(String[] args) {
+ ConcurrentHashMap map = new ConcurrentHashMap<>();
+ map.put("1", "A");
+ map.put("2", "B");
+ Iterator iterator = map.values().iterator();
+ System.out.println("a." + iterator.next());
+ map.remove("2");
+ System.out.println("b." + iterator.next());
+ System.out.println(iterator.next());
+ }
+
+}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncService.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncService.java
deleted file mode 100644
index 8613180..0000000
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncService.java
+++ /dev/null
@@ -1,36 +0,0 @@
-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.nio.channels.Channel;
-import java.util.Map;
-
-/**
- * SyncService
- *
- * @author tmy
- * @date 2022-05-08 09:42
- */
-public class SyncService {
-
- public SyncService(String clientAddr, Channel channel) {
-
- }
-
- public void fullSync() {
- Map routerUserMap = RouterUserStorage.getInstance().getRouterUserMap();
- for (Map.Entry entry : routerUserMap.entrySet()) {
- Long key = entry.getKey();
- RouterUser value = entry.getValue();
- Jackson.json().serialize(value);
-
- }
- }
-
- public static void main(String[] args) {
-
- }
-
-}
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 ec22424..9285706 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
@@ -69,6 +69,22 @@ public class SyncClient {
clientChannel.attr(key).get();
}
+ /**
+ * 全量同步数据
+ */
+ public void syncLogByHash() {
+
+ }
+
+ /**
+ * 通过计算一致性 hash 同步数据
+ * @param currentAddr 当前新节点地址,hash(currentAddr)
+ */
+ public void syncLogByHash(String currentAddr) {
+ SyncCmd syncByHash = SyncCmd.syncByHash(currentAddr);
+ clientChannel.writeAndFlush(syncByHash);
+ }
+
public void close() {
this.group.shutdownGracefully();
}
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/server/SyncLogByHashService.java
new file mode 100644
index 0000000..d0fd556
--- /dev/null
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java
@@ -0,0 +1,29 @@
+package net.sopod.soim.router.datasync.server;
+
+import io.netty.util.AttributeKey;
+
+/**
+ * SyncLogPushService
+ *
+ * @author tmy
+ * @date 2022-05-10 00:30
+ */
+
+public class SyncLogByHashService {
+
+ public static final AttributeKey ATTR_KEY = AttributeKey
+ .valueOf(SyncLogByHashService.class, "SYNC_LOG_BY_HASH_SERVICE");
+
+ private final String newNodeAddr;
+
+ public SyncLogByHashService(String newNodeAddr) {
+ this.newNodeAddr = newNodeAddr;
+
+ }
+
+ public void startPush() {
+ // TODO 开始数据推送
+
+ }
+
+}
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 979e7f1..5a9a238 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
@@ -64,11 +64,11 @@ public class SyncServer {
}
return;
}
- logger.info("sync-server listening at {}...", port);
+ logger.info("im-router sync-server listening at {}...", port);
});
}
- public void shutdown() {
+ public void close() {
if (boss != null) {
boss.shutdownGracefully();
}
@@ -82,8 +82,6 @@ public class SyncServer {
.start(9999, err -> {
logger.error("服务启动失败");
});
-
-
}
}
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 21a3105..68a1c07 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
@@ -25,20 +25,37 @@ public class SyncCmd {
public static final int PONG = 2;
/**
- * 同步命令
+ * 全量同步命令
*/
- public static final Integer SYNC = 3;
+ public static final Integer SYNC_FULL = 3;
+
+ /**
+ * 计算一致性hash同步数据
+ */
+ public static final int SYNC_BY_HASH = 4;
+
+ /**
+ * 一致性hash同步数据,收到后ACK响应
+ */
+ public static final int SYNC_BY_HASH_ACK = 5;
/**
* 同步结束命令
*/
- public static final Integer SYNC_END = 4;
+ public static final Integer SYNC_END = 8;
/**
* SyncLog 推送命令
*/
public static final Integer SYNC_LOG = 10;
+ public static SyncCmd syncByHash(String hashNode) {
+ SyncCmd syncCmd = new SyncCmd();
+ syncCmd.setCmdType(SYNC_BY_HASH);
+ syncCmd.setParam1(hashNode);
+ return syncCmd;
+ }
+
private Integer cmdType;
private String param1;
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncStatus.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncStatus.java
new file mode 100644
index 0000000..1f179d3
--- /dev/null
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncStatus.java
@@ -0,0 +1,16 @@
+package net.sopod.soim.router.datasync.server.data;
+
+/**
+ * SyncStatus
+ *
+ * @author tmy
+ * @date 2022-05-10 00:18
+ */
+public enum SyncStatus {
+
+ NONE, // 无动作
+ FULL_SYNCING, // 数据全量同步中
+ SYNC_CHANGE_LOG, // 全量同步结束,同步更改日志中
+ SYNC_FINISH, // 数据同步结束,不在推送
+
+}
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 171c837..aaeea0f 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,6 +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.server.data.SyncCmd;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -18,6 +19,7 @@ public class SyncCmdServerHandler extends SimpleChannelInboundHandler {
@Override
protected void channelRead0(ChannelHandlerContext ctx, SyncCmd syncCmd) throws Exception {
+ logger.info("read syncCmd: {}", syncCmd);
switch (syncCmd.getCmdType()) {
case SyncCmd.PING:
this.handlePing(ctx, syncCmd);
@@ -25,6 +27,9 @@ public class SyncCmdServerHandler extends SimpleChannelInboundHandler {
case SyncCmd.PONG:
this.handlePong(ctx, syncCmd);
break;
+ case SyncCmd.SYNC_BY_HASH:
+ this.handleReqSyncByHash(ctx, syncCmd);
+ break;
}
}
@@ -39,4 +44,13 @@ public class SyncCmdServerHandler extends SimpleChannelInboundHandler {
logger.info("pong: {}", ctx.channel());
}
+ private void handleReqSyncByHash(ChannelHandlerContext ctx, SyncCmd syncCmd) {
+ String clientAddr = syncCmd.getParam1();
+ // 绑定数据同步服务
+ SyncLogByHashService syncLogByHashService = new SyncLogByHashService(clientAddr);
+ ctx.channel().attr(SyncLogByHashService.ATTR_KEY).set(syncLogByHashService);
+ // 开始数据同步
+ syncLogByHashService.startPush();
+ }
+
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncClientSession.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncClientSession.java
index 33a8160..1a401a5 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncClientSession.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncClientSession.java
@@ -1,7 +1,12 @@
package net.sopod.soim.router.datasync.server.session;
+import net.sopod.soim.common.constant.AppConstant;
import net.sopod.soim.router.datasync.server.SyncClient;
+import org.apache.commons.lang3.tuple.Pair;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import java.util.List;
import java.util.concurrent.ConcurrentHashMap;
/**
@@ -10,8 +15,11 @@ import java.util.concurrent.ConcurrentHashMap;
* @author tmy
* @date 2022-05-09 17:45
*/
+@Deprecated
public class SyncClientSession {
+ private static final Logger logger = LoggerFactory.getLogger(SyncClientSession.class);
+
private static final SyncClientSession INSTANCE = new SyncClientSession();
public static SyncClientSession getInstance() {
@@ -21,20 +29,20 @@ public class SyncClientSession {
/**
* 服务端addr,连接 channel
*/
- public final ConcurrentHashMap serverChannel = new ConcurrentHashMap<>();
+ public final ConcurrentHashMap syncLogServers = new ConcurrentHashMap<>();
public void connect(String host, int port) {
SyncClient client = new SyncClient();
try {
client.connect(host, port);
- serverChannel.put(host + ":" + port, client);
+ syncLogServers.put(host + ":" + port, client);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
public void close(String addr) {
- SyncClient syncClient = serverChannel.get(addr);
+ SyncClient syncClient = syncLogServers.get(addr);
if (syncClient != null) {
syncClient.close();
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncServerSession.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncServerSession.java
index eb17871..97bd405 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncServerSession.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncServerSession.java
@@ -1,31 +1,50 @@
package net.sopod.soim.router.datasync.server.session;
+import net.sopod.soim.router.datasync.server.SyncServer;
import org.apache.dubbo.common.utils.ConcurrentHashSet;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import java.nio.channels.Channel;
import java.util.Set;
/**
* SyncServerSession
+ * 作为客户端或服务端主动操作连接
+ * 服务端场景:
+ * 1.获取客户端连接(备份服务器类型,新增服务器同步类型)
+ * 客户端场景:
+ * 1.即将删除:推送到其他几个节点服务器
+ * 2.新增节点:连接几个节点服务器发送拉取数据命令
*
* @author tmy
* @date 2022-05-09 14:49
*/
public class SyncServerSession {
+ private static final Logger logger = LoggerFactory.getLogger(SyncServerSession.class);
+
private static final SyncServerSession INSTANCE = new SyncServerSession();
public static SyncServerSession getInstance() {
return INSTANCE;
}
- /** 已连接客户端 */
+ private SyncServer syncServer;
+
+ /**
+ * 已连接客户端
+ */
public final Set clients = new ConcurrentHashSet<>();
- /** 备份数据节点 */
+ /**
+ * 备份数据节点
+ */
public final Set backupClients = new ConcurrentHashSet<>();
- /** 新增节点发送 Sync 命令后添加到该连接集合 */
+ /**
+ * 新增节点发送 Sync 命令后添加到该连接集合
+ */
public final Set newNodeClients = new ConcurrentHashSet<>();
public void addClient(Channel channel) {
@@ -42,4 +61,16 @@ public class SyncServerSession {
newNodeClients.remove(channel);
}
+ public void start(int port) {
+ this.syncServer = new SyncServer();
+ syncServer.start(port, err -> {
+ throw new IllegalStateException("备份服务启动失败", err);
+ });
+
+ }
+
+ public void close() {
+ this.syncServer.close();
+ }
+
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncSession.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncSession.java
deleted file mode 100644
index 3e24c29..0000000
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncSession.java
+++ /dev/null
@@ -1,61 +0,0 @@
-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
- * 作为客户端或服务端主动操作连接
- * 服务端场景:
- * 1.获取客户端连接(备份服务器类型,新增服务器同步类型)
- * 客户端场景:
- * 1.即将删除:推送到其他几个节点服务器
- * 2.新增节点:连接几个节点服务器发送拉取数据命令
- *
- * @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) {
-
- }
-
- }
-
-}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java b/im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java
index a475660..9c7a9c2 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java
@@ -13,7 +13,7 @@ import net.sopod.soim.router.api.model.RegistryRes;
import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.api.service.UserRouteService;
import net.sopod.soim.router.cache.RouterUserStorage;
-import net.sopod.soim.router.config.ImRouterAppContextHolder;
+import net.sopod.soim.router.config.AppContextHolder;
import org.apache.dubbo.config.annotation.DubboReference;
import org.apache.dubbo.config.annotation.DubboService;
import org.apache.dubbo.rpc.RpcContext;
@@ -57,7 +57,7 @@ public class UserRouteServiceImpl implements UserRouteService {
// 接口返回 im_router_id,后续调用 im-router 负载均衡指向当前router服务
return new RegistryRes()
.setSuccess(true)
- .setImRouterId(ImRouterAppContextHolder.IM_ROUTER_ID);
+ .setImRouterId(AppContextHolder.IM_ROUTER_ID);
}
@Override