From 0241f42cf8858df3e0c3843703b54bf2dbaeb0fc Mon Sep 17 00:00:00 2001 From: tangmingyou <234767776@qq.com> Date: Tue, 10 May 2022 00:39:11 +0800 Subject: [PATCH] =?UTF-8?q?im-router=E6=95=B0=E6=8D=AE=E5=90=8C=E6=AD=A5?= =?UTF-8?q?=E6=9C=8D=E5=8A=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- im-common/pom.xml | 4 + .../router/api/route/ConsistentHashTest.java | 3 +- .../api/route/UidConsistentHashSelector.java | 51 +++++++++-- ...ntextHolder.java => AppContextHolder.java} | 36 ++++++-- .../router/config/ImRouterAppOnReady.java | 74 ++++++++++++--- .../listener/ImRouterAPIExportListener.java | 11 +-- .../datasync/SyncLogMigrateService.java | 90 +++++++++++++++++++ .../soim/router/datasync/SyncService.java | 36 -------- .../router/datasync/server/SyncClient.java | 16 ++++ .../datasync/server/SyncLogByHashService.java | 29 ++++++ .../router/datasync/server/SyncServer.java | 6 +- .../router/datasync/server/data/SyncCmd.java | 23 ++++- .../datasync/server/data/SyncStatus.java | 16 ++++ .../server/handler/SyncCmdServerHandler.java | 14 +++ .../server/session/SyncClientSession.java | 14 ++- .../server/session/SyncServerSession.java | 37 +++++++- .../datasync/server/session/SyncSession.java | 61 ------------- .../router/service/UserRouteServiceImpl.java | 4 +- 18 files changed, 385 insertions(+), 140 deletions(-) rename im-service/im-router/src/main/java/net/sopod/soim/router/config/{ImRouterAppContextHolder.java => AppContextHolder.java} (52%) create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java delete mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncService.java create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncStatus.java delete mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncSession.java 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