Browse Source

启动时数据同步

master
tangmingyou 4 years ago
parent
commit
57fae90f58
  1. 3
      README.md
  2. 8
      im-common/src/main/java/net/sopod/soim/common/util/Collects.java
  3. 27
      im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUserStorage.java
  4. 60
      im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java
  5. 65
      im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java
  6. 9
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java
  7. 47
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncStorage.java
  8. 7
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncTypes.java
  9. 3
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java
  10. 139
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java
  11. 13
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java
  12. 25
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java
  13. 10
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogServerHandler.java
  14. 2
      im-service/im-router/src/main/resources/application.yml

3
README.md

@ -10,6 +10,7 @@
- (05-12~05-12) 备用计划:router 层 id 号段负载均衡
- (05-13~05-14) dubbo 服务异步处理
- 功能开发:
- entry-http entry 节点获取功能(05-13~05-13)
- (05-13~05-13) 用户注册功能
- (05-14~05-14) 好友列表(在线状态:批量uid一致性hash, router查询)
- 用户查询
@ -32,7 +33,7 @@
- (05-25~05-25) 压测:压测开发
- (05-26~05-26) client: 控制台完善,grallvm 打包
- 后续:
- 通讯加密
- 通讯加密,服务链路SSL
- websocket 网关
- web 页面开发
- 异/同设备,多地登录

8
im-common/src/main/java/net/sopod/soim/common/util/Collects.java

@ -12,6 +12,14 @@ import java.util.function.Function;
*/
public class Collects {
public static boolean isEmpty(@Nullable Map<?, ?> map) {
return map == null || map.isEmpty();
}
public static boolean isNotEmpty(@Nullable Map<?, ?> map) {
return !isEmpty(map);
}
public static boolean isEmpty(@Nullable Collection<?> collection) {
return collection == null || collection.isEmpty();
}

27
im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUserStorage.java vendored

@ -2,12 +2,11 @@ package net.sopod.soim.router.cache;
import net.sf.cglib.proxy.Enhancer;
import net.sopod.soim.common.util.StringUtil;
import net.sopod.soim.router.datasync.DataChangeTrigger;
import net.sopod.soim.router.datasync.DataSyncProxyFactory;
import net.sopod.soim.router.datasync.SyncTypes;
import net.sopod.soim.router.datasync.*;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Iterator;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@ -20,7 +19,7 @@ import java.util.concurrent.ConcurrentHashMap;
* @author tmy
* @date 2022-05-02 14:07
*/
public class RouterUserStorage {
public class RouterUserStorage extends DataSyncStorage<RouterUser> {
private static final Logger logger = LoggerFactory.getLogger(RouterUserStorage.class);
@ -32,17 +31,22 @@ public class RouterUserStorage {
return INSTANCE;
}
public RouterUserStorage() {
super.registry(SyncTypes.ROUTER_USER, this);
}
public RouterUser put(Long uid, RouterUser routerUser) {
// TODO 这里克隆一个代理对象
if (Enhancer.isEnhanced(routerUser.getClass())) {
routerUserMap.put(uid, routerUser);
return routerUser;
}
RouterUser proxyRouterUser = DataSyncProxyFactory.newProxyInstance(SyncTypes.ROUTER_USER);
// 创建代理对象
RouterUser proxyRouterUser = DataSyncProxyFactory.newProxyInstance(SyncTypes.ROUTER_USER, routerUser);
routerUserMap.put(uid, proxyRouterUser);
// 新增数据触发
DataChangeTrigger.instance().onAdd(SyncTypes.ROUTER_USER, routerUser);
return routerUserMap.put(uid, routerUser);
super.onDataAdd(proxyRouterUser);
return proxyRouterUser;
}
public RouterUser get(Long uid) {
@ -51,7 +55,7 @@ public class RouterUserStorage {
public RouterUser remove(Long uid) {
if (uid != null) {
DataChangeTrigger.instance().onRemove(SyncTypes.ROUTER_USER, StringUtil.toString(uid));
super.onDataRemove(StringUtil.toString(uid));
return routerUserMap.remove(uid);
}
return null;
@ -61,4 +65,9 @@ public class RouterUserStorage {
return routerUserMap;
}
@Override
public Iterator<RouterUser> getFullDataIterator() {
return routerUserMap.values().iterator();
}
}

60
im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java

@ -1,12 +1,20 @@
package net.sopod.soim.router.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.pojo.Instance;
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.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.ApplicationContext;
import java.util.Collections;
import java.util.List;
import java.util.Properties;
import java.util.concurrent.CopyOnWriteArrayList;
/**
@ -20,6 +28,8 @@ import java.util.concurrent.CopyOnWriteArrayList;
*/
public class AppContextHolder {
private static final Logger logger = LoggerFactory.getLogger(AppContextHolder.class);
private static ApplicationContext applicationContext;
/**
@ -27,6 +37,7 @@ public class AppContextHolder {
*/
private static final List<URL> registryInvokerUrls = new CopyOnWriteArrayList<>();
private static String discoveryAddr;
private static String appAddr;
private static String appHost;
private static int appPort;
@ -59,6 +70,14 @@ public class AppContextHolder {
AppContextHolder.appPort = port;
}
public static void setServiceDiscoveryRegistryAddr(String discoveryAddr) {
AppContextHolder.discoveryAddr = discoveryAddr;
}
public static String getDiscoveryAddr() {
return discoveryAddr;
}
public static String getAppAddr() {
return appAddr;
}
@ -71,4 +90,45 @@ public class AppContextHolder {
return appPort;
}
private static NamingService namingService;
/**
* 获取 im-router 集群节点信息
* TODO 注册应用关闭
*/
public static List<Instance> getClusterImRouterInstance() {
String discoveryAddr;
if (null == (discoveryAddr = AppContextHolder.getDiscoveryAddr())) {
return Collections.emptyList();
}
if (namingService == null) {
synchronized (AppContextHolder.class) {
if (namingService == null) {
Properties properties = new Properties();
properties.put("serverAddr", discoveryAddr);
// 获取当前服务实例
try {
namingService = NacosFactory.createNamingService(properties);
} catch (NacosException e) {
throw new IllegalStateException("nacos实例"+discoveryAddr+"连接失败", e);
}
}
}
}
try {
List<Instance> allInstances = namingService.getAllInstances(AppConstant.APP_IM_ROUTER_NAME);
return allInstances;
} catch (NacosException e) {
throw new IllegalStateException(AppConstant.APP_IM_ROUTER_NAME + "集群信息获取失败:", e);
} finally {
if (namingService != null) {
try {
namingService.shutDown();
} catch (NacosException e) {
logger.error("查询{}服务NamingServer关闭失败!", AppConstant.APP_IM_ROUTER_NAME, e);
}
}
}
}
}

65
im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java

@ -8,6 +8,8 @@ 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;
import net.sopod.soim.router.cache.RouterUserStorage;
import net.sopod.soim.router.datasync.SyncLogMigrateService;
import net.sopod.soim.router.datasync.server.session.SyncServerSession;
import org.apache.commons.lang3.tuple.ImmutablePair;
@ -47,6 +49,7 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
@Override
public void onApplicationEvent(ApplicationReadyEvent event) {
AppContextHolder.setApplicationContext(event.getApplicationContext());
this.init();
// 启动数据同步服务
this.startSyncServer();
@ -57,10 +60,46 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
// 服务已可用进行注册
if (registryNow) {
this.createMockData();
this.doRegistry();
}
}
private void createMockData() {
for (long i = 10000L; i < 11000L; i++) {
RouterUser routerUser = new RouterUser()
.setUid(i)
.setAccount("京城" + i)
.setIsOnline(i % 2 == 0)
.setImEntryAddr("127.0.0.1:1313");
RouterUserStorage.getInstance().put(i, routerUser);
}
}
/**
* 获取注册中心地址
*/
private void init() {
RegistryManager registryManager = ApplicationModel.defaultModel().getBeanFactory()
.getBean(RegistryManager.class);
Collection<Registry> registries = registryManager.getRegistries();
if (Collects.isNotEmpty(registries)) {
for (Registry registry : registries) {
// isServiceDiscovery(): true是注册应用(im-router)的registry, false是注册服务接口的registry
if (registry.isAvailable()
&& registry.isServiceDiscovery()) {
String discoveryAddr = registry.getUrl().getAddress();
AppContextHolder.setServiceDiscoveryRegistryAddr(discoveryAddr);
logger.info("service discovery registry addr: {}", discoveryAddr);
break;
}
}
}
if (AppContextHolder.getDiscoveryAddr() == null) {
logger.info("未解析到服务注册中心: 将直接启动不进行数据迁移");
}
}
private void startSyncServer() {
SyncServerSession.getInstance()
.start(AppContextHolder.getAppPort() + SYNC_SERVER_PORT_OFFSET);
@ -101,7 +140,7 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
public boolean checkClusterEnvironment(SyncLogMigrateService syncLogMigrateService) {
// TODO 加分布式锁,同一时间单一节点进行数据同步
List<Instance> clusterImEntryInstance = getClusterImEntryInstance();
List<Instance> clusterImEntryInstance = AppContextHolder.getClusterImRouterInstance();
if (Collects.isEmpty(clusterImEntryInstance)) {
logger.info("当前集群无{}节点,直接启动", AppConstant.APP_IM_ROUTER_NAME);
return true;
@ -134,28 +173,6 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
return false;
}
private List<Instance> getClusterImEntryInstance() {
String serverAddr = "124.222.131.236:3848";
Properties properties = new Properties();
properties.put("serverAddr", serverAddr);
// 获取当前服务实例
NamingService namingService = null;
try {
namingService = NacosFactory.createNamingService(properties);
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);
}
}
}
}
/**
* 一致性hash使用了虚拟节点会导致迁移多个数据节点
*
@ -169,6 +186,8 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
* 4.数据迁移完成后注册dubbo服务
* 5.定时n秒后或添加provider filter服务第一次被调用m秒后依次调用所有迁移数据节点可清空用户数据用户数据已移过来了失败不用管
* 可能依次调用清空数据过程中又会有数据更新需要订阅
*
* 关闭时检查默认有备份注册备份无备份推送尝试数据到其他节点
*/
private void nameServerTestCode() {
String serverAddr = "124.222.131.236:3848";

9
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java

@ -18,6 +18,8 @@ public class DataChangeTrigger {
private static DataChangeTrigger INSTANCE;
// TODO subscribes 订阅 selector.select(dataKey);
public static DataChangeTrigger instance() {
if (INSTANCE == null) {
synchronized (DataChangeTrigger.class) {
@ -74,4 +76,11 @@ public class DataChangeTrigger {
return seqCounterMap.computeIfAbsent(dataKey, key -> new AtomicInteger());
}
/**
* 删除暂存数据字段等...
*/
public void clean() {
}
}

47
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncStorage.java

@ -0,0 +1,47 @@
package net.sopod.soim.router.datasync;
import net.sopod.soim.common.util.StringUtil;
import net.sopod.soim.router.cache.RouterUser;
import java.util.HashMap;
import java.util.Iterator;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* DataSyncStorage
*
* @author tmy
* @date 2022-05-10 10:23
*/
public abstract class DataSyncStorage<T extends DataSync> {
private static ConcurrentHashMap<DataSyncStorage<? extends DataSync>, SyncTypes.SyncType<? extends DataSync>> STORAGES = new ConcurrentHashMap<>();
public static Map<DataSyncStorage<? extends DataSync>, SyncTypes.SyncType<? extends DataSync>> getStorages() {
return new HashMap<>(STORAGES);
}
protected SyncTypes.SyncType<T> syncType;
/**
* 注册可数据数据容器
*/
protected void registry(SyncTypes.SyncType<T> syncType, DataSyncStorage<T> storage) {
STORAGES.put(storage, syncType);
this.syncType = syncType;
}
public abstract Iterator<T> getFullDataIterator();
public void onDataAdd(T data) {
// TODO...
DataChangeTrigger.instance().onAdd(syncType, data);
}
public void onDataRemove(String dataKey) {
// TODO...
DataChangeTrigger.instance().onRemove(SyncTypes.ROUTER_USER, dataKey);
}
}

7
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncTypes.java

@ -52,6 +52,11 @@ public class SyncTypes {
public boolean removeData(String uid) {
return null != RouterUserStorage.getInstance().remove(Long.valueOf(uid));
}
@Override
public int onceSyncSize() {
return 30;
}
};
public static abstract class SyncType<T extends DataSync> {
@ -92,6 +97,8 @@ public class SyncTypes {
public abstract boolean removeData(String key);
public abstract int onceSyncSize();
@Nullable
@SuppressWarnings("unchecked")
static <T extends DataSync> SyncType<T> getSyncType(int ordinal) {

3
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java

@ -6,6 +6,8 @@ import io.netty.channel.ChannelInitializer;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
import io.netty.util.AttributeKey;
import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.datasync.SyncTypes;
@ -37,6 +39,7 @@ public class SyncClient {
@Override
protected void initChannel(SocketChannel channel) throws Exception {
channel.pipeline()
.addLast(new LoggingHandler(LogLevel.INFO))
.addLast(new SyncCmdCodec())
.addLast(new SyncLogEncoder())
.addLast(new SyncCmdClientHandler())

139
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java

@ -1,6 +1,21 @@
package net.sopod.soim.router.datasync.server;
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.datasync.DataSync;
import net.sopod.soim.router.datasync.DataSyncStorage;
import net.sopod.soim.router.datasync.SyncTypes;
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;
/**
* SyncLogPushService
@ -11,19 +26,139 @@ import io.netty.util.AttributeKey;
public class SyncLogByHashService {
private static final Logger logger = LoggerFactory.getLogger(SyncLogByHashService.class);
public static final AttributeKey<SyncLogByHashService> ATTR_KEY = AttributeKey
.valueOf(SyncLogByHashService.class, "SYNC_LOG_BY_HASH_SERVICE");
private final WeakReference<Channel> clientChannel;
private final String newNodeAddr;
public SyncLogByHashService(String newNodeAddr) {
private final UidConsistentHashSelector<String> selector;
Set<String> syncedDataKeys = new HashSet<>();
public SyncLogByHashService(Channel clientChannel, String newNodeAddr) {
this.clientChannel = new WeakReference<>(clientChannel);
this.newNodeAddr = newNodeAddr;
// 构建hash环匹配要迁移的数据
Map<String, String> twoNodes = new HashMap<>();
twoNodes.put(newNodeAddr, newNodeAddr);
twoNodes.put(AppContextHolder.getAppAddr(), AppContextHolder.getAppAddr());
selector = new UidConsistentHashSelector<>(twoNodes, twoNodes.hashCode());
}
public static void main(String[] args) {
Map<String, Pair<String, AtomicInteger>> 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<Pair<String, AtomicInteger>> selector = new UidConsistentHashSelector<>(twoNodes, twoNodes.hashCode());
for (int i = 10000; i < 11000; i++) {
Pair<String, AtomicInteger> pair = selector.select(i + "");
pair.getRight().incrementAndGet();
}
System.out.println(twoNodes);
}
private Map<DataSyncStorage<? extends DataSync>, SyncTypes.SyncType<? extends DataSync>> storages;
private DataSyncStorage<? extends DataSync> curStorage;
private Iterator<? extends DataSync> fullDataIterator;
private int totalCount = 0;
public void startPush() {
// TODO 开始数据推送
// 开始数据推送
this.storages = DataSyncStorage.getStorages();
if (Collects.isEmpty(this.storages)) {
this.syncFinishSuccess();
return;
}
this.curStorage = storages.keySet().iterator().next();
this.fullDataIterator = this.curStorage.getFullDataIterator();
this.pushNextBatch();
}
@SuppressWarnings("unchecked")
public void pushNextBatch() {
if (this.fullDataIterator == null) {
return;
}
int count = 0;
SyncTypes.SyncType<DataSync> syncType = (SyncTypes.SyncType<DataSync>)storages.get(curStorage);
SyncLog.AddLog<DataSync> addLog = SyncLog.addLog(0, syncType);
while(fullDataIterator.hasNext()) {
DataSync data = fullDataIterator.next();
String dataKey = syncType.getDataKey(data);
// 保留数据不迁移
if (!newNodeAddr.equals(selector.select(dataKey))) {
continue;
}
logger.info("sync data: {}", dataKey);
// 迁移数据
// TODO 数据加锁
addLog.addData(data);
// TODO DataChangeTrigger.instance().subscribe(syncType, dataKey)
// 注册更改日志
syncedDataKeys.add(dataKey);
// TODO 解锁
count ++; totalCount ++;
// 批量发送数据
if (count >= syncType.onceSyncSize()) {
Channel channel = clientChannel.get();
if (channel == null || !channel.isActive()) {
this.syncFinishFail();
return;
}
channel.writeAndFlush(addLog);
// 重新创建对象
addLog = null;
break;
}
}
// 最后一点数据
if (addLog != null) {
Channel channel = clientChannel.get();
// 发送数据
if (channel == null || !channel.isActive()) {
this.syncFinishFail();
return;
}
channel.writeAndFlush(addLog);
}
this.checkNextIterator();
}
/**
* 获取下一个 storage
*/
private void checkNextIterator() {
if (this.fullDataIterator != null && this.fullDataIterator.hasNext()) {
return;
}
if (this.curStorage != null) {
storages.remove(this.curStorage);
}
if (Collects.isEmpty(storages)) {
this.curStorage = null;
this.fullDataIterator = null;
this.syncFinishSuccess();
return;
}
this.curStorage = storages.keySet().iterator().next();
this.fullDataIterator = this.curStorage.getFullDataIterator();
}
private void syncFinishFail() {
logger.error("sync push finish fail!");
}
private void syncFinishSuccess() {
logger.info("sync push finish and success!");
}
}

13
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java

@ -30,6 +30,9 @@ public class SyncCmdServerHandler extends SimpleChannelInboundHandler<SyncCmd> {
case SyncCmd.SYNC_BY_HASH:
this.handleReqSyncByHash(ctx, syncCmd);
break;
case SyncCmd.SYNC_BY_HASH_ACK:
this.handleReqSyncByHashAck(ctx, syncCmd);
break;
}
}
@ -47,10 +50,18 @@ public class SyncCmdServerHandler extends SimpleChannelInboundHandler<SyncCmd> {
private void handleReqSyncByHash(ChannelHandlerContext ctx, SyncCmd syncCmd) {
String clientAddr = syncCmd.getParam1();
// 绑定数据同步服务
SyncLogByHashService syncLogByHashService = new SyncLogByHashService(clientAddr);
SyncLogByHashService syncLogByHashService = new SyncLogByHashService(ctx.channel(), clientAddr);
ctx.channel().attr(SyncLogByHashService.ATTR_KEY).set(syncLogByHashService);
// 开始数据同步
syncLogByHashService.startPush();
}
/**
* 推送数据响应推送下一批数据
*/
private void handleReqSyncByHashAck(ChannelHandlerContext ctx, SyncCmd syncCmd) {
SyncLogByHashService syncLogByHashService = ctx.channel().attr(SyncLogByHashService.ATTR_KEY).get();
syncLogByHashService.pushNextBatch();
}
}

25
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java

@ -2,7 +2,16 @@ package net.sopod.soim.router.datasync.server.handler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import net.sopod.soim.router.datasync.DataSync;
import net.sopod.soim.router.datasync.DataSyncStorage;
import net.sopod.soim.router.datasync.SyncTypes;
import net.sopod.soim.router.datasync.server.codec.CodecUtil;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import net.sopod.soim.router.datasync.server.data.SyncLog;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
/**
* SyncLogHandler
@ -13,9 +22,21 @@ import net.sopod.soim.router.datasync.server.data.SyncLog;
*/
public class SyncLogClientHandler extends SimpleChannelInboundHandler<SyncLog> {
@Override
protected void channelRead0(ChannelHandlerContext channelHandlerContext, SyncLog syncLog) throws Exception {
private static final Logger logger = LoggerFactory.getLogger(SyncLogClientHandler.class);
@Override
protected void channelRead0(ChannelHandlerContext ctx, SyncLog syncLog) {
List<byte[]> bytesList = syncLog.getSerializeDataCollect();
SyncTypes.SyncType<DataSync> syncType = SyncTypes.getSyncType(syncLog.getSyncType());
for (byte[] bytes : bytesList) {
DataSync instance = CodecUtil.decode(bytes, syncType.dataType());
syncType.addData(instance);
System.out.println(instance);
}
// 同步完成响应
SyncCmd syncCmd = new SyncCmd().setCmdType(SyncCmd.SYNC_BY_HASH_ACK);
ctx.writeAndFlush(syncCmd);
logger.info("storage: {}", DataSyncStorage.getStorages());
}
}

10
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogServerHandler.java

@ -5,6 +5,7 @@ import io.netty.channel.SimpleChannelInboundHandler;
import net.sopod.soim.router.datasync.DataSync;
import net.sopod.soim.router.datasync.SyncTypes;
import net.sopod.soim.router.datasync.server.codec.CodecUtil;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import net.sopod.soim.router.datasync.server.data.SyncLog;
import java.util.List;
@ -19,14 +20,17 @@ import java.util.List;
public class SyncLogServerHandler extends SimpleChannelInboundHandler<SyncLog> {
@Override
protected void channelRead0(ChannelHandlerContext channelHandlerContext, SyncLog syncLog) throws Exception {
protected void channelRead0(ChannelHandlerContext ctx, SyncLog syncLog) throws Exception {
List<byte[]> bytesList = syncLog.getSerializeDataCollect();
SyncTypes.SyncType<DataSync> syncType = SyncTypes.getSyncType(syncLog.getSyncType());
for (byte[] bytes : bytesList) {
SyncTypes.SyncType<DataSync> syncType = SyncTypes.getSyncType(syncLog.getSyncType());
DataSync instance = CodecUtil.decode(bytes, syncType.dataType());
syncType.addData(instance);
System.out.println(instance);
}
System.out.println("read: " + syncLog);
// 同步完成响应
SyncCmd syncCmd = new SyncCmd().setCmdType(SyncCmd.SYNC_BY_HASH_ACK);
ctx.writeAndFlush(syncCmd);
}
}

2
im-service/im-router/src/main/resources/application.yml

@ -20,7 +20,7 @@ dubbo:
group: so-im
protocol:
name: dubbo
port: 3031
port: 3032
consumer:
check: false
# filter: invoke_im_entry_filter # 调用im-entry时设置调用地址,配合im_entry_loadbalance路由到用户对应连接的im-entry

Loading…
Cancel
Save