Browse Source

im-router数据同步服务

master
tangmingyou 4 years ago
parent
commit
0241f42cf8
  1. 4
      im-common/pom.xml
  2. 3
      im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ConsistentHashTest.java
  3. 51
      im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java
  4. 36
      im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java
  5. 72
      im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java
  6. 11
      im-service/im-router/src/main/java/net/sopod/soim/router/config/listener/ImRouterAPIExportListener.java
  7. 90
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java
  8. 36
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncService.java
  9. 16
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java
  10. 29
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java
  11. 6
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java
  12. 23
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java
  13. 16
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncStatus.java
  14. 14
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java
  15. 14
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncClientSession.java
  16. 37
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncServerSession.java
  17. 61
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncSession.java
  18. 4
      im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java

4
im-common/pom.xml

@ -26,6 +26,10 @@
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
</dependency>
<dependency>
<groupId>io.netty</groupId>
<artifactId>netty-common</artifactId>

3
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();
}

51
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<V> {
*/
private static final int VIRTUAL_NODE_SIZE = 120;
private final TreeMap<Long, V> virtualNodeMap;
private final TreeMap<Long, Pair<String, V>> virtualNodeMap;
private final int identityHashCode;
@ -33,7 +34,17 @@ public class UidConsistentHashSelector<V> {
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<String, V> lastNode;
if (null == (lastNode = this.virtualNodeMap.get(hash))) {
this.virtualNodeMap.put(hash, ImmutablePair.of(serverAddr, value));
} else {
// hash 冲突取排序小的节点
List<String> 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<V> {
throw new IllegalStateException("im-router consistent hash route, ctx uid can not be null!");
}
long hash = hash(uid, 0);
Map.Entry<Long, V> entry = virtualNodeMap.ceilingEntry(hash);
Map.Entry<Long, Pair<String, V>> entry = virtualNodeMap.ceilingEntry(hash);
if (entry == null) {
entry = virtualNodeMap.firstEntry();
}
return entry.getValue();
return entry.getValue().getRight();
}
/**
* 计算加入新节点需要迁移数据的节点
* @param newNode 新节点地址
*/
public Set<V> selectMigrateNodes(String newNode) {
Set<V> 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<Long, Pair<String, V>> entry = this.virtualNodeMap.ceilingEntry(hash);
if (entry == null) {
entry = this.virtualNodeMap.firstEntry();
}
// hash 冲突,取排序最小的一个(老节点)
if (entry.getKey().equals(hash)) {
List<String> 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) {

36
im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppContextHolder.java → 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<URL> 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> T getBean(Class<T> 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;
}
}

72
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,14 +39,32 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
private static final Logger logger = LoggerFactory.getLogger(ImRouterAppOnReady.class);
private static final int SYNC_SERVER_PORT_OFFSET = 1000;
/**
* 同步一致性hash临近节点数据
*/
@Override
public void onApplicationEvent(ApplicationReadyEvent event) {
AppContextHolder.setApplicationContext(event.getApplicationContext());
// 启动数据同步服务
this.startSyncServer();
SyncLogMigrateService syncLogMigrateService = event.getApplicationContext().getBean(SyncLogMigrateService.class);
// 检查集群状态
boolean registryNow = this.checkClusterEnvironment(syncLogMigrateService);
// 服务已可用进行注册
if (registryNow) {
this.doRegistry();
}
}
private void startSyncServer() {
SyncServerSession.getInstance()
.start(AppContextHolder.getAppPort() + SYNC_SERVER_PORT_OFFSET);
}
/**
* 注册 im-router 的API接口服务
@ -51,7 +73,7 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
RegistryManager registryManager = ApplicationModel.defaultModel().getBeanFactory()
.getBean(RegistryManager.class);
Collection<Registry> registries = registryManager.getRegistries();
List<URL> registryInvokerUrls = ImRouterAppContextHolder.getRegistryInvokerUrls();
List<URL> registryInvokerUrls = AppContextHolder.getRegistryInvokerUrls();
if (Collects.isNotEmpty(registries)
&& Collects.isNotEmpty(registryInvokerUrls)) {
for (Registry registry : registries) {
@ -59,7 +81,7 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
// 添加 im-router 服务id参数,生成新的 url
URL url = invokerUrl.addParameter(
DubboConstant.IM_ROUTER_ID_KEY,
ImRouterAppContextHolder.IM_ROUTER_ID
AppContextHolder.IM_ROUTER_ID
);
registry.register(url);
}
@ -74,18 +96,42 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
/**
* 检查集群情况
* @return 是否可立即注册服务
*/
public void checkClusterEnvironment() {
public boolean checkClusterEnvironment(SyncLogMigrateService syncLogMigrateService) {
// TODO 加分布式锁,同一时间单一节点进行数据同步
List<Instance> 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<String, Instance> addrInstanceMap = Collects.collect2Map(clusterImEntryInstance,
Instance::toInetAddr, // 服务地址, 如: 192.168.56.1:3031
new HashMap<>(Collects.mapCapacity(clusterImEntryInstance.size()))
);
UidConsistentHashSelector<Instance> selector = new UidConsistentHashSelector<>(addrInstanceMap, addrInstanceMap.hashCode());
Set<Instance> 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<Pair<String, Integer>> 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<Instance> getClusterImEntryInstance() {
@ -99,6 +145,14 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
return namingService.getAllInstances(AppConstant.APP_IM_ROUTER_NAME);
} catch (NacosException e) {
throw new IllegalStateException("集群状态检查失败:", e);
} finally {
if (namingService != null) {
try {
namingService.shutDown();
} catch (NacosException e) {
logger.info("查询{}服务namingServer关闭失败!", AppConstant.APP_IM_ROUTER_NAME, e);
}
}
}
}

11
im-service/im-router/src/main/java/net/sopod/soim/router/config/listener/ImRouterAPIExportListener.java

@ -1,6 +1,6 @@
package net.sopod.soim.router.config.listener;
import net.sopod.soim.router.config.ImRouterAppContextHolder;
import net.sopod.soim.router.config.AppContextHolder;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.rpc.Exporter;
import org.apache.dubbo.rpc.ExporterListener;
@ -24,10 +24,11 @@ public class ImRouterAPIExportListener implements ExporterListener {
public void exported(Exporter<?> 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());
}
}
}

90
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<Pair<String, Integer>> migrateHosts;
private LinkedList<Pair<String, Integer>> curMigrateHosts;
/**
* 同步数据
* @param migrateHosts 同步数据节点
*/
public void migrateSyncLog(List<Pair<String, Integer>> 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<String, Integer> 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<Long, RouterUser> routerUserMap = RouterUserStorage.getInstance().getRouterUserMap();
for (Map.Entry<Long, RouterUser> entry : routerUserMap.entrySet()) {
Long key = entry.getKey();
RouterUser value = entry.getValue();
Jackson.json().serialize(value);
}
}
public static void main(String[] args) {
ConcurrentHashMap<String, String> map = new ConcurrentHashMap<>();
map.put("1", "A");
map.put("2", "B");
Iterator<String> iterator = map.values().iterator();
System.out.println("a." + iterator.next());
map.remove("2");
System.out.println("b." + iterator.next());
System.out.println(iterator.next());
}
}

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

@ -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<Long, RouterUser> routerUserMap = RouterUserStorage.getInstance().getRouterUserMap();
for (Map.Entry<Long, RouterUser> entry : routerUserMap.entrySet()) {
Long key = entry.getKey();
RouterUser value = entry.getValue();
Jackson.json().serialize(value);
}
}
public static void main(String[] args) {
}
}

16
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();
}

29
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<SyncLogByHashService> 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 开始数据推送
}
}

6
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("服务启动失败");
});
}
}

23
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;

16
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, // 数据同步结束,不在推送
}

14
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<SyncCmd> {
@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<SyncCmd> {
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<SyncCmd> {
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();
}
}

14
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<String, SyncClient> serverChannel = new ConcurrentHashMap<>();
public final ConcurrentHashMap<String, SyncClient> 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();
}

37
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<Channel> clients = new ConcurrentHashSet<>();
/** 备份数据节点 */
/**
* 备份数据节点
*/
public final Set<Channel> backupClients = new ConcurrentHashSet<>();
/** 新增节点发送 Sync 命令后添加到该连接集合 */
/**
* 新增节点发送 Sync 命令后添加到该连接集合
*/
public final Set<Channel> 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();
}
}

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

@ -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) {
}
}
}

4
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

Loading…
Cancel
Save