Browse Source

im-router数据同步开发

master
tangmingyou 4 years ago
parent
commit
72c76a14bc
  1. 2
      im-client/src/main/java/net/sopod/soim/client/config/ClientConfig.java
  2. 6
      im-common/pom.xml
  3. 2
      im-common/src/main/java/net/sopod/soim/common/util/Collects.java
  4. 2
      im-common/src/main/java/net/sopod/soim/common/util/netty/FastThreadLocalThreadFactory.java
  5. 3
      im-entry/src/main/java/net/sopod/soim/entry/server/EntryServer.java
  6. 3
      im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java
  7. 2
      im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java
  8. 40
      im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUser.java
  9. 20
      im-service/im-router/src/main/java/net/sopod/soim/router/cache/SoImUserCache.java
  10. 45
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java
  11. 104
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
  12. 10
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncDataIncr.java
  13. 17
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLog.java
  14. 19
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogInboundHandler.java
  15. 42
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java
  16. 33
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServerInboundHandler.java

2
im-client/src/main/java/net/sopod/soim/client/config/ClientConfig.java

@ -17,6 +17,6 @@ public class ClientConfig {
private String host = "127.0.0.1";
private Integer port = 8087;
private Integer port = 8089;
}

6
im-common/pom.xml

@ -26,6 +26,12 @@
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
</dependency>
<dependency>
<groupId>io.netty</groupId>
<artifactId>netty-common</artifactId>
<version>${netty.version}</version>
<scope>provided</scope>
</dependency>
</dependencies>
</project>

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

@ -21,7 +21,7 @@ public class Collects {
}
public static boolean isEmpty(@Nullable Object[] arr) {
return arr == null || arr.length > 0;
return arr == null || arr.length == 0;
}
public static boolean isNotEmpty(@Nullable Object[] arr) {

2
im-entry/src/main/java/net/sopod/soim/entry/util/FastThreadLocalThreadFactory.java → im-common/src/main/java/net/sopod/soim/common/util/netty/FastThreadLocalThreadFactory.java

@ -1,4 +1,4 @@
package net.sopod.soim.entry.util;
package net.sopod.soim.common.util.netty;
import io.netty.util.concurrent.FastThreadLocalThread;

3
im-entry/src/main/java/net/sopod/soim/entry/server/EntryServer.java

@ -7,9 +7,8 @@ import io.netty.channel.ChannelOption;
import io.netty.channel.WriteBufferWaterMark;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.nio.NioServerSocketChannel;
import io.netty.util.concurrent.DefaultThreadFactory;
import net.sopod.soim.common.constant.Consts;
import net.sopod.soim.entry.util.FastThreadLocalThreadFactory;
import net.sopod.soim.common.util.netty.FastThreadLocalThreadFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

3
im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java

@ -20,7 +20,8 @@ public class ImEntryInitializer extends ChannelInitializer<SocketChannel> {
private static final Logger logger = LoggerFactory.getLogger(ImEntryInitializer.class);
@Override
protected void initChannel(SocketChannel socketChannel) throws Exception {
protected void initChannel(SocketChannel socketChannel) {
logger.info("init channel: {}", socketChannel);
LogLevel logLevel = logger.isDebugEnabled() ? LogLevel.DEBUG : LogLevel.INFO;
ChannelPipeline pipeline = socketChannel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel))

2
im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java

@ -4,7 +4,7 @@ import com.lmax.disruptor.*;
import com.lmax.disruptor.dsl.Disruptor;
import com.lmax.disruptor.dsl.ProducerType;
import net.sopod.soim.common.util.ImClock;
import net.sopod.soim.entry.util.FastThreadLocalThreadFactory;
import net.sopod.soim.common.util.netty.FastThreadLocalThreadFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

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

@ -1,8 +1,17 @@
package net.sopod.soim.router.cache;
import lombok.Data;
import lombok.EqualsAndHashCode;
import lombok.experimental.Accessors;
import net.sopod.soim.router.datasync.DataSync;
import net.sopod.soim.router.datasync.annotation.SyncIgnore;
import java.lang.reflect.Field;
import java.lang.reflect.Modifier;
class A {
private String aName;
}
/**
* RouterUser
@ -10,9 +19,12 @@ import net.sopod.soim.router.datasync.DataSync;
* @author tmy
* @date 2022-04-28 11:11
*/
@EqualsAndHashCode(callSuper = true)
@Data
@Accessors(chain = true)
public class RouterUser implements DataSync {
public class RouterUser extends A implements DataSync {
public static final int a = 1;
private long uid;
@ -25,4 +37,30 @@ public class RouterUser implements DataSync {
private String imEntryAddr;
@SyncIgnore
public RouterUser setUid(long uid) {
this.uid = uid;
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()
}
}

20
im-service/im-router/src/main/java/net/sopod/soim/router/cache/SoImUserCache.java vendored

@ -1,5 +1,10 @@
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 org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -26,6 +31,15 @@ public class SoImUserCache {
}
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);
// 新增数据触发
DataChangeTrigger.instance().onAdd(SyncTypes.ROUTER_USER, routerUser);
return routerUserMap.put(uid, routerUser);
}
@ -34,7 +48,11 @@ public class SoImUserCache {
}
public RouterUser remove(Long uid) {
return routerUserMap.remove(uid);
if (uid != null) {
DataChangeTrigger.instance().onRemove(SyncTypes.ROUTER_USER, StringUtil.toString(uid));
return routerUserMap.remove(uid);
}
return null;
}
public Map<Long, RouterUser> getRouterUserMap() {

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

@ -3,6 +3,7 @@ package net.sopod.soim.router.datasync;
import net.sopod.soim.router.datasync.server.SyncLog;
import java.util.Queue;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicInteger;
@ -15,28 +16,62 @@ import java.util.concurrent.atomic.AtomicInteger;
*/
public class DataChangeTrigger {
private static final DataChangeTrigger INSTANCE = new DataChangeTrigger();
private static DataChangeTrigger INSTANCE;
public static DataChangeTrigger instance() {
if (INSTANCE == null) {
synchronized (DataChangeTrigger.class) {
if (INSTANCE == null) {
INSTANCE = new DataChangeTrigger();
}
}
}
return INSTANCE;
}
private AtomicInteger seqCounter = new AtomicInteger();
private final ConcurrentHashMap<String, AtomicInteger> seqCounterMap = new ConcurrentHashMap<>(128);
private Queue<SyncLog> logQueue = new ConcurrentLinkedQueue<>();
private final Queue<SyncLog> logQueue = new ConcurrentLinkedQueue<>();
public <T extends DataSync> void onUpdate(SyncTypes.SyncType<T> syncType, String dataKey, String method, Object[] args) {
// 序列化 args,避免后续更改
// logQueue.add()
SyncLog.UpdateLog<T> updateLog = SyncLog.updateLog(getSeq(dataKey), syncType)
.setDataKey(dataKey)
.setMethod(method)
.setArgs(args);
publishLog(updateLog);
}
public <T extends DataSync> void onAdd(SyncTypes.SyncType<T> syncType, T data) {
String dataKey = syncType.getDataKey(data);
AtomicInteger seqCounter = getSeqCounter(dataKey);
// 序列化 data,避免后续更改
SyncLog.AddLog<T> addLog = SyncLog.addLog(seqCounter.getAndIncrement(), syncType)
.addData(data);
publishLog(addLog);
}
/**
* 数据删除日志
*/
public <T extends DataSync> void onRemove(SyncTypes.SyncType<T> syncType, String dataKey) {
SyncLog.RemoveLog<T> removeLog = SyncLog.removeLog(getSeq(dataKey), syncType)
.setDataKey(dataKey);
publishLog(removeLog);
}
private void publishLog(SyncLog log) {
// TODO 判断无订阅者跳过,新增节点同步按 dataKey 单独订阅每一个数据更新
logQueue.add(log);
System.out.println("publish log: "+log);
}
private int getSeq(String dataKey) {
return getSeqCounter(dataKey).getAndIncrement();
}
private AtomicInteger getSeqCounter(String dataKey) {
return seqCounterMap.computeIfAbsent(dataKey, key -> new AtomicInteger());
}
}

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

@ -3,15 +3,15 @@ package net.sopod.soim.router.datasync;
import net.sf.cglib.proxy.Enhancer;
import net.sf.cglib.proxy.MethodInterceptor;
import net.sf.cglib.proxy.MethodProxy;
import net.sopod.soim.common.util.Collects;
import net.sopod.soim.common.util.ImClock;
import net.sopod.soim.common.util.Jackson;
import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.datasync.annotation.SyncIgnore;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.lang.reflect.Constructor;
import java.lang.reflect.InvocationTargetException;
import java.lang.reflect.Method;
import java.lang.reflect.*;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
@ -31,8 +31,12 @@ public class DataSyncProxyFactory {
*/
private static final Map<Class<? extends DataSync>, Set<String>> typeUpdaterMethodsCache = new ConcurrentHashMap<>();
@SuppressWarnings("unchecked")
public static <T extends DataSync> T newProxyInstance(SyncTypes.SyncType<T> syncType) {
return newProxyInstance(syncType, null);
}
@SuppressWarnings("unchecked")
public static <T extends DataSync> T newProxyInstance(SyncTypes.SyncType<T> syncType, T source) {
Class<T> type = syncType.dataType();
T instance;
try {
@ -78,7 +82,41 @@ public class DataSyncProxyFactory {
Enhancer enhancer = new Enhancer();
enhancer.setSuperclass(type);
enhancer.setCallback(new DataSyncProxyCallback<>(syncType, updaterMethods));
return (T) enhancer.create();
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 <T> void cloneFields(T source, T proxyObj) throws Exception {
cloneFields0(source.getClass(), source, proxyObj);
}
/**
* 复制对象属性值
* TODO 缓存 <class,fields> 避免每次clone从方法区查找
*/
private static <T> void cloneFields0(Class<?> clazz, T source, T target) throws Exception {
Field[] fields = clazz.getDeclaredFields();
if (Collects.isNotEmpty(fields)) {
for (Field field : fields) {
if (!Modifier.isFinal(field.getModifiers())) {
field.setAccessible(true);
field.set(target, field.get(source));
System.out.println(field.getName() + ":" + field.get(source));
}
}
}
Class<?> superClazz = clazz.getSuperclass();
if (superClazz != Object.class) {
cloneFields0(superClazz, source, target);
}
}
public static class DataSyncProxyCallback<T extends DataSync> implements MethodInterceptor {
@ -103,7 +141,7 @@ public class DataSyncProxyFactory {
return methodProxy.invokeSuper(instance, args);
}
// TODO 记录更新操作(方法和参数),查询数据订阅者,异步同步数据
// 记录数据更新操作(方法和参数)
DataChangeTrigger.instance().onUpdate(syncType, syncType.getDataKey((T) instance), methodName, args);
System.out.println("intercept invoke:" + methodName);
@ -112,54 +150,16 @@ public class DataSyncProxyFactory {
}
public static interface A {
default void setName(String name) {
System.out.println("setName: " + name);
}
}
public static class B implements A {
private int age;
public void setAge(int age) {
this.age = age;
}
@Override
public void setName(String name) {
System.out.println("over setName: " + name);
}
}
public static void main(String[] args) {
RouterUser user1 = new RouterUser();
user1.setUid(10086L);
user1.setAccount("日月光");
RouterUser routerUser = newProxyInstance(SyncTypes.ROUTER_USER, user1);
System.out.println(routerUser);
routerUser.setOnlineTime(ImClock.millis());
System.out.println(routerUser);
RouterUser routerUser = newProxyInstance(SyncTypes.ROUTER_USER);
routerUser.setAccount("日月光");
System.out.println(routerUser.getAccount());
// Enhancer enhancer = new Enhancer();
// enhancer.setSuperclass(B.class);
// enhancer.setCallback(new MethodInterceptor() {
// @Override
// public Object intercept(Object o, Method method, Object[] args, MethodProxy methodProxy) throws Throwable {
// System.out.println("proxy method: " + method.getName());
// return methodProxy.invokeSuper(o, args);
// }
// });
// B b = (B) enhancer.create();
// b.setAge(12);
// b.setName("沧海");
//
// System.out.println(b.age);
//
// for (Method method : B.class.getMethods()) {
// System.out.println(method.getName() + "-" + method.hashCode() + ": " + method.getDeclaringClass());
// }
// System.out.println(Serializable.class.isAssignableFrom(String.class));
// System.out.println(Serializable.class.isAssignableFrom(Integer.class));
// System.out.println(Integer.class.isAssignableFrom(Serializable.class));
}
}

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

@ -1,10 +0,0 @@
package net.sopod.soim.router.datasync.server;
/**
* SyncDataIncr
*
* @author tmy
* @date 2022-05-05 22:33
*/
public class SyncDataIncr {
}

17
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLog.java

@ -43,8 +43,8 @@ public class SyncLog implements Serializable {
public static final int OPT_REMOVE = 2;
public static final int OPT_UPDATE = 3;
public static <T extends DataSync> AddLog<T> addLog(SyncTypes.SyncType<T> syncType) {
return new AddLog<>(syncType);
public static <T extends DataSync> AddLog<T> addLog(int logSeq, SyncTypes.SyncType<T> syncType) {
return new AddLog<>(logSeq, syncType);
}
public static <T extends DataSync> UpdateLog<T> updateLog(int logSeq, SyncTypes.SyncType<T> syncType) {
@ -79,7 +79,7 @@ public class SyncLog implements Serializable {
/** ================ 数据id标示:删除,更新用 ===================== */
protected String dataKey;
/** ================ 更新数据:类,更新方法,更新方法序列化参数(避免修改) ===================== */
/** ================ 更新数据:类,更新方法,更新方法序列化参数(避免修改) ===================== */
protected String clazz;
protected String method;
@ -280,9 +280,9 @@ public class SyncLog implements Serializable {
.setAccount("画中")
.setImEntryAddr("127.0.0.2")
.setIsOnline(true);
AddLog<RouterUser> addLog = addLog(SyncTypes.ROUTER_USER)
.addSerializeData(user1)
.addSerializeData(user2);
AddLog<RouterUser> addLog = addLog(0, SyncTypes.ROUTER_USER)
.addData(user1)
.addData(user2);
// String json = Jackson.json().serialize(addLog);
// System.out.println(json);
// System.out.println(json.getBytes(StandardCharsets.UTF_8).length);
@ -322,9 +322,10 @@ public class SyncLog implements Serializable {
* @param <T>
*/
public static class AddLog<T extends DataSync> extends SyncLog {
AddLog(SyncTypes.SyncType<T> syncType) {
AddLog(int logSeq, SyncTypes.SyncType<T> syncType) {
super(syncType);
this.operateType = OPT_ADD;
this.logSeq = logSeq;
}
public AddLog<T> setSerializeDataCollect(List<String> serializeDataCollect) {
@ -332,7 +333,7 @@ public class SyncLog implements Serializable {
return this;
}
public AddLog<T> addSerializeData(T data) {
public AddLog<T> addData(T data) {
if (this.serializeDataCollect == null) {
this.serializeDataCollect = new ArrayList<>();
}

19
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogInboundHandler.java

@ -1,19 +0,0 @@
package net.sopod.soim.router.datasync.server;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
/**
* SyncLogInboundHandler
*
* @author tmy
* @date 2022-05-05 14:57
*/
public class SyncLogInboundHandler extends SimpleChannelInboundHandler<SyncLog> {
@Override
protected void channelRead0(ChannelHandlerContext channelHandlerContext, SyncLog syncLog) throws Exception {
}
}

42
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java

@ -2,15 +2,19 @@ package net.sopod.soim.router.datasync.server;
import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.Channel;
import io.netty.channel.ChannelFutureListener;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.ChannelPipeline;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.nio.NioServerSocketChannel;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
import net.sopod.soim.common.util.netty.FastThreadLocalThreadFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.function.Consumer;
/**
* SyncServer
*
@ -21,15 +25,24 @@ public class SyncServer {
private static final Logger logger = LoggerFactory.getLogger(SyncServer.class);
public SyncServer() {
NioEventLoopGroup boss = new NioEventLoopGroup(1);
NioEventLoopGroup worker = new NioEventLoopGroup(2);
private final int port;
private NioEventLoopGroup boss;
private NioEventLoopGroup worker;
public SyncServer(int port) {
this.port = port;
}
public void start(Consumer<Throwable> onFail) throws InterruptedException {
this.boss = new NioEventLoopGroup(1, new FastThreadLocalThreadFactory("sync-server-boss-%d", Thread.NORM_PRIORITY));
this.worker = new NioEventLoopGroup(2, new FastThreadLocalThreadFactory("sync-server-worker-%d", Thread.NORM_PRIORITY));
ServerBootstrap serverBoot = new ServerBootstrap()
.group(boss, worker)
.channel(NioServerSocketChannel.class)
.childHandler(new ChannelInitializer<>() {
@Override
protected void initChannel(Channel channel) throws Exception {
protected void initChannel(Channel channel) {
LogLevel logLevel = logger.isDebugEnabled() ? LogLevel.DEBUG
: logger.isInfoEnabled() ? LogLevel.INFO
: logger.isWarnEnabled() ? LogLevel.WARN
@ -37,14 +50,27 @@ public class SyncServer {
ChannelPipeline pipeline = channel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel))
.addLast(new SyncLogDataCodec())
.addLast(new SyncLogInboundHandler());
.addLast(new SyncServerInboundHandler());
}
});
serverBoot.bind(8080);
serverBoot.bind(port).addListener((ChannelFutureListener) future -> {
if (!future.isSuccess()) {
if (onFail != null) {
onFail.accept(future.cause());
}
return;
}
logger.info("sync-server listening at {}...", port);
});
}
public void start() {
public void shutdown() {
if (boss != null) {
boss.shutdownGracefully();
}
if (worker != null) {
worker.shutdownGracefully();
}
}
}

33
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServerInboundHandler.java

@ -0,0 +1,33 @@
package net.sopod.soim.router.datasync.server;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
/**
* SyncLogInboundHandler
*
* @author tmy
* @date 2022-05-05 14:57
*/
public class SyncServerInboundHandler extends ChannelInboundHandlerAdapter {
public SyncServerInboundHandler() {
}
@Override
public void channelActive(ChannelHandlerContext ctx) throws Exception {
super.channelActive(ctx);
}
@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception {
super.channelInactive(ctx);
}
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
super.exceptionCaught(ctx, cause);
}
}
Loading…
Cancel
Save