diff --git a/README.md b/README.md index 904a367..8d5ac8d 100644 --- a/README.md +++ b/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 页面开发 - 异/同设备,多地登录 diff --git a/im-common/src/main/java/net/sopod/soim/common/util/Collects.java b/im-common/src/main/java/net/sopod/soim/common/util/Collects.java index 8b1c9c8..9409d33 100644 --- a/im-common/src/main/java/net/sopod/soim/common/util/Collects.java +++ b/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(); } diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUserStorage.java b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUserStorage.java index 26866fb..229c700 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUserStorage.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUserStorage.java @@ -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 { 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 getFullDataIterator() { + return routerUserMap.values().iterator(); + } + } 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/AppContextHolder.java index 0524950..8853c4d 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/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 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 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 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); + } + } + } + } + } 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 59741c0..70eddaf 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java @@ -8,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 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 clusterImEntryInstance = getClusterImEntryInstance(); + List 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 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 new AtomicInteger()); } + /** + * 删除暂存数据字段等... + */ + public void clean() { + + } + } diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncStorage.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncStorage.java new file mode 100644 index 0000000..7ec0669 --- /dev/null +++ b/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 { + + private static ConcurrentHashMap, SyncTypes.SyncType> STORAGES = new ConcurrentHashMap<>(); + + public static Map, SyncTypes.SyncType> getStorages() { + return new HashMap<>(STORAGES); + } + + protected SyncTypes.SyncType syncType; + + /** + * 注册可数据数据容器 + */ + protected void registry(SyncTypes.SyncType syncType, DataSyncStorage storage) { + STORAGES.put(storage, syncType); + this.syncType = syncType; + } + + public abstract Iterator 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); + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncTypes.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncTypes.java index 5ef33ed..6b8d547 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncTypes.java +++ b/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 { @@ -92,6 +97,8 @@ public class SyncTypes { public abstract boolean removeData(String key); + public abstract int onceSyncSize(); + @Nullable @SuppressWarnings("unchecked") static SyncType getSyncType(int ordinal) { 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 9285706..a02b248 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 @@ -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()) diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java index 9725bc1..b6783d4 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogByHashService.java @@ -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 ATTR_KEY = AttributeKey .valueOf(SyncLogByHashService.class, "SYNC_LOG_BY_HASH_SERVICE"); + private final WeakReference clientChannel; + private final String newNodeAddr; - public SyncLogByHashService(String newNodeAddr) { + private final UidConsistentHashSelector selector; + + Set syncedDataKeys = new HashSet<>(); + + public SyncLogByHashService(Channel clientChannel, String newNodeAddr) { + this.clientChannel = new WeakReference<>(clientChannel); this.newNodeAddr = newNodeAddr; + // 构建hash环匹配要迁移的数据 + Map 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> 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); } + private Map, SyncTypes.SyncType> storages; + + private DataSyncStorage curStorage; + + private Iterator 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 syncType = (SyncTypes.SyncType)storages.get(curStorage); + SyncLog.AddLog 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!"); } } 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 aaeea0f..8a5a792 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 @@ -30,6 +30,9 @@ public class SyncCmdServerHandler extends SimpleChannelInboundHandler { 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 { 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(); + } + } 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 ec92d7e..d7a4a73 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 @@ -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 { - @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 bytesList = syncLog.getSerializeDataCollect(); + SyncTypes.SyncType 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()); } } diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogServerHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogServerHandler.java index 808edf3..0f20e13 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogServerHandler.java +++ b/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 { @Override - protected void channelRead0(ChannelHandlerContext channelHandlerContext, SyncLog syncLog) throws Exception { + protected void channelRead0(ChannelHandlerContext ctx, SyncLog syncLog) throws Exception { List bytesList = syncLog.getSerializeDataCollect(); + SyncTypes.SyncType syncType = SyncTypes.getSyncType(syncLog.getSyncType()); for (byte[] bytes : bytesList) { - SyncTypes.SyncType 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); } } diff --git a/im-service/im-router/src/main/resources/application.yml b/im-service/im-router/src/main/resources/application.yml index c0e5565..2cd8da0 100644 --- a/im-service/im-router/src/main/resources/application.yml +++ b/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