From a816ddc99483835d52c45cd945b5494ae17c6d67 Mon Sep 17 00:00:00 2001
From: tangmingyou <234767776@qq.com>
Date: Thu, 12 May 2022 23:27:15 +0800
Subject: [PATCH] =?UTF-8?q?im-router=E6=96=B0=E8=8A=82=E7=82=B9=E6=95=B0?=
=?UTF-8?q?=E6=8D=AE=E5=90=8C=E6=AD=A5?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
README.md | 3 +-
.../net/sopod/soim/common/util/Jackson.java | 8 ---
im-service-api/pom.xml | 41 +++++++++++
.../router/datasync/DataChangeTrigger.java | 57 +++++++++------
.../router/datasync/DataSyncProxyFactory.java | 71 ++++++++++---------
...hService.java => SyncLogByHashPusher.java} | 52 +++++++-------
.../datasync/SyncLogMigrateService.java | 28 ++++----
.../router/datasync/server/SyncClient.java | 42 +----------
.../router/datasync/server/data/SyncCmd.java | 11 ++-
.../server/handler/SyncCmdClientHandler.java | 2 +-
.../server/handler/SyncCmdServerHandler.java | 27 +++++--
.../server/handler/SyncLogClientHandler.java | 4 +-
.../src/main/resources/application.yml | 3 +-
pom.xml | 1 +
14 files changed, 191 insertions(+), 159 deletions(-)
rename im-service/im-router/src/main/java/net/sopod/soim/router/datasync/{SyncLogByHashService.java => SyncLogByHashPusher.java} (78%)
diff --git a/README.md b/README.md
index 8d5ac8d..2ef8c26 100644
--- a/README.md
+++ b/README.md
@@ -2,7 +2,7 @@
### 开发计划
-- (05-08~05-11)router 一致性 hash,新增、删除、备份处理
+- (05-08~05-11)router 一致性 hash,新增、删除、(备份处理)[滞后]
- 数据同步删除后,保留id一段时间,有请求进行重定向/转发
- 新增:重算hash,发起数据同步,接收其他节点推送数据及数据更改日志,注册服务,其他节点删除数据
- 正常删除节点:重算 hash,数据和更改日志推送给其他节点,取消注册,关闭服务
@@ -43,6 +43,7 @@
一致性hash算法能否解决数据迁移的问题?
https://www.zhihu.com/question/521159623
+dubbo + protobuf 兼容到 2.7.15
### TODO
diff --git a/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java b/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java
index acd57d1..4797141 100644
--- a/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java
+++ b/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java
@@ -318,14 +318,6 @@ public class Jackson {
return objectMapper.getTypeFactory().constructCollectionLikeType(List.class, elementClass);
}
- public ObjectMapper getObjectMapper() {
- return objectMapper;
- }
-
- public TypeFactory getTypeFactory() {
- return objectMapper.getTypeFactory();
- }
-
/**
* @param content 序列化内容
* @param type 方法参数类型等可能包含泛型类型
diff --git a/im-service-api/pom.xml b/im-service-api/pom.xml
index ac25479..cd41810 100644
--- a/im-service-api/pom.xml
+++ b/im-service-api/pom.xml
@@ -18,4 +18,45 @@
im-logic-user-api
+
+
+
+ kr.motd.maven
+ os-maven-plugin
+ 1.6.1
+
+
+
+
+
+ org.xolstice.maven.plugins
+ protobuf-maven-plugin
+ 0.6.1
+
+ com.google.protobuf:protoc:3.19.4:exe:${os.detected.classifier}
+ src/main/java
+ false
+
+
+
+
+ dubbo
+ org.apache.dubbo
+ dubbo-compiler
+ ${dubbo.compiler.version}
+ org.apache.dubbo.gen.dubbo.Dubbo3Generator
+
+
+
+
+
+
+ compile
+ test-compile
+
+
+
+
+
+
\ No newline at end of file
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java
index 098fb22..cd81d24 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java
@@ -1,13 +1,13 @@
package net.sopod.soim.router.datasync;
-import net.sopod.soim.router.config.ImRouterAppOnReady;
import net.sopod.soim.router.datasync.server.data.SyncLog;
+import org.apache.dubbo.common.utils.ConcurrentHashSet;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.atomic.AtomicInteger;
/**
@@ -46,15 +46,15 @@ public class DataChangeTrigger {
/**
* 监听者列表
*/
- private final List subscribes = new ArrayList<>();
+ private final List subscribes = new CopyOnWriteArrayList<>();
- private List getSubscribes(String dataKey) {
+ private List getSubscribes(SyncTypes.SyncType> syncType, String dataKey) {
if (subscribes.isEmpty()) {
return Collections.emptyList();
}
List acceptSubs = new ArrayList<>(5);
for (SyncLogSubscribe subscribe : subscribes) {
- if (subscribe.accept(dataKey)) {
+ if (subscribe.accept(syncType, dataKey)) {
acceptSubs.add(subscribe);
}
}
@@ -63,7 +63,7 @@ public class DataChangeTrigger {
public void onUpdate(SyncTypes.SyncType syncType, String dataKey, String method, Object[] args) {
List acceptSubs;
- if ((acceptSubs = getSubscribes(dataKey)).isEmpty()) {
+ if ((acceptSubs = getSubscribes(syncType, dataKey)).isEmpty()) {
return;
}
// 序列化 args,避免后续更改
@@ -71,21 +71,21 @@ public class DataChangeTrigger {
.setDataKey(dataKey)
.setMethod(method)
.setArgs(args);
- publishLog(updateLog, acceptSubs);
+ publishLog(syncType, updateLog, acceptSubs);
}
public void onAdd(SyncTypes.SyncType syncType, T data) {
List acceptSubs;
String dataKey = syncType.getDataKey(data);
- if ((acceptSubs = getSubscribes(dataKey)).isEmpty()) {
+ if ((acceptSubs = getSubscribes(syncType, dataKey)).isEmpty()) {
return;
}
AtomicInteger seqCounter = getSeqCounter(dataKey);
// 序列化 data,避免后续更改
SyncLog.AddLog addLog = SyncLog.addLog(seqCounter.getAndIncrement(), syncType)
.addData(data);
- publishLog(addLog, acceptSubs);
+ publishLog(syncType, addLog, acceptSubs);
}
/**
@@ -93,18 +93,22 @@ public class DataChangeTrigger {
*/
public void onRemove(SyncTypes.SyncType syncType, String dataKey) {
List acceptSubs;
- if ((acceptSubs = getSubscribes(dataKey)).isEmpty()) {
+ if ((acceptSubs = getSubscribes(syncType, dataKey)).isEmpty()) {
return;
}
SyncLog.RemoveLog removeLog = SyncLog.removeLog(getSeq(dataKey), syncType)
.setDataKey(dataKey);
- publishLog(removeLog, acceptSubs);
+ publishLog(syncType, removeLog, acceptSubs);
}
- private void publishLog(SyncLog log, List acceptSubs) {
+ /**
+ * 为监听者推送数据更新消息
+ * TODO 异步处理
+ */
+ private void publishLog(SyncTypes.SyncType> syncType, SyncLog log, List acceptSubs) {
for (SyncLogSubscribe acceptSub : acceptSubs) {
try {
- acceptSub.onSyncLog(log);
+ acceptSub.onSyncLog(syncType, log);
} catch (Exception e) {
logger.error("监听者 {} 执行错误", acceptSub, e);
}
@@ -122,10 +126,14 @@ public class DataChangeTrigger {
/**
* 添加更改日志监听者
*/
- public void addSubscribe(SyncLogSubscribe syncLogSubscribe) {
+ public void subscribe(SyncLogSubscribe syncLogSubscribe) {
subscribes.add(syncLogSubscribe);
}
+ public void unsubscribe(SyncLogSubscribe syncLogSubscribe) {
+ subscribes.remove(syncLogSubscribe);
+ }
+
/**
* 删除暂存数据字段等...
*/
@@ -135,31 +143,34 @@ public class DataChangeTrigger {
public static interface SyncLogSubscribe {
- public abstract boolean accept(String dataKey);
+ public abstract boolean accept(SyncTypes.SyncType> syncType, String dataKey);
- public abstract void onSyncLog(SyncLog syncLog);
+ public abstract void onSyncLog(SyncTypes.SyncType> syncType, SyncLog syncLog);
}
public static abstract class DataKeySyncLogSubscribe implements SyncLogSubscribe {
- private final Set subscribeDataKeys;
+ protected final ConcurrentHashMap, Set> syncedTypeDataKyes;
public DataKeySyncLogSubscribe() {
- this.subscribeDataKeys = new HashSet<>(1024);
+ this.syncedTypeDataKyes = new ConcurrentHashMap<>(6);
}
@Override
- public boolean accept(String dataKey) {
- return subscribeDataKeys.contains(dataKey);
+ public boolean accept(SyncTypes.SyncType> syncType, String dataKey) {
+ return syncedTypeDataKyes.getOrDefault(syncType, Collections.emptySet()).contains(dataKey);
}
- public void addSubscribeDataKey(String dataKey) {
- subscribeDataKeys.add(dataKey);
+ public void addSubscribeDataKey(SyncTypes.SyncType> syncType, String dataKey) {
+ syncedTypeDataKyes.computeIfAbsent(
+ syncType,
+ type -> new ConcurrentHashSet<>()
+ ).add(dataKey);
}
@Override
- public abstract void onSyncLog(SyncLog syncLog);
+ public abstract void onSyncLog(SyncTypes.SyncType> syncType, SyncLog syncLog);
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
index 5e0a02d..df84b21 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
@@ -34,6 +34,8 @@ public class DataSyncProxyFactory {
*/
private static final Map, Set> typeUpdaterMethodsCache = new ConcurrentHashMap<>();
+ private static final Map, DataSyncProxyCallback>> syncTypeCallbackCache = new ConcurrentHashMap<>();
+
public static T newProxyInstance(SyncTypes.SyncType syncType) {
return newProxyInstance(syncType, null);
}
@@ -50,53 +52,56 @@ public class DataSyncProxyFactory {
throw new IllegalStateException(type + "实例创建失败,没有无参构造函数", e);
}
- // 获取查询更新方法列表
- Set updaterMethods = typeUpdaterMethodsCache.computeIfAbsent(type, dataType -> {
- Method[] methods = dataType.getMethods();
- Set updaterMethodNames = new HashSet<>();
- for (Method m : methods) {
- SyncIgnore syncIgnore = m.getDeclaredAnnotation(SyncIgnore.class);
- String methodName = m.getName();
- if (syncIgnore != null) {
- continue;
- }
- boolean isUpdater = m.getDeclaringClass() != Object.class
- && instance.isUpdateMethod(methodName);
- if (isUpdater) {
- if (updaterMethodNames.contains(methodName)) {
- throw new IllegalStateException(dataType + "重复的数据更新方法名" + methodName);
- }
- // 检查参数可序列化
- Class>[] paramTypes = m.getParameterTypes();
- for (int i = 0; i < paramTypes.length; i++) {
- if (!Jackson.json().canSerialize(paramTypes[i])) {
- throw new IllegalStateException(String.format("类:%s 更新方法:%s 第%di个参数不可序列化",
- instance.getClass().getName(), methodName, i + 1));
- }
- }
- updaterMethodNames.add(methodName);
- }
- }
- logger.info("{} 更新方法列表: {}", dataType, updaterMethodNames);
- return updaterMethodNames;
- });
-
// 创建代理对象
Enhancer enhancer = new Enhancer();
enhancer.setSuperclass(type);
- enhancer.setCallback(new DataSyncProxyCallback<>(syncType, updaterMethods));
+ DataSyncProxyCallback> dataSyncProxyCallback = syncTypeCallbackCache
+ .computeIfAbsent(syncType, sType -> {
+ Set updaterMethods = getUpdaterMethods(instance, type);
+ return new DataSyncProxyCallback<>(syncType, updaterMethods);
+ });
+ enhancer.setCallback(dataSyncProxyCallback);
T proxyObj = (T) enhancer.create();
if (source != null) {
try {
cloneFields(source, proxyObj);
} catch (Exception e) {
- // logger.error("克隆对象失败:", e);
throw new IllegalStateException("克隆对象失败", e);
}
}
return proxyObj;
}
+ private static Set getUpdaterMethods(T instance, Class dataType) {
+ Method[] methods = dataType.getMethods();
+ Set updaterMethodNames = new HashSet<>();
+ for (Method m : methods) {
+ SyncIgnore syncIgnore = m.getDeclaredAnnotation(SyncIgnore.class);
+ String methodName = m.getName();
+ if (syncIgnore != null) {
+ continue;
+ }
+ boolean isUpdater = m.getDeclaringClass() != Object.class
+ && instance.isUpdateMethod(methodName);
+ if (isUpdater) {
+ if (updaterMethodNames.contains(methodName)) {
+ throw new IllegalStateException(dataType + "重复的数据更新方法名" + methodName);
+ }
+ // 检查参数可序列化
+ Class>[] paramTypes = m.getParameterTypes();
+ for (int i = 0; i < paramTypes.length; i++) {
+ if (!Jackson.json().canSerialize(paramTypes[i])) {
+ throw new IllegalStateException(String.format("类:%s 更新方法:%s 第%di个参数不可序列化",
+ instance.getClass().getName(), methodName, i + 1));
+ }
+ }
+ updaterMethodNames.add(methodName);
+ }
+ }
+ logger.info("{} 更新方法列表: {}", dataType, updaterMethodNames);
+ return updaterMethodNames;
+ }
+
private static void cloneFields(T source, T proxyObj) throws Exception {
cloneFields0(source.getClass(), source, proxyObj);
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashService.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashPusher.java
similarity index 78%
rename from im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashService.java
rename to im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashPusher.java
index 778aecd..9c09047 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashService.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashPusher.java
@@ -7,13 +7,11 @@ import net.sopod.soim.router.api.route.UidConsistentHashSelector;
import net.sopod.soim.router.config.AppContextHolder;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import net.sopod.soim.router.datasync.server.data.SyncLog;
-import org.apache.commons.lang3.tuple.Pair;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.lang.ref.WeakReference;
import java.util.*;
-import java.util.concurrent.atomic.AtomicInteger;
/**
* SyncLogByHashService
@@ -21,22 +19,23 @@ import java.util.concurrent.atomic.AtomicInteger;
* @author tmy
* @date 2022-05-10 00:30
*/
-public class SyncLogByHashService extends DataChangeTrigger.DataKeySyncLogSubscribe {
+public class SyncLogByHashPusher extends DataChangeTrigger.DataKeySyncLogSubscribe {
- private static final Logger logger = LoggerFactory.getLogger(SyncLogByHashService.class);
+ private static final Logger logger = LoggerFactory.getLogger(SyncLogByHashPusher.class);
- public static final AttributeKey ATTR_KEY = AttributeKey
- .valueOf(SyncLogByHashService.class, "SYNC_LOG_BY_HASH_SERVICE");
+ public static final AttributeKey ATTR_KEY = AttributeKey
+ .valueOf(SyncLogByHashPusher.class, "SYNC_LOG_BY_HASH_PUSHER");
private final WeakReference clientChannel;
+ /**
+ * 新im-router节点地址
+ */
private final String newNodeAddr;
private final UidConsistentHashSelector selector;
- // Set syncedDataKeys = new HashSet<>();
-
- public SyncLogByHashService(Channel clientChannel, String newNodeAddr) {
+ public SyncLogByHashPusher(Channel clientChannel, String newNodeAddr) {
this.clientChannel = new WeakReference<>(clientChannel);
this.newNodeAddr = newNodeAddr;
// 构建hash环匹配要迁移的数据
@@ -46,7 +45,7 @@ public class SyncLogByHashService extends DataChangeTrigger.DataKeySyncLogSubscr
selector = new UidConsistentHashSelector<>(twoNodes, twoNodes.hashCode());
// 添加数据变化监听
- DataChangeTrigger.instance().addSubscribe(this);
+ DataChangeTrigger.instance().subscribe(this);
}
private Map, SyncTypes.SyncType extends DataSync>> storages;
@@ -92,7 +91,7 @@ public class SyncLogByHashService extends DataChangeTrigger.DataKeySyncLogSubscr
// TODO DataChangeTrigger.instance().subscribe(syncType, dataKey)
// 注册更改日志
- super.addSubscribeDataKey(dataKey);
+ super.addSubscribeDataKey(syncType, dataKey);
// TODO 解锁
count ++; totalCount ++;
@@ -152,32 +151,35 @@ public class SyncLogByHashService extends DataChangeTrigger.DataKeySyncLogSubscr
Channel channel = clientChannel.get();
if (channel != null && channel.isActive()) {
SyncCmd syncEndCmd = new SyncCmd();
- syncEndCmd.setCmdType(SyncCmd.SYNC_END);
+ syncEndCmd.setCmdType(SyncCmd.SYNC_FULL_END);
channel.writeAndFlush(syncEndCmd);
}
}
+ public void releaseResource() {
+ // 移除当前环境已同步的数据
+ int count = 0;
+ for (Map.Entry, Set> entry : syncedTypeDataKyes.entrySet()) {
+ SyncTypes.SyncType> syncType = entry.getKey();
+ for (String dataKey : entry.getValue()) {
+ boolean removed = syncType.removeData(dataKey);
+ if (removed) {
+ count++;
+ }
+ }
+ }
+ logger.info("{}条已同步数据已移除", count);
+ }
+
/**
* 修改数据日志
*/
@Override
- public void onSyncLog(SyncLog syncLog) {
+ public void onSyncLog(SyncTypes.SyncType> syncType, SyncLog syncLog) {
Channel channel = clientChannel.get();
if (channel != null && channel.isActive()) {
channel.writeAndFlush(syncLog);
}
}
- public static void main(String[] args) {
- Map> twoNodes = new HashMap<>();
- twoNodes.put("192.168.101.69:3032", Pair.of("192.168.101.69:3032", new AtomicInteger()));
- twoNodes.put("192.168.101.69:3031", Pair.of("192.168.101.69:3031", new AtomicInteger()));
- UidConsistentHashSelector> selector = new UidConsistentHashSelector<>(twoNodes, twoNodes.hashCode());
- for (int i = 10000; i < 11000; i++) {
- Pair pair = selector.select(i + "");
- pair.getRight().incrementAndGet();
- }
- System.out.println(twoNodes);
- }
-
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java
index 83ab332..575c11c 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
@@ -7,13 +7,13 @@ 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 net.sopod.soim.router.datasync.server.data.SyncCmd;
import org.apache.commons.lang3.tuple.Pair;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
import java.util.*;
-import java.util.concurrent.ConcurrentHashMap;
/**
* SyncLogMigrateService
@@ -56,13 +56,11 @@ public class SyncLogMigrateService {
throw new IllegalStateException("当前有客户端正在全量同步数据");
}
if (this.curMigrateHosts.isEmpty()) {
- return false;
- }
- Pair nextHost = this.curMigrateHosts.removeFirst();
- if (nextHost == null) {
+ // 同步结束
this.allHostSyncFinish();
return true;
}
+ Pair nextHost = this.curMigrateHosts.removeFirst();
try {
this.curClient = new SyncClient();
this.curClient.connect(nextHost.getLeft(), nextHost.getRight());
@@ -96,6 +94,15 @@ public class SyncLogMigrateService {
logger.info("所有节点数据同步完成:执行注册服务....");
AppContextHolder.doRegistry();
logger.info("注册服务成功");
+ // 通知服务端同步结束,关闭服务端连接
+ for (SyncClient migrateClient : migrateClients) {
+ SyncCmd syncFinishCmd = new SyncCmd();
+ syncFinishCmd.setCmdType(SyncCmd.SYNC_FINISH_CLOSE);
+ migrateClient.write(syncFinishCmd);
+ }
+ for (SyncClient migratedClient : migrateClients) {
+ migratedClient.close();
+ }
}
@Deprecated
@@ -108,15 +115,4 @@ public class SyncLogMigrateService {
}
}
- 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/server/SyncClient.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java
index f9e1fd1..2f47777 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
@@ -68,21 +68,6 @@ public class SyncClient {
clientChannel.write(syncLog);
}
- public void setAttr(AttributeKey key, T val) {
- clientChannel.attr(key).set(val);
- }
-
- public void getAttr(AttributeKey key) {
- clientChannel.attr(key).get();
- }
-
- /**
- * 全量同步数据
- */
- public void syncLogByHash() {
-
- }
-
/**
* 通过计算一致性 hash 同步数据
* @param currentAddr 当前新节点地址,hash(currentAddr)
@@ -93,33 +78,8 @@ public class SyncClient {
}
public void close() {
+ this.clientChannel.close();
this.group.shutdownGracefully();
}
-
- public static void main(String[] args) throws InterruptedException {
- SyncClient client = new SyncClient();
- client.connect("127.0.0.1", 9999);
- // SyncCmd syncCmd = new SyncCmd().setCmdType(SyncCmd.SYNC_LOG);
- // new SyncLog()
- RouterUser user = new RouterUser()
- .setUid(12312L)
- .setAccount("蓝水云烟")
- .setImEntryAddr("127.0.0.1")
- .setIsOnline(false)
- .setOnlineTime(10086L);
- RouterUser user2 = new RouterUser()
- .setUid(10010L)
- .setAccount("百战成诗")
- .setImEntryAddr("192.168.1.101")
- .setIsOnline(true)
- .setOnlineTime(16161L);
- SyncLog.AddLog addLog = SyncLog.addLog(1, SyncTypes.ROUTER_USER)
- .addData(user)
- .addData(user2);
- client.clientChannel.writeAndFlush(addLog);
-
- Thread.sleep(10000);
- client.close();
- }
}
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 f23b9a3..dc297fc 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
@@ -18,7 +18,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
@Accessors(chain = true)
public class SyncCmd {
- private static final short MAGIC = 0x7a21;
+ private static final short MAGIC = 0x544D;
public static final int PING = 1;
@@ -40,9 +40,14 @@ public class SyncCmd {
public static final int SYNC_BY_HASH_ACK = 5;
/**
- * 同步结束命令
+ * 全量同步结束命令
*/
- public static final int SYNC_END = 8;
+ public static final int SYNC_FULL_END = 8;
+
+ /**
+ * 结束同步命令,释放资源,关闭连接
+ */
+ public static final int SYNC_FINISH_CLOSE = 9;
/**
* SyncLog 推送命令
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 f4ff225..eb483e5 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
@@ -27,7 +27,7 @@ public class SyncCmdClientHandler extends SimpleChannelInboundHandler {
case SyncCmd.PONG:
this.handlePong(ctx, syncCmd);
break;
- case SyncCmd.SYNC_END:
+ case SyncCmd.SYNC_FULL_END:
this.handleSyncEnd(ctx, syncCmd);
break;
}
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 b8d0937..42632ec 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java
@@ -2,7 +2,8 @@ package net.sopod.soim.router.datasync.server.handler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
-import net.sopod.soim.router.datasync.SyncLogByHashService;
+import net.sopod.soim.router.datasync.DataChangeTrigger;
+import net.sopod.soim.router.datasync.SyncLogByHashPusher;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -33,6 +34,9 @@ public class SyncCmdServerHandler extends SimpleChannelInboundHandler {
case SyncCmd.SYNC_BY_HASH_ACK:
this.handleReqSyncByHashAck(ctx, syncCmd);
break;
+ case SyncCmd.SYNC_FINISH_CLOSE:
+ this.handleSyncFinishClose(ctx, syncCmd);
+ break;
}
}
@@ -50,18 +54,29 @@ public class SyncCmdServerHandler extends SimpleChannelInboundHandler {
private void handleReqSyncByHash(ChannelHandlerContext ctx, SyncCmd syncCmd) {
String clientAddr = syncCmd.getParam1();
// 绑定数据同步服务
- SyncLogByHashService syncLogByHashService = new SyncLogByHashService(ctx.channel(), clientAddr);
- ctx.channel().attr(SyncLogByHashService.ATTR_KEY).set(syncLogByHashService);
+ SyncLogByHashPusher syncLogByHashPusher = new SyncLogByHashPusher(ctx.channel(), clientAddr);
+ ctx.channel().attr(SyncLogByHashPusher.ATTR_KEY).set(syncLogByHashPusher);
// 开始数据同步
- syncLogByHashService.startPush();
+ syncLogByHashPusher.startPush();
}
/**
* 推送数据响应,推送下一批数据
*/
private void handleReqSyncByHashAck(ChannelHandlerContext ctx, SyncCmd syncCmd) {
- SyncLogByHashService syncLogByHashService = ctx.channel().attr(SyncLogByHashService.ATTR_KEY).get();
- syncLogByHashService.pushNextBatch();
+ SyncLogByHashPusher syncLogByHashPusher = ctx.channel().attr(SyncLogByHashPusher.ATTR_KEY).get();
+ syncLogByHashPusher.pushNextBatch();
+ }
+
+ private void handleSyncFinishClose(ChannelHandlerContext ctx, SyncCmd syncCmd) {
+ SyncLogByHashPusher syncLogByHashPusher = ctx.channel().attr(SyncLogByHashPusher.ATTR_KEY).get();
+ // 取消修改监听
+ DataChangeTrigger.instance().unsubscribe(syncLogByHashPusher);
+ // 移除已迁移的数据
+ syncLogByHashPusher.releaseResource();
+ // 断开连接
+ ctx.channel().close();
+ logger.info("channel@{} 同步结束,资源已释放", ctx.channel().id());
}
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java
index e90cd11..281bb98 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java
@@ -100,8 +100,10 @@ public class SyncLogClientHandler extends SimpleChannelInboundHandler {
method.invoke(data, params);
logger.info("invoke:{}, args:{}", method.getName(), params);
} catch (IllegalAccessException | InvocationTargetException e) {
- e.printStackTrace();
+ logger.error("failed invoke {}.{} args:{} ", dataType, method.getName(), params, e);
}
}
+
+
}
diff --git a/im-service/im-router/src/main/resources/application.yml b/im-service/im-router/src/main/resources/application.yml
index 2cd8da0..7812230 100644
--- a/im-service/im-router/src/main/resources/application.yml
+++ b/im-service/im-router/src/main/resources/application.yml
@@ -20,7 +20,8 @@ dubbo:
group: so-im
protocol:
name: dubbo
- port: 3032
+ port: 3034
+ # serialization: protobuf
consumer:
check: false
# filter: invoke_im_entry_filter # 调用im-entry时设置调用地址,配合im_entry_loadbalance路由到用户对应连接的im-entry
diff --git a/pom.xml b/pom.xml
index 2e9c276..e3038e5 100644
--- a/pom.xml
+++ b/pom.xml
@@ -43,6 +43,7 @@
5.1.0
2.0.3
0.9.1
+ 0.0.3