Browse Source

im-router数据同步方案

master
tangmingyou 4 years ago
parent
commit
cb3a6c9db9
  1. 27
      im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUser.java
  2. 10
      im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUserStorage.java
  3. 6
      im-service/im-router/src/main/java/net/sopod/soim/router/config/filter/InvokeImEntryFilter.java
  4. 7
      im-service/im-router/src/main/java/net/sopod/soim/router/config/loadbalance/ImEntryServerAddressLoadBalance.java
  5. 54
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
  6. 8
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncTypes.java
  7. 4
      im-service/im-router/src/main/java/net/sopod/soim/router/service/OnlineUserServiceImpl.java
  8. 6
      im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java
  9. 4
      im-service/im-router/src/main/java/net/sopod/soim/router/util/RpcContextUtil.java

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

@ -9,20 +9,15 @@ import net.sopod.soim.router.datasync.annotation.SyncIgnore;
import java.lang.reflect.Field; import java.lang.reflect.Field;
import java.lang.reflect.Modifier; import java.lang.reflect.Modifier;
class A {
private String aName;
}
/** /**
* RouterUser * RouterUser
* *
* @author tmy * @author tmy
* @date 2022-04-28 11:11 * @date 2022-04-28 11:11
*/ */
@EqualsAndHashCode(callSuper = true)
@Data @Data
@Accessors(chain = true) @Accessors(chain = true)
public class RouterUser extends A implements DataSync { public class RouterUser implements DataSync {
public static final int a = 1; public static final int a = 1;
@ -43,24 +38,4 @@ public class RouterUser extends A implements DataSync {
return this; return this;
} }
public static void main(String[] args) {
for (Field field : RouterUser.class.getDeclaredFields()) {
System.out.println(field.getName());
}
System.out.println("=============");
for (Field field : RouterUser.class.getDeclaredFields()) {
System.out.println(field.getName());
}
System.out.println("=============");
for (Field field : RouterUser.class.getSuperclass().getDeclaredFields()) {
System.out.println(field.getName());
}
System.out.println(RouterUser.class.getSuperclass().getSuperclass());
for (Field field : RouterUser.class.getSuperclass().getSuperclass().getDeclaredFields()) {
System.out.println(field.getName());
}
// Modifier.isFinal()
}
} }

10
im-service/im-router/src/main/java/net/sopod/soim/router/cache/SoImUserCache.java → im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUserStorage.java vendored

@ -14,19 +14,21 @@ import java.util.concurrent.ConcurrentHashMap;
/** /**
* ImUserCache * ImUserCache
* 在线用户属性缓存 * 在线用户属性缓存
* 分段[cache0,cache2,...cache31] id取模定位, lockMap0[okId1, okId2] 加读写锁避免序列化时加锁影响全部数据
* 数据序列化时get加锁删除加锁新增加锁代理类更新方法加锁
* *
* @author tmy * @author tmy
* @date 2022-05-02 14:07 * @date 2022-05-02 14:07
*/ */
public class SoImUserCache { public class RouterUserStorage {
private static final Logger logger = LoggerFactory.getLogger(SoImUserCache.class); private static final Logger logger = LoggerFactory.getLogger(RouterUserStorage.class);
private static final SoImUserCache INSTANCE = new SoImUserCache(); private static final RouterUserStorage INSTANCE = new RouterUserStorage();
private final ConcurrentHashMap<Long, RouterUser> routerUserMap = new ConcurrentHashMap<>(128); private final ConcurrentHashMap<Long, RouterUser> routerUserMap = new ConcurrentHashMap<>(128);
public static SoImUserCache getInstance() { public static RouterUserStorage getInstance() {
return INSTANCE; return INSTANCE;
} }

6
im-service/im-router/src/main/java/net/sopod/soim/router/config/filter/InvokeImEntryFilter.java

@ -2,7 +2,7 @@ package net.sopod.soim.router.config.filter;
import net.sopod.soim.common.constant.DubboConstant; import net.sopod.soim.common.constant.DubboConstant;
import net.sopod.soim.router.cache.RouterUser; import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.cache.SoImUserCache; import net.sopod.soim.router.cache.RouterUserStorage;
import org.apache.dubbo.rpc.*; import org.apache.dubbo.rpc.*;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
@ -30,8 +30,8 @@ public class InvokeImEntryFilter implements Filter {
// 获取请求链路用户id // 获取请求链路用户id
String ctxUid = invocation.getAttachment(DubboConstant.CTX_UID); String ctxUid = invocation.getAttachment(DubboConstant.CTX_UID);
if (ctxUid != null) { if (ctxUid != null) {
SoImUserCache soImUserCache = SoImUserCache.getInstance(); RouterUserStorage routerUserStorage = RouterUserStorage.getInstance();
RouterUser routerUser = soImUserCache.get(Long.valueOf(ctxUid)); RouterUser routerUser = routerUserStorage.get(Long.valueOf(ctxUid));
if (routerUser != null) { if (routerUser != null) {
imEntryServerAddr = routerUser.getImEntryAddr(); imEntryServerAddr = routerUser.getImEntryAddr();
} }

7
im-service/im-router/src/main/java/net/sopod/soim/router/config/loadbalance/ImEntryServerAddressLoadBalance.java

@ -1,10 +1,9 @@
package net.sopod.soim.router.config.loadbalance; package net.sopod.soim.router.config.loadbalance;
import net.sopod.soim.common.constant.DubboConstant;
import net.sopod.soim.common.util.StringUtil; import net.sopod.soim.common.util.StringUtil;
import net.sopod.soim.logic.common.util.RpcContextUtil; import net.sopod.soim.logic.common.util.RpcContextUtil;
import net.sopod.soim.router.cache.RouterUser; import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.cache.SoImUserCache; import net.sopod.soim.router.cache.RouterUserStorage;
import org.apache.dubbo.common.URL; import org.apache.dubbo.common.URL;
import org.apache.dubbo.rpc.Invocation; import org.apache.dubbo.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker; import org.apache.dubbo.rpc.Invoker;
@ -34,8 +33,8 @@ public class ImEntryServerAddressLoadBalance extends AbstractLoadBalance {
if (StringUtil.isEmpty(ctxUid)) { if (StringUtil.isEmpty(ctxUid)) {
throw new IllegalCallerException("调用im-entry服务接口,上下文uid未指定"); throw new IllegalCallerException("调用im-entry服务接口,上下文uid未指定");
} }
SoImUserCache soImUserCache = SoImUserCache.getInstance(); RouterUserStorage routerUserStorage = RouterUserStorage.getInstance();
RouterUser routerUser = soImUserCache.get(Long.valueOf(ctxUid)); RouterUser routerUser = routerUserStorage.get(Long.valueOf(ctxUid));
if (routerUser != null) { if (routerUser != null) {
entryAddr = routerUser.getImEntryAddr(); entryAddr = routerUser.getImEntryAddr();
} }

54
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java

@ -19,6 +19,9 @@ import java.util.concurrent.ConcurrentHashMap;
* BiSyncProxyManager * BiSyncProxyManager
* 设置属性的时候将设置方法同步处理 * 设置属性的时候将设置方法同步处理
* *
* 分段[cache0,cache2,...cache31] id取模定位, lockMap0[okId1, okId2] 加读写锁避免序列化时加锁影响全部数据
* 数据序列化时get加锁删除加锁新增加锁代理类更新方法加锁
*
* @author tmy * @author tmy
* @date 2022-05-04 17:10 * @date 2022-05-04 17:10
*/ */
@ -109,7 +112,6 @@ public class DataSyncProxyFactory {
if (!Modifier.isFinal(field.getModifiers())) { if (!Modifier.isFinal(field.getModifiers())) {
field.setAccessible(true); field.setAccessible(true);
field.set(target, field.get(source)); field.set(target, field.get(source));
System.out.println(field.getName() + ":" + field.get(source));
} }
} }
} }
@ -151,14 +153,48 @@ public class DataSyncProxyFactory {
} }
public static void main(String[] args) { public static void main(String[] args) {
RouterUser user1 = new RouterUser(); // RouterUser user1 = new RouterUser();
user1.setUid(10086L); // user1.setUid(10086L);
user1.setAccount("日月光"); // user1.setAccount("日月光");
//
RouterUser routerUser = newProxyInstance(SyncTypes.ROUTER_USER, user1); // RouterUser routerUser = newProxyInstance(SyncTypes.ROUTER_USER, user1);
System.out.println(routerUser); // System.out.println(routerUser);
routerUser.setOnlineTime(ImClock.millis()); // routerUser.setOnlineTime(ImClock.millis());
System.out.println(routerUser); // System.out.println(routerUser);
ConcurrentHashMap<String, RouterUser> map = new ConcurrentHashMap<>();
map.put("10081", new RouterUser().setAccount("阿基过天玺"));
map.put("10082", new RouterUser().setAccount("家国"));
map.put("10083", new RouterUser().setAccount("天下"));
Collection<RouterUser> values = map.values();
map.put("10084", new RouterUser().setAccount("晚风"));
new Thread(() -> {
for (int i = 0; i < 100; i++) {
map.put("100" + i, new RouterUser().setAccount("灯" + i));
if (i % 10 == 0) {
try {
Thread.sleep(100);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}).start();
int i = 0;
for (RouterUser value : values) {
System.out.println(value);
i++;
if (i % 3 == 0) {
try {
Thread.sleep(100);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
System.out.println(values.getClass());
System.out.println(values);
} }

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

@ -2,7 +2,7 @@ package net.sopod.soim.router.datasync;
import net.sopod.soim.common.util.StringUtil; import net.sopod.soim.common.util.StringUtil;
import net.sopod.soim.router.cache.RouterUser; import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.cache.SoImUserCache; import net.sopod.soim.router.cache.RouterUserStorage;
import javax.annotation.Nullable; import javax.annotation.Nullable;
import java.util.ArrayList; import java.util.ArrayList;
@ -39,18 +39,18 @@ public class SyncTypes {
@Override @Override
public RouterUser getData(String uid) { public RouterUser getData(String uid) {
return SoImUserCache.getInstance().get(Long.valueOf(uid)); return RouterUserStorage.getInstance().get(Long.valueOf(uid));
} }
@Override @Override
public boolean addData(RouterUser data) { public boolean addData(RouterUser data) {
SoImUserCache.getInstance().put(data.getUid(), data); RouterUserStorage.getInstance().put(data.getUid(), data);
return true; return true;
} }
@Override @Override
public boolean removeData(String uid) { public boolean removeData(String uid) {
return null != SoImUserCache.getInstance().remove(Long.valueOf(uid)); return null != RouterUserStorage.getInstance().remove(Long.valueOf(uid));
} }
}; };

4
im-service/im-router/src/main/java/net/sopod/soim/router/service/OnlineUserServiceImpl.java

@ -2,7 +2,7 @@ package net.sopod.soim.router.service;
import net.sopod.soim.entry.api.service.OnlineUserService; import net.sopod.soim.entry.api.service.OnlineUserService;
import net.sopod.soim.router.cache.RouterUser; import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.cache.SoImUserCache; import net.sopod.soim.router.cache.RouterUserStorage;
import org.apache.dubbo.config.annotation.DubboService; import org.apache.dubbo.config.annotation.DubboService;
/** /**
@ -16,7 +16,7 @@ public class OnlineUserServiceImpl implements OnlineUserService {
@Override @Override
public String getImEntryAddrByUid(Long uid) { public String getImEntryAddrByUid(Long uid) {
RouterUser routerUser = SoImUserCache.getInstance().get(uid); RouterUser routerUser = RouterUserStorage.getInstance().get(uid);
return routerUser == null ? null : routerUser.getImEntryAddr(); return routerUser == null ? null : routerUser.getImEntryAddr();
} }

6
im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java

@ -12,7 +12,7 @@ import net.sopod.soim.logic.common.model.UserInfo;
import net.sopod.soim.router.api.model.RegistryRes; import net.sopod.soim.router.api.model.RegistryRes;
import net.sopod.soim.router.cache.RouterUser; import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.api.service.UserRouteService; import net.sopod.soim.router.api.service.UserRouteService;
import net.sopod.soim.router.cache.SoImUserCache; import net.sopod.soim.router.cache.RouterUserStorage;
import net.sopod.soim.router.config.ImRouterAppContextHolder; import net.sopod.soim.router.config.ImRouterAppContextHolder;
import org.apache.dubbo.config.annotation.DubboReference; import org.apache.dubbo.config.annotation.DubboReference;
import org.apache.dubbo.config.annotation.DubboService; import org.apache.dubbo.config.annotation.DubboService;
@ -53,7 +53,7 @@ public class UserRouteServiceImpl implements UserRouteService {
.setIsOnline(Boolean.TRUE) .setIsOnline(Boolean.TRUE)
.setOnlineTime(ImClock.millis()) .setOnlineTime(ImClock.millis())
.setImEntryAddr(imEntryAddr); .setImEntryAddr(imEntryAddr);
SoImUserCache.getInstance().put(uid, routerUser); RouterUserStorage.getInstance().put(uid, routerUser);
// 接口返回 im_router_id,后续调用 im-router 负载均衡指向当前router服务 // 接口返回 im_router_id,后续调用 im-router 负载均衡指向当前router服务
return new RegistryRes() return new RegistryRes()
.setSuccess(true) .setSuccess(true)
@ -65,7 +65,7 @@ public class UserRouteServiceImpl implements UserRouteService {
//logger.info("client context uid: {}", ServerContext.getContextUid()); //logger.info("client context uid: {}", ServerContext.getContextUid());
logger.info("client context uid: {}", RpcContext.getServiceContext().getAttachment(DubboConstant.CTX_UID)); logger.info("client context uid: {}", RpcContext.getServiceContext().getAttachment(DubboConstant.CTX_UID));
Stream<RouterUser> stream = SoImUserCache.getInstance().getRouterUserMap().values().stream(); Stream<RouterUser> stream = RouterUserStorage.getInstance().getRouterUserMap().values().stream();
if (!StringUtil.isEmpty(keyword)) { if (!StringUtil.isEmpty(keyword)) {
// 根据关键词过滤 // 根据关键词过滤
stream = stream.filter(user -> user.getAccount().contains(keyword)); stream = stream.filter(user -> user.getAccount().contains(keyword));

4
im-service/im-router/src/main/java/net/sopod/soim/router/util/RpcContextUtil.java

@ -2,7 +2,7 @@ package net.sopod.soim.router.util;
import net.sopod.soim.common.constant.DubboConstant; import net.sopod.soim.common.constant.DubboConstant;
import net.sopod.soim.router.cache.RouterUser; import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.cache.SoImUserCache; import net.sopod.soim.router.cache.RouterUserStorage;
import org.apache.dubbo.rpc.RpcContext; import org.apache.dubbo.rpc.RpcContext;
import java.util.Objects; import java.util.Objects;
@ -18,7 +18,7 @@ public class RpcContextUtil {
public static boolean setImEntryRouteServerAddrByUid(Long uid) { public static boolean setImEntryRouteServerAddrByUid(Long uid) {
Objects.requireNonNull(uid, "设置im-entry路由参数uid不能为空"); Objects.requireNonNull(uid, "设置im-entry路由参数uid不能为空");
RpcContext.getServiceContext().setAttachment(DubboConstant.CTX_UID, uid); RpcContext.getServiceContext().setAttachment(DubboConstant.CTX_UID, uid);
RouterUser routerUser = SoImUserCache.getInstance().get(uid); RouterUser routerUser = RouterUserStorage.getInstance().get(uid);
String imEntryAddr; String imEntryAddr;
if (routerUser == null if (routerUser == null
|| null == (imEntryAddr = routerUser.getImEntryAddr())) { || null == (imEntryAddr = routerUser.getImEntryAddr())) {

Loading…
Cancel
Save