From 4b04f85111ed3b47afd4a56059e43d3270129ced Mon Sep 17 00:00:00 2001 From: tangmingyou <234767776@qq.com> Date: Fri, 10 Jun 2022 07:47:33 +0800 Subject: [PATCH] im-entry-http monitor --- .../soim/common/constant/AppConstant.java | 10 +++ im-das-api/im-das-message-api/pom.xml | 4 ++ .../entry/http/config/ApplicationOnReady.java | 65 +++++++++++++++++++ .../soim/entry/config/EntryServerConfig.java | 4 ++ .../config/ImEntryAPIExporterListener.java | 6 +- .../soim/entry/config/ImEntryAppContext.java | 47 ++++++++++++++ .../entry/config/ImEntryAppContextHolder.java | 33 ---------- .../soim/entry/config/ImEntryAppOnReady.java | 52 +++++++++++++++ .../SpringApplicationContextInitialed.java | 26 -------- .../handlers/auth/ReqTokenAuthHandler.java | 4 +- .../soim/entry/registry/RegistryService.java | 1 + .../sopod/soim/entry/server/EntryServer.java | 8 +++ .../soim/entry/server/EntryServerRunner.java | 3 +- .../entry/service/TextChatServiceImpl.java | 4 +- im-service-api/im-logic-message-api/pom.xml | 12 ++++ im-service/im-router/pom.xml | 8 +-- ...extHolder.java => ImRouterAppContext.java} | 22 +++---- .../router/config/ImRouterAppOnReady.java | 22 +++---- .../listener/ImRouterAPIExportListener.java | 8 +-- .../router/datasync/SyncLogByHashPusher.java | 4 +- .../datasync/SyncLogMigrateService.java | 6 +- .../server/handler/SyncCmdClientHandler.java | 4 +- .../router/service/UserRouteServiceImpl.java | 14 +--- 23 files changed, 250 insertions(+), 117 deletions(-) create mode 100644 im-entry-http/src/main/java/net/sopod/soim/entry/http/config/ApplicationOnReady.java create mode 100644 im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppContext.java delete mode 100644 im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppContextHolder.java create mode 100644 im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppOnReady.java delete mode 100644 im-entry/src/main/java/net/sopod/soim/entry/config/SpringApplicationContextInitialed.java rename im-service/im-router/src/main/java/net/sopod/soim/router/config/{AppContextHolder.java => ImRouterAppContext.java} (89%) diff --git a/im-common/src/main/java/net/sopod/soim/common/constant/AppConstant.java b/im-common/src/main/java/net/sopod/soim/common/constant/AppConstant.java index 37ee4a1..9eace09 100644 --- a/im-common/src/main/java/net/sopod/soim/common/constant/AppConstant.java +++ b/im-common/src/main/java/net/sopod/soim/common/constant/AppConstant.java @@ -20,4 +20,14 @@ public interface AppConstant { String APP_IM_HTTP_ENTRY_NAME = "im-http-entry"; String APP_IM_DAS_USER_NAME = "im-das-user"; + /** + * im-entry 端口偏移量 + */ + int IM_ENTRY_SERVER_OFFSET = 1000; + + /** + * im-router 端口偏移量 + */ + int IM_ROUTER_SYNC_SERVER_OFFSET = 1000; + } diff --git a/im-das-api/im-das-message-api/pom.xml b/im-das-api/im-das-message-api/pom.xml index 79d9421..eb28bda 100644 --- a/im-das-api/im-das-message-api/pom.xml +++ b/im-das-api/im-das-message-api/pom.xml @@ -40,6 +40,10 @@ org.xerial.snappy snappy-java + + org.mybatis + mybatis + \ No newline at end of file diff --git a/im-entry-http/src/main/java/net/sopod/soim/entry/http/config/ApplicationOnReady.java b/im-entry-http/src/main/java/net/sopod/soim/entry/http/config/ApplicationOnReady.java new file mode 100644 index 0000000..7b8e691 --- /dev/null +++ b/im-entry-http/src/main/java/net/sopod/soim/entry/http/config/ApplicationOnReady.java @@ -0,0 +1,65 @@ +package net.sopod.soim.entry.http.config; + +import com.alibaba.nacos.api.NacosFactory; +import com.alibaba.nacos.api.exception.NacosException; +import com.alibaba.nacos.api.naming.NamingService; +import com.alibaba.nacos.api.naming.listener.Event; +import com.alibaba.nacos.api.naming.listener.NamingEvent; +import com.alibaba.nacos.api.naming.pojo.Instance; +import net.sopod.soim.common.constant.AppConstant; +import org.apache.dubbo.common.URL; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.context.event.ApplicationReadyEvent; +import org.springframework.context.ApplicationListener; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.Ordered; + +import java.util.List; +import java.util.Properties; + +/** + * ApplicationOnReady + * + * @author tmy + * @date 2022-06-09 11:16 + */ +@Configuration +public class ApplicationOnReady implements ApplicationListener, Ordered { + + private static final Logger logger = LoggerFactory.getLogger(ApplicationOnReady.class); + + private final String serverAddress; + + public ApplicationOnReady(@Value("${dubbo.registry.address}") String serverAddress) { + this.serverAddress = URL.valueOf(serverAddress).getAddress(); + } + + @Override + public void onApplicationEvent(ApplicationReadyEvent readyEvent) { + // NamingServer 获取 im-entry 实例,连接上, + Properties properties = new Properties(); + properties.put("serverAddr", serverAddress); + + NamingService namingService = null; + try { + namingService = NacosFactory.createNamingService(properties); + List allInstances = namingService.getAllInstances(AppConstant.APP_IM_ENTRY_NAME); + logger.info("im-entry instance...: {}", allInstances); + namingService.subscribe(AppConstant.APP_IM_ENTRY_NAME, (Event event) -> { + NamingEvent ne = (NamingEvent) event; + logger.info("namingEvent: {}", ne); + }); + + } catch (NacosException e) { + e.printStackTrace(); + } + } + + @Override + public int getOrder() { + return Ordered.HIGHEST_PRECEDENCE; + } + +} diff --git a/im-entry/src/main/java/net/sopod/soim/entry/config/EntryServerConfig.java b/im-entry/src/main/java/net/sopod/soim/entry/config/EntryServerConfig.java index 4de00cc..2bfbbf1 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/config/EntryServerConfig.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/config/EntryServerConfig.java @@ -23,6 +23,10 @@ public class EntryServerConfig { /** entry 所在服务器 ip */ private String ip = "127.0.0.1"; + /** + * 使用dubbo服务端口偏移量 + */ + @Deprecated private Integer port = 8088; private String nacosAddr; diff --git a/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAPIExporterListener.java b/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAPIExporterListener.java index f89a8b3..a85755b 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAPIExporterListener.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAPIExporterListener.java @@ -23,9 +23,9 @@ public class ImEntryAPIExporterListener implements ExporterListener { public void exported(Exporter exporter) throws RpcException { URL invokerUrl = exporter.getInvoker().getUrl(); if (!InjvmProtocol.NAME.equals(invokerUrl.getProtocol())) { - if (ImEntryAppContextHolder.getDubboAppServiceAddr() == null) { - ImEntryAppContextHolder.setDubboAppServiceAddr(invokerUrl.getAddress()); - logger.info("im-entry registry serverAddr: {}", invokerUrl.getAddress()); + if (ImEntryAppContext.getEntryAppAddr() == null) { + ImEntryAppContext.setAppServiceAddr(invokerUrl.getHost(), invokerUrl.getPort()); + logger.info("im-entry dubbo service addr: {}", invokerUrl.getAddress()); } } } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppContext.java b/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppContext.java new file mode 100644 index 0000000..2320f26 --- /dev/null +++ b/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppContext.java @@ -0,0 +1,47 @@ +package net.sopod.soim.entry.config; + +import org.springframework.context.ApplicationContext; + +/** + * ContextHandler + * + * @author tmy + * @date 2022-04-28 15:06 + */ +public class ImEntryAppContext { + + private static ApplicationContext applicationContext; + + private static String appAddr; + + private static String appHost; + + private static int appPort; + + public static void setContext(ApplicationContext applicationContext) { + ImEntryAppContext.applicationContext = applicationContext; + } + + public static T getBean(Class beanType) { + return applicationContext.getBean(beanType); + } + + public static void setAppServiceAddr(String host, int port) { + ImEntryAppContext.appAddr = host + ":" + port; + ImEntryAppContext.appHost = host; + ImEntryAppContext.appPort = port; + } + + public static String getEntryAppAddr() { + return appAddr; + } + + public static String getAppHost() { + return appHost; + } + + public static int getAppPort() { + return appPort; + } + +} diff --git a/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppContextHolder.java b/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppContextHolder.java deleted file mode 100644 index 9619240..0000000 --- a/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppContextHolder.java +++ /dev/null @@ -1,33 +0,0 @@ -package net.sopod.soim.entry.config; - -import org.springframework.context.ApplicationContext; - -/** - * ContextHandler - * - * @author tmy - * @date 2022-04-28 15:06 - */ -public class ImEntryAppContextHolder { - - private static ApplicationContext applicationContext; - - private static String dubboAppServiceAddr; - - public static void setContext(ApplicationContext applicationContext) { - ImEntryAppContextHolder.applicationContext = applicationContext; - } - - public static T getBean(Class beanType) { - return applicationContext.getBean(beanType); - } - - public static void setDubboAppServiceAddr(String dubboAppServiceAddr) { - ImEntryAppContextHolder.dubboAppServiceAddr = dubboAppServiceAddr; - } - - public static String getDubboAppServiceAddr() { - return dubboAppServiceAddr; - } - -} diff --git a/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppOnReady.java b/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppOnReady.java new file mode 100644 index 0000000..09ef44a --- /dev/null +++ b/im-entry/src/main/java/net/sopod/soim/entry/config/ImEntryAppOnReady.java @@ -0,0 +1,52 @@ +package net.sopod.soim.entry.config; + +import net.sopod.soim.common.constant.AppConstant; +import net.sopod.soim.entry.registry.ProtoMessageHandlerRegistry; +import net.sopod.soim.entry.server.EntryServer; +import net.sopod.soim.entry.worker.WorkerGroup; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.boot.context.event.ApplicationReadyEvent; +import org.springframework.context.ApplicationListener; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.Ordered; + +/** + * ImEntryAppOnReady + * + * @author tmy + * @date 2022-06-09 10:45 + */ +@Configuration +public class ImEntryAppOnReady implements ApplicationListener, Ordered { + + private static final Logger logger = LoggerFactory.getLogger(ImEntryAppOnReady.class); + + @Override + public void onApplicationEvent(ApplicationReadyEvent event) { + ImEntryAppContext.setContext(event.getApplicationContext()); + // 注册 protobuf 消息 handler + ProtoMessageHandlerRegistry.registerHandlerWithApplicationContext(event.getApplicationContext()); + + // 启动 entry-server + EntryServerConfig config = ImEntryAppContext.getBean(EntryServerConfig.class); + + EntryServer entryServer = new EntryServer( + config.getName(), + // 使用dubbo 服务端口偏移量 + ImEntryAppContext.getAppPort() + AppConstant.IM_ENTRY_SERVER_OFFSET); + entryServer.startServer(err -> { + logger.error("EntryServer 启动失败:", err); + }); + Runtime.getRuntime().addShutdownHook(new Thread(entryServer::shutdown)); + + // 启动消息消费队列组 + WorkerGroup.init(config.getWorkerSize()); + } + + @Override + public int getOrder() { + return Ordered.HIGHEST_PRECEDENCE; + } + +} diff --git a/im-entry/src/main/java/net/sopod/soim/entry/config/SpringApplicationContextInitialed.java b/im-entry/src/main/java/net/sopod/soim/entry/config/SpringApplicationContextInitialed.java deleted file mode 100644 index c225747..0000000 --- a/im-entry/src/main/java/net/sopod/soim/entry/config/SpringApplicationContextInitialed.java +++ /dev/null @@ -1,26 +0,0 @@ -package net.sopod.soim.entry.config; - -import net.sopod.soim.entry.registry.ProtoMessageHandlerRegistry; -import org.springframework.beans.BeansException; -import org.springframework.context.ApplicationContext; -import org.springframework.context.ApplicationContextAware; -import org.springframework.context.annotation.Configuration; - -/** - * ApplicationContextInitialed - * - * @author tmy - * @date 2022-04-10 22:20 - */ -@Configuration -public class SpringApplicationContextInitialed implements ApplicationContextAware { - - @Override - public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { - ImEntryAppContextHolder.setContext(applicationContext); - - // 注册 protobuf 消息 handler - ProtoMessageHandlerRegistry.registerHandlerWithApplicationContext(applicationContext); - } - -} diff --git a/im-entry/src/main/java/net/sopod/soim/entry/handlers/auth/ReqTokenAuthHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/handlers/auth/ReqTokenAuthHandler.java index 12f01b2..4d42aef 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/handlers/auth/ReqTokenAuthHandler.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/handlers/auth/ReqTokenAuthHandler.java @@ -1,7 +1,7 @@ package net.sopod.soim.entry.handlers.auth; import com.google.protobuf.MessageLite; -import net.sopod.soim.entry.config.ImEntryAppContextHolder; +import net.sopod.soim.entry.config.ImEntryAppContext; import net.sopod.soim.entry.server.handler.ImContext; import net.sopod.soim.entry.server.handler.NetUserMessageHandler; import net.sopod.soim.entry.server.session.Account; @@ -49,7 +49,7 @@ public class ReqTokenAuthHandler extends NetUserMessageHandlerim-logic-common ${soim.version} + + + + + + + + + + + + diff --git a/im-service/im-router/pom.xml b/im-service/im-router/pom.xml index 342d88c..2e444d1 100644 --- a/im-service/im-router/pom.xml +++ b/im-service/im-router/pom.xml @@ -89,10 +89,10 @@ - - - - + + org.xerial.snappy + snappy-java + 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/ImRouterAppContext.java similarity index 89% rename from im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java rename to im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppContext.java index 9c3c4ae..730ccf6 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/ImRouterAppContext.java @@ -32,9 +32,9 @@ import java.util.concurrent.CopyOnWriteArrayList; * @author tmy * @date 2022-05-04 09:21 */ -public class AppContextHolder { +public class ImRouterAppContext { - private static final Logger logger = LoggerFactory.getLogger(AppContextHolder.class); + private static final Logger logger = LoggerFactory.getLogger(ImRouterAppContext.class); private static ApplicationContext applicationContext; @@ -55,7 +55,7 @@ public class AppContextHolder { } public static void setApplicationContext(ApplicationContext applicationContext) { - AppContextHolder.applicationContext = applicationContext; + ImRouterAppContext.applicationContext = applicationContext; } public static T getBean(Class type) { @@ -71,13 +71,13 @@ public class AppContextHolder { } public static void setAppServiceAddr(String host, int port) { - AppContextHolder.appAddr = host + ":" + port; - AppContextHolder.appHost = host; - AppContextHolder.appPort = port; + ImRouterAppContext.appAddr = host + ":" + port; + ImRouterAppContext.appHost = host; + ImRouterAppContext.appPort = port; } public static void setServiceDiscoveryRegistryAddr(String discoveryAddr) { - AppContextHolder.discoveryAddr = discoveryAddr; + ImRouterAppContext.discoveryAddr = discoveryAddr; } public static String getDiscoveryAddr() { @@ -104,11 +104,11 @@ public class AppContextHolder { */ public static List getClusterImRouterInstance() { String discoveryAddr; - if (null == (discoveryAddr = AppContextHolder.getDiscoveryAddr())) { + if (null == (discoveryAddr = ImRouterAppContext.getDiscoveryAddr())) { return Collections.emptyList(); } if (namingService == null) { - synchronized (AppContextHolder.class) { + synchronized (ImRouterAppContext.class) { if (namingService == null) { Properties properties = new Properties(); properties.put("serverAddr", discoveryAddr); @@ -143,7 +143,7 @@ public class AppContextHolder { RegistryManager registryManager = ApplicationModel.defaultModel().getBeanFactory() .getBean(RegistryManager.class); Collection registries = registryManager.getRegistries(); - List registryInvokerUrls = AppContextHolder.getRegistryInvokerUrls(); + List registryInvokerUrls = ImRouterAppContext.getRegistryInvokerUrls(); if (Collects.isNotEmpty(registries) && Collects.isNotEmpty(registryInvokerUrls)) { for (Registry registry : registries) { @@ -151,7 +151,7 @@ public class AppContextHolder { // 添加 im-router 服务id参数,生成新的 url URL url = invokerUrl.addParameter( DubboConstant.IM_ROUTER_ID_KEY, - AppContextHolder.IM_ROUTER_ID + ImRouterAppContext.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 b92cbb5..dd3f497 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 @@ -5,7 +5,6 @@ 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.router.api.route.UidConsistentHashSelector; import net.sopod.soim.router.cache.RouterUser; @@ -14,7 +13,6 @@ 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; @@ -29,7 +27,6 @@ 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; @@ -45,14 +42,12 @@ public class ImRouterAppOnReady implements ApplicationListener clusterImEntryInstance = AppContextHolder.getClusterImRouterInstance(); + List clusterImEntryInstance = ImRouterAppContext.getClusterImRouterInstance(); if (Collects.isEmpty(clusterImEntryInstance)) { logger.info("当前集群无{}节点,直接启动", AppConstant.APP_IM_ROUTER_NAME); return true; @@ -153,7 +148,7 @@ public class ImRouterAppOnReady implements ApplicationListener(Collects.mapCapacity(clusterImEntryInstance.size())) ); UidConsistentHashSelector selector = new UidConsistentHashSelector<>(addrInstanceMap, addrInstanceMap.hashCode()); - Set migrateNodes = selector.selectMigrateNodes(AppContextHolder.getAppAddr()); + Set migrateNodes = selector.selectMigrateNodes(ImRouterAppContext.getAppAddr()); logger.info("需迁移数据{}节点: {} of {}", AppConstant.APP_IM_ROUTER_NAME, migrateNodes.size(), @@ -163,7 +158,7 @@ public class ImRouterAppOnReady implements ApplicationListener> migrateHosts = migrateNodes.stream() - .map(instance -> ImmutablePair.of(instance.getIp(), instance.getPort() + SYNC_SERVER_PORT_OFFSET)) + .map(instance -> ImmutablePair.of(instance.getIp(), instance.getPort() + AppConstant.IM_ROUTER_SYNC_SERVER_OFFSET)) .collect(Collectors.toList()); // 开始同步数据 syncLogMigrateService.migrateSyncLog(migrateHosts); @@ -186,6 +181,7 @@ public class ImRouterAppOnReady implements ApplicationListener exporter) throws RpcException { URL invokerUrl = exporter.getInvoker().getUrl(); if (!InjvmProtocol.NAME.equals(invokerUrl.getProtocol())) { - AppContextHolder.addRegistryInvokerUrl(invokerUrl); - if (AppContextHolder.getAppAddr() == null) { + ImRouterAppContext.addRegistryInvokerUrl(invokerUrl); + if (ImRouterAppContext.getAppAddr() == null) { // 保存服务注册地址 - AppContextHolder.setAppServiceAddr(invokerUrl.getHost(), invokerUrl.getPort()); + ImRouterAppContext.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/SyncLogByHashPusher.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashPusher.java index 9c09047..4c56af3 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashPusher.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashPusher.java @@ -4,7 +4,7 @@ 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.config.ImRouterAppContext; import net.sopod.soim.router.datasync.server.data.SyncCmd; import net.sopod.soim.router.datasync.server.data.SyncLog; import org.slf4j.Logger; @@ -41,7 +41,7 @@ public class SyncLogByHashPusher extends DataChangeTrigger.DataKeySyncLogSubscri // 构建hash环匹配要迁移的数据 Map twoNodes = new HashMap<>(); twoNodes.put(newNodeAddr, newNodeAddr); - twoNodes.put(AppContextHolder.getAppAddr(), AppContextHolder.getAppAddr()); + twoNodes.put(ImRouterAppContext.getAppAddr(), ImRouterAppContext.getAppAddr()); selector = new UidConsistentHashSelector<>(twoNodes, twoNodes.hashCode()); // 添加数据变化监听 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 575c11c..7d851eb 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 @@ -5,7 +5,7 @@ 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.config.ImRouterAppContext; import net.sopod.soim.router.datasync.server.SyncClient; import net.sopod.soim.router.datasync.server.data.SyncCmd; import org.apache.commons.lang3.tuple.Pair; @@ -64,7 +64,7 @@ public class SyncLogMigrateService { try { this.curClient = new SyncClient(); this.curClient.connect(nextHost.getLeft(), nextHost.getRight()); - this.curClient.syncLogByHash(AppContextHolder.getAppAddr()); + this.curClient.syncLogByHash(ImRouterAppContext.getAppAddr()); } catch (InterruptedException e) { logger.error("节点{}:{}连接失败, 跳过!", nextHost.getLeft(), nextHost.getRight(), e); // 同步下一个节点 @@ -92,7 +92,7 @@ public class SyncLogMigrateService { */ private void allHostSyncFinish() { logger.info("所有节点数据同步完成:执行注册服务...."); - AppContextHolder.doRegistry(); + ImRouterAppContext.doRegistry(); logger.info("注册服务成功"); // 通知服务端同步结束,关闭服务端连接 for (SyncClient migrateClient : migrateClients) { 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 eb483e5..090cda7 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,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.config.AppContextHolder; +import net.sopod.soim.router.config.ImRouterAppContext; import net.sopod.soim.router.datasync.SyncLogMigrateService; import net.sopod.soim.router.datasync.server.data.SyncCmd; import org.slf4j.Logger; @@ -46,7 +46,7 @@ public class SyncCmdClientHandler extends SimpleChannelInboundHandler { private void handleSyncEnd(ChannelHandlerContext ctx, SyncCmd syncCmd) { logger.info("sync end: {}", ctx.channel()); - SyncLogMigrateService migrateService = AppContextHolder.getBean(SyncLogMigrateService.class); + SyncLogMigrateService migrateService = ImRouterAppContext.getBean(SyncLogMigrateService.class); migrateService.curHostSyncEnd(); } 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 f0c5115..ca68e75 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 @@ -3,23 +3,18 @@ package net.sopod.soim.router.service; import net.sopod.soim.common.constant.DubboConstant; import net.sopod.soim.common.util.ImClock; import net.sopod.soim.common.util.StringUtil; -import net.sopod.soim.das.common.config.LogicTables; -import net.sopod.soim.das.message.api.entity.ImMessage; import net.sopod.soim.das.user.api.model.entity.ImUser; -import net.sopod.soim.das.message.api.service.DasMQMessagePersistentService; import net.sopod.soim.das.user.api.service.FriendDas; import net.sopod.soim.das.user.api.service.UserDas; import net.sopod.soim.entry.api.service.OnlineUserService; import net.sopod.soim.entry.api.service.TextChatService; import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator; -import net.sopod.soim.logic.common.model.TextChat; import net.sopod.soim.logic.common.model.UserInfo; -import net.sopod.soim.logic.common.util.RpcContextUtil; 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.RouterUser; import net.sopod.soim.router.cache.RouterUserStorage; -import net.sopod.soim.router.config.AppContextHolder; +import net.sopod.soim.router.config.ImRouterAppContext; import org.apache.dubbo.config.annotation.DubboReference; import org.apache.dubbo.config.annotation.DubboService; import org.apache.dubbo.rpc.RpcContext; @@ -58,9 +53,6 @@ public class UserRouteServiceImpl implements UserRouteService { @Resource private SegmentIdGenerator segmentIdGenerator; - @Resource - private DasMQMessagePersistentService dasMQMessagePersistentService; - @Override public RegistryRes registryUserEntry(Long uid, String imEntryAddr) { ImUser imUser = userDas.getUserById(uid); @@ -74,7 +66,7 @@ public class UserRouteServiceImpl implements UserRouteService { // 接口返回 im_router_id,后续调用 im-router 负载均衡指向当前router服务 return new RegistryRes() .setSuccess(true) - .setImRouterId(AppContextHolder.IM_ROUTER_ID); + .setImRouterId(ImRouterAppContext.IM_ROUTER_ID); } @Override