Browse Source

im-router数据更改日志同步

master
tangmingyou 4 years ago
parent
commit
078f0e0b40
  1. 2
      im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java
  2. 6
      im-common/pom.xml
  3. 19
      im-common/src/main/java/net/sopod/soim/common/util/HashAlgorithms.java
  4. 73
      im-common/src/main/java/net/sopod/soim/common/util/Jackson.java
  5. 29
      im-common/src/main/java/net/sopod/soim/common/util/Reflects.java
  6. 35
      im-common/src/main/java/net/sopod/soim/common/util/cache/LinkedMapLRUCache.java
  7. 28
      im-common/src/main/java/net/sopod/soim/common/util/netty/Varint32FrameCodec.java
  8. 2
      im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java
  9. 4
      im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ConsistentHashTest.java
  10. 24
      im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java
  11. 32
      im-service/im-router/src/main/java/net/sopod/soim/router/config/AppContextHolder.java
  12. 45
      im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java
  13. 100
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java
  14. 80
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
  15. 59
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogByHashService.java
  16. 48
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java
  17. 4
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java
  18. 4
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java
  19. 1
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java
  20. 2
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java
  21. 15
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncLog.java
  22. 11
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java
  23. 2
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java
  24. 81
      im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java

2
im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java

@ -10,6 +10,7 @@ import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;
import net.sopod.soim.client.logger.Logger;
import net.sopod.soim.common.util.netty.Varint32FrameCodec;
import net.sopod.soim.core.net.ImMessageCodec;
import net.sopod.soim.data.msg.auth.Auth;
import net.sopod.soim.data.msg.chat.Chat;
@ -55,6 +56,7 @@ public class SoImSession {
@Override
protected void initChannel(SocketChannel ch) throws Exception {
ch.pipeline()
.addLast(new Varint32FrameCodec())
.addLast(new ImMessageCodec())
.addLast(messageDispatcher);
}

6
im-common/pom.xml

@ -36,6 +36,12 @@
<version>${netty.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>io.netty</groupId>
<artifactId>netty-codec</artifactId>
<version>${netty.version}</version>
<scope>provided</scope>
</dependency>
</dependencies>
</project>

19
im-common/src/main/java/net/sopod/soim/common/util/HashAlgorithms.java

@ -35,25 +35,6 @@ public class HashAlgorithms {
return hashCode & 0xffffffffL;
}
public static long md5Hash(String value, int number) {
MessageDigest md5;
try {
md5 = MessageDigest.getInstance("MD5");
} catch (NoSuchAlgorithmException e) {
throw new RuntimeException("MD5 not supported", e);
}
md5.reset();
byte[] keyBytes = value.getBytes(StandardCharsets.UTF_8);
md5.update(keyBytes);
byte[] digest = md5.digest();
// dubbo md5 hash algorithms
return (((long) (digest[3 + number * 4] & 0xFF) << 24)
| ((long) (digest[2 + number * 4] & 0xFF) << 16)
| ((long) (digest[1 + number * 4] & 0xFF) << 8)
| (digest[number * 4] & 0xFF))
& 0xFFFFFFFFL;
}
/**
* CRC系列算法本身并非查表可是查表是它的一种最快的实现方式以下是CRC32的实现

73
im-common/src/main/java/net/sopod/soim/common/util/Jackson.java

@ -10,6 +10,7 @@ import com.fasterxml.jackson.databind.*;
import com.fasterxml.jackson.databind.module.SimpleModule;
import com.fasterxml.jackson.databind.type.CollectionLikeType;
import com.fasterxml.jackson.databind.type.MapType;
import com.fasterxml.jackson.databind.type.TypeFactory;
import com.fasterxml.jackson.datatype.jsr310.deser.LocalDateDeserializer;
import com.fasterxml.jackson.datatype.jsr310.deser.LocalDateTimeDeserializer;
import com.fasterxml.jackson.datatype.jsr310.deser.LocalTimeDeserializer;
@ -20,6 +21,8 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.lang.reflect.ParameterizedType;
import java.lang.reflect.Type;
import java.text.SimpleDateFormat;
import java.time.LocalDate;
import java.time.LocalDateTime;
@ -312,8 +315,76 @@ public class Jackson {
}
private CollectionLikeType getListType(Class<?> elementClass) {
objectMapper.getTypeFactory().constructCollectionLikeType(List.class, getMapType(String.class, Object.class));
return objectMapper.getTypeFactory().constructCollectionLikeType(List.class, elementClass);
}
public ObjectMapper getObjectMapper() {
return objectMapper;
}
public TypeFactory getTypeFactory() {
return objectMapper.getTypeFactory();
}
/**
* @param content 序列化内容
* @param type 方法参数类型等可能包含泛型类型
* @return 解析结果
*/
public Object readWithType(String content, Type type) {
if (content == null || content.length() == 0) {
return null;
} else {
JavaType javaType = getJavaType(type);
try {
return objectMapper.readValue(content, javaType);
} catch (JsonProcessingException e) {
e.printStackTrace();
return null;
}
}
}
public Object readWithType(byte[] bytes, Type type) {
if (bytes == null || bytes.length == 0) {
return null;
} else {
JavaType javaType = getJavaType(type);
try {
return objectMapper.readValue(bytes, javaType);
} catch (IOException e) {
e.printStackTrace();
return null;
}
}
}
/**
* 解析任意层数泛型对象
* @param type 泛型类型
* @return jackson JavaType
*/
private JavaType getJavaType(Type type) {
//判断是否带有泛型
if (type instanceof ParameterizedType) {
Type[] actualTypeArguments = ((ParameterizedType) type).getActualTypeArguments();
System.out.println(Arrays.asList(actualTypeArguments));
//获取泛型类型
Class<?> rowClass = (Class<?>) ((ParameterizedType) type).getRawType();
JavaType[] javaTypes = new JavaType[actualTypeArguments.length];
for (int i = 0; i < actualTypeArguments.length; i++) {
//泛型也可能带有泛型,递归获取
javaTypes[i] = getJavaType(actualTypeArguments[i]);
}
return objectMapper.getTypeFactory().constructParametricType(rowClass, javaTypes);
} else {
//简单类型直接用该类构建JavaType
Class<?> cla = (Class<?>) type;
// TypeFactory.defaultInstance()
return objectMapper.getTypeFactory().constructParametricType(cla, new JavaType[0]);
}
}
}

29
im-common/src/main/java/net/sopod/soim/common/util/Reflects.java

@ -1,5 +1,7 @@
package net.sopod.soim.common.util;
import java.lang.reflect.Method;
import java.lang.reflect.Modifier;
import java.lang.reflect.Type;
import java.util.ArrayList;
import java.util.Arrays;
@ -56,4 +58,31 @@ public class Reflects {
return Arrays.asList(genericNames);
}
public static List<Method> getNonstaticMethods(Class<?> clazz, String name) {
return getNonstaticMethods(clazz, name, -1);
}
/**
* @param clazz
* @param name 方法名称
* @param paramSize 大于0匹配参数个数
* @return 匹配到的方法集合
*/
public static List<Method> getNonstaticMethods(Class<?> clazz, String name, int paramSize) {
List<Method> methodList = new ArrayList<>();
Method[] methods = clazz.getMethods();
for (Method method : methods) {
if (method.getName().equals(name)
&& !Modifier.isStatic(method.getModifiers())) {
if (paramSize < 0) {
methodList.add(method);
} else if (method.getParameterCount() == paramSize) {
// 参数个数匹配
methodList.add(method);
}
}
}
return methodList;
}
}

35
im-common/src/main/java/net/sopod/soim/common/util/cache/LinkedMapLRUCache.java vendored

@ -0,0 +1,35 @@
package net.sopod.soim.common.util.cache;
import java.util.LinkedHashMap;
import java.util.Map;
/**
* LinkedMapLRUCache
*
* @author tmy
* @date 2022-05-11 17:46
*/
public class LinkedMapLRUCache<K, V> extends LinkedHashMap<K, V> {
/**
* 缓存容量
*/
private final int maxEntries;
public LinkedMapLRUCache(int maxEntries) {
this.maxEntries = maxEntries;
}
/**
* 通过重写removeEldestEntry方法加入一定的条件满足条件返回true
*
* @param eldest 大链表头节点
* @return true表示允许移除头节点false表示不允许移除头节点
*/
@Override
protected boolean removeEldestEntry(Map.Entry eldest) {
//如果节点数量大于LRU缓存容量,那么返回true
return size() > maxEntries;
}
}

28
im-common/src/main/java/net/sopod/soim/common/util/netty/Varint32FrameCodec.java

@ -0,0 +1,28 @@
package net.sopod.soim.common.util.netty;
import io.netty.channel.CombinedChannelDuplexHandler;
import io.netty.handler.codec.protobuf.ProtobufVarint32FrameDecoder;
import io.netty.handler.codec.protobuf.ProtobufVarint32LengthFieldPrepender;
/**
* 帧编解码器
* 接收字节可能会分段到达添加帧编解码器修改应用程序代码以不调用 readVarInt writeVarInt
* 接收到的 ByteBuf 的大小将始终与发送的 ByteBuf 的大小匹配因为解码器负责聚合任何片段
* https://netty.io/4.0/xref/io/netty/example/worldclock/WorldClockServerInitializer.html
*
* {@link ProtobufVarint32FrameDecoder} 它确保管道中的下一个处理程序一次总是接收准确的一帧字节
* {@link ProtobufVarint32LengthFieldPrepender} 负责将帧长度 header 添加到所有传出消息
*
* @author tmy
* @date 2022-05-11 17:06
*/
public class Varint32FrameCodec extends CombinedChannelDuplexHandler<
ProtobufVarint32FrameDecoder,
ProtobufVarint32LengthFieldPrepender> {
public Varint32FrameCodec() {
// FrameDecoder, LengthFieldPrepender
super(new ProtobufVarint32FrameDecoder(), new ProtobufVarint32LengthFieldPrepender());
}
}

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

@ -5,6 +5,7 @@ import io.netty.channel.ChannelPipeline;
import io.netty.channel.socket.SocketChannel;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
import net.sopod.soim.common.util.netty.Varint32FrameCodec;
import net.sopod.soim.core.net.ImMessageCodec;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -25,6 +26,7 @@ public class ImEntryInitializer extends ChannelInitializer<SocketChannel> {
LogLevel logLevel = logger.isDebugEnabled() ? LogLevel.DEBUG : LogLevel.INFO;
ChannelPipeline pipeline = socketChannel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel))
.addLast(new Varint32FrameCodec())
.addLast(new ImMessageCodec())
.addLast(new InboundImMessageHandler());
}

4
im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ConsistentHashTest.java

@ -51,10 +51,6 @@ public class ConsistentHashTest {
return HashAlgorithms.md5Hash(value);
}
public static long hash(String value, int number) {
return HashAlgorithms.md5Hash(value, number);
}
private static void testConsistentHash() {
String[] nodes = {
"192.168.31.156:3032",

24
im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java

@ -4,6 +4,9 @@ import net.sopod.soim.common.util.HashAlgorithms;
import org.apache.commons.lang3.tuple.ImmutablePair;
import org.apache.commons.lang3.tuple.Pair;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.util.*;
/**
@ -91,7 +94,26 @@ public class UidConsistentHashSelector<V> {
}
private static long hash(String value, int number) {
return HashAlgorithms.md5Hash(value, number);
return md5Hash(value, number);
}
private static long md5Hash(String value, int number) {
MessageDigest md5;
try {
md5 = MessageDigest.getInstance("MD5");
} catch (NoSuchAlgorithmException e) {
throw new RuntimeException("MD5 not supported", e);
}
md5.reset();
byte[] keyBytes = value.getBytes(StandardCharsets.UTF_8);
md5.update(keyBytes);
byte[] digest = md5.digest();
// dubbo md5 hash algorithms
return (((long) (digest[3 + number * 4] & 0xFF) << 24)
| ((long) (digest[2 + number * 4] & 0xFF) << 16)
| ((long) (digest[1 + number * 4] & 0xFF) << 8)
| (digest[number * 4] & 0xFF))
& 0xFFFFFFFFL;
}
public int getIdentityHashCode() {

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

@ -5,13 +5,19 @@ 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.constant.DubboConstant;
import net.sopod.soim.common.util.Collects;
import net.sopod.soim.common.util.HashAlgorithms;
import net.sopod.soim.common.util.StringUtil;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.registry.Registry;
import org.apache.dubbo.registry.support.RegistryManager;
import org.apache.dubbo.rpc.model.ApplicationModel;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.ApplicationContext;
import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Properties;
@ -116,8 +122,7 @@ public class AppContextHolder {
}
}
try {
List<Instance> allInstances = namingService.getAllInstances(AppConstant.APP_IM_ROUTER_NAME);
return allInstances;
return namingService.getAllInstances(AppConstant.APP_IM_ROUTER_NAME);
} catch (NacosException e) {
throw new IllegalStateException(AppConstant.APP_IM_ROUTER_NAME + "集群信息获取失败:", e);
} finally {
@ -131,4 +136,27 @@ public class AppContextHolder {
}
}
/**
* 注册 im-router 的API接口服务
*/
public static void doRegistry() {
RegistryManager registryManager = ApplicationModel.defaultModel().getBeanFactory()
.getBean(RegistryManager.class);
Collection<Registry> registries = registryManager.getRegistries();
List<URL> registryInvokerUrls = AppContextHolder.getRegistryInvokerUrls();
if (Collects.isNotEmpty(registries)
&& Collects.isNotEmpty(registryInvokerUrls)) {
for (Registry registry : registries) {
for (URL invokerUrl : registryInvokerUrls) {
// 添加 im-router 服务id参数,生成新的 url
URL url = invokerUrl.addParameter(
DubboConstant.IM_ROUTER_ID_KEY,
AppContextHolder.IM_ROUTER_ID
);
registry.register(url);
}
}
}
}
}

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

@ -27,6 +27,10 @@ import org.springframework.context.annotation.Configuration;
import org.springframework.core.Ordered;
import java.util.*;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors;
/**
@ -58,13 +62,15 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
// 检查集群状态
boolean registryNow = this.checkClusterEnvironment(syncLogMigrateService);
// 服务已可用进行注册
// 服务已可用立即进行注册
if (registryNow) {
this.createMockData();
this.doRegistry();
AppContextHolder.doRegistry();
}
}
private AtomicLong uidCounter = new AtomicLong(10000);
private void createMockData() {
for (long i = 10000L; i < 11000L; i++) {
RouterUser routerUser = new RouterUser()
@ -74,6 +80,18 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
.setImEntryAddr("127.0.0.1:1313");
RouterUserStorage.getInstance().put(i, routerUser);
}
Executors.newSingleThreadScheduledExecutor()
.scheduleWithFixedDelay(() -> {
for (int i = 0; i < 10; i++) {
RouterUser routerUser = RouterUserStorage.getInstance().get(uidCounter.getAndIncrement());
if (routerUser == null) {
return;
}
routerUser.setAccount("changeAccount:" + uidCounter.get());
}
logger.info("change log....: {}", uidCounter.get());
}, 30, 10, TimeUnit.SECONDS);
}
/**
@ -86,6 +104,7 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
if (Collects.isNotEmpty(registries)) {
for (Registry registry : registries) {
// isServiceDiscovery(): true是注册应用(im-router)的registry, false是注册服务接口的registry
logger.info("registry: {}, {}", registry.getUrl().getAddress(), registry.isServiceDiscovery());
if (registry.isAvailable()
&& registry.isServiceDiscovery()) {
String discoveryAddr = registry.getUrl().getAddress();
@ -105,28 +124,6 @@ public class ImRouterAppOnReady implements ApplicationListener<ApplicationReadyE
.start(AppContextHolder.getAppPort() + SYNC_SERVER_PORT_OFFSET);
}
/**
* 注册 im-router 的API接口服务
*/
private void doRegistry() {
RegistryManager registryManager = ApplicationModel.defaultModel().getBeanFactory()
.getBean(RegistryManager.class);
Collection<Registry> registries = registryManager.getRegistries();
List<URL> registryInvokerUrls = AppContextHolder.getRegistryInvokerUrls();
if (Collects.isNotEmpty(registries)
&& Collects.isNotEmpty(registryInvokerUrls)) {
for (Registry registry : registries) {
for (URL invokerUrl : registryInvokerUrls) {
// 添加 im-router 服务id参数,生成新的 url
URL url = invokerUrl.addParameter(
DubboConstant.IM_ROUTER_ID_KEY,
AppContextHolder.IM_ROUTER_ID
);
registry.register(url);
}
}
}
}
@Override
public int getOrder() {

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

@ -1,8 +1,11 @@
package net.sopod.soim.router.datasync;
import net.sopod.soim.router.config.ImRouterAppOnReady;
import net.sopod.soim.router.datasync.server.data.SyncLog;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Queue;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicInteger;
@ -16,6 +19,8 @@ import java.util.concurrent.atomic.AtomicInteger;
*/
public class DataChangeTrigger {
private static final Logger logger = LoggerFactory.getLogger(DataChangeTrigger.class);
private static DataChangeTrigger INSTANCE;
// TODO subscribes 订阅 selector.select(dataKey);
@ -31,41 +36,79 @@ public class DataChangeTrigger {
return INSTANCE;
}
private final ConcurrentHashMap<String, AtomicInteger> seqCounterMap = new ConcurrentHashMap<>(128);
/**
* <dataKey, 数据日志序列号器>
* TODO 优化
*
*/
private final ConcurrentHashMap<String, AtomicInteger> seqCounterMap = new ConcurrentHashMap<>(1024);
private final Queue<SyncLog> logQueue = new ConcurrentLinkedQueue<>();
/**
* 监听者列表
*/
private final List<SyncLogSubscribe> subscribes = new ArrayList<>();
private List<SyncLogSubscribe> getSubscribes(String dataKey) {
if (subscribes.isEmpty()) {
return Collections.emptyList();
}
List<SyncLogSubscribe> acceptSubs = new ArrayList<>(5);
for (SyncLogSubscribe subscribe : subscribes) {
if (subscribe.accept(dataKey)) {
acceptSubs.add(subscribe);
}
}
return acceptSubs;
}
public <T extends DataSync> void onUpdate(SyncTypes.SyncType<T> syncType, String dataKey, String method, Object[] args) {
List<SyncLogSubscribe> acceptSubs;
if ((acceptSubs = getSubscribes(dataKey)).isEmpty()) {
return;
}
// 序列化 args,避免后续更改
SyncLog.UpdateLog<T> updateLog = SyncLog.updateLog(getSeq(dataKey), syncType)
.setDataKey(dataKey)
.setMethod(method)
.setArgs(args);
publishLog(updateLog);
publishLog(updateLog, acceptSubs);
}
public <T extends DataSync> void onAdd(SyncTypes.SyncType<T> syncType, T data) {
List<SyncLogSubscribe> acceptSubs;
String dataKey = syncType.getDataKey(data);
if ((acceptSubs = getSubscribes(dataKey)).isEmpty()) {
return;
}
AtomicInteger seqCounter = getSeqCounter(dataKey);
// 序列化 data,避免后续更改
SyncLog.AddLog<T> addLog = SyncLog.addLog(seqCounter.getAndIncrement(), syncType)
.addData(data);
publishLog(addLog);
publishLog(addLog, acceptSubs);
}
/**
* 数据删除日志
*/
public <T extends DataSync> void onRemove(SyncTypes.SyncType<T> syncType, String dataKey) {
List<SyncLogSubscribe> acceptSubs;
if ((acceptSubs = getSubscribes(dataKey)).isEmpty()) {
return;
}
SyncLog.RemoveLog<T> removeLog = SyncLog.removeLog(getSeq(dataKey), syncType)
.setDataKey(dataKey);
publishLog(removeLog);
publishLog(removeLog, acceptSubs);
}
private void publishLog(SyncLog log) {
// TODO 判断无订阅者跳过,新增节点同步按 dataKey 单独订阅每一个数据更新
logQueue.add(log);
System.out.println("publish log: "+log);
private void publishLog(SyncLog log, List<SyncLogSubscribe> acceptSubs) {
for (SyncLogSubscribe acceptSub : acceptSubs) {
try {
acceptSub.onSyncLog(log);
} catch (Exception e) {
logger.error("监听者 {} 执行错误", acceptSub, e);
}
}
}
private int getSeq(String dataKey) {
@ -76,6 +119,13 @@ public class DataChangeTrigger {
return seqCounterMap.computeIfAbsent(dataKey, key -> new AtomicInteger());
}
/**
* 添加更改日志监听者
*/
public void addSubscribe(SyncLogSubscribe syncLogSubscribe) {
subscribes.add(syncLogSubscribe);
}
/**
* 删除暂存数据字段等...
*/
@ -83,4 +133,34 @@ public class DataChangeTrigger {
}
public static interface SyncLogSubscribe {
public abstract boolean accept(String dataKey);
public abstract void onSyncLog(SyncLog syncLog);
}
public static abstract class DataKeySyncLogSubscribe implements SyncLogSubscribe {
private final Set<String> subscribeDataKeys;
public DataKeySyncLogSubscribe() {
this.subscribeDataKeys = new HashSet<>(1024);
}
@Override
public boolean accept(String dataKey) {
return subscribeDataKeys.contains(dataKey);
}
public void addSubscribeDataKey(String dataKey) {
subscribeDataKeys.add(dataKey);
}
@Override
public abstract void onSyncLog(SyncLog syncLog);
}
}

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

@ -145,57 +145,49 @@ public class DataSyncProxyFactory {
// 记录数据更新操作(方法和参数)
DataChangeTrigger.instance().onUpdate(syncType, syncType.getDataKey((T) instance), methodName, args);
System.out.println("intercept invoke:" + methodName);
// System.out.println("intercept invoke:" + methodName);
return methodProxy.invokeSuper(instance, args);
}
}
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);
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();
}
}
public static class A {
public static String name() {
return "1";
}
System.out.println(values.getClass());
System.out.println(values);
public void age() {}
}
public static class B extends A {
public void some() {}
public static void what() {}
}
public static void main(String[] args) {
// Enhancer enhancer = new Enhancer();
// enhancer.setSuperclass(B.class);
// enhancer.setCallback(new MethodInterceptor() {
// @Override
// public Object intercept(Object o, Method method, Object[] objects, MethodProxy methodProxy) throws Throwable {
// System.out.println("intercept:" + method.getName());
// return methodProxy.invokeSuper(o, objects);
// }
// });
// B b = (B)enhancer.create();
// b.some();
// b.age();
// b.name();
// b.what();
// b.equals(b);
// B.name();
for (Method m : B.class.getDeclaredMethods()) {
System.out.println(m.getName());
}
System.out.println("=============================");
for (Method m : B.class.getMethods()) {
System.out.println(m.getName());
}
}
}

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

@ -1,13 +1,11 @@
package net.sopod.soim.router.datasync.server;
package net.sopod.soim.router.datasync;
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.SyncCmd;
import net.sopod.soim.router.datasync.server.data.SyncLog;
import org.apache.commons.lang3.tuple.Pair;
import org.slf4j.Logger;
@ -18,13 +16,12 @@ import java.util.*;
import java.util.concurrent.atomic.AtomicInteger;
/**
* SyncLogPushService
* SyncLogByHashService
*
* @author tmy
* @date 2022-05-10 00:30
*/
public class SyncLogByHashService {
public class SyncLogByHashService extends DataChangeTrigger.DataKeySyncLogSubscribe {
private static final Logger logger = LoggerFactory.getLogger(SyncLogByHashService.class);
@ -37,7 +34,7 @@ public class SyncLogByHashService {
private final UidConsistentHashSelector<String> selector;
Set<String> syncedDataKeys = new HashSet<>();
// Set<String> syncedDataKeys = new HashSet<>();
public SyncLogByHashService(Channel clientChannel, String newNodeAddr) {
this.clientChannel = new WeakReference<>(clientChannel);
@ -47,18 +44,9 @@ public class SyncLogByHashService {
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);
// 添加数据变化监听
DataChangeTrigger.instance().addSubscribe(this);
}
private Map<DataSyncStorage<? extends DataSync>, SyncTypes.SyncType<? extends DataSync>> storages;
@ -102,8 +90,10 @@ public class SyncLogByHashService {
// TODO 数据加锁
addLog.addData(data);
// TODO DataChangeTrigger.instance().subscribe(syncType, dataKey)
// 注册更改日志
syncedDataKeys.add(dataKey);
super.addSubscribeDataKey(dataKey);
// TODO 解锁
count ++; totalCount ++;
@ -159,6 +149,35 @@ public class SyncLogByHashService {
private void syncFinishSuccess() {
logger.info("sync push finish and success!");
Channel channel = clientChannel.get();
if (channel != null && channel.isActive()) {
SyncCmd syncEndCmd = new SyncCmd();
syncEndCmd.setCmdType(SyncCmd.SYNC_END);
channel.writeAndFlush(syncEndCmd);
}
}
/**
* 修改数据日志
*/
@Override
public void onSyncLog(SyncLog syncLog) {
Channel channel = clientChannel.get();
if (channel != null && channel.isActive()) {
channel.writeAndFlush(syncLog);
}
}
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);
}
}

48
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncLogMigrateService.java

@ -12,10 +12,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
import java.util.Iterator;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
/**
@ -33,12 +30,16 @@ public class SyncLogMigrateService {
private volatile List<Pair<String, Integer>> migrateHosts;
private LinkedList<Pair<String, Integer>> curMigrateHosts;
private List<SyncClient> migrateClients = new ArrayList<>();
private SyncClient curClient;
/**
* 同步数据
* @param migrateHosts 同步数据节点
*/
public void migrateSyncLog(List<Pair<String, Integer>> migrateHosts) {
Preconditions.checkState(migrateHosts != null, "当前正在进行数据同步");
Preconditions.checkState(this.migrateHosts == null, "当前正在进行数据同步");
this.migrateHosts = migrateHosts;
this.curMigrateHosts = new LinkedList<>(migrateHosts);
logger.info("开始连接{}节点同步服务...", AppConstant.APP_IM_ROUTER_NAME);
@ -50,14 +51,22 @@ public class SyncLogMigrateService {
* @return 是否还有下一个同步数据节点
*/
public synchronized boolean syncNextHost() {
if (this.curClient != null) {
// TODO 优化多节点同时同步数据
throw new IllegalStateException("当前有客户端正在全量同步数据");
}
if (this.curMigrateHosts.isEmpty()) {
return false;
}
Pair<String, Integer> nextHost = this.curMigrateHosts.removeFirst();
if (nextHost == null) {
this.allHostSyncFinish();
return true;
}
try {
SyncClient client = new SyncClient();
client.connect(nextHost.getLeft(), nextHost.getRight());
client.syncLogByHash(AppContextHolder.getAppAddr());
this.curClient = new SyncClient();
this.curClient.connect(nextHost.getLeft(), nextHost.getRight());
this.curClient.syncLogByHash(AppContextHolder.getAppAddr());
} catch (InterruptedException e) {
logger.error("节点{}:{}连接失败, 跳过!", nextHost.getLeft(), nextHost.getRight(), e);
// 同步下一个节点
@ -66,6 +75,29 @@ public class SyncLogMigrateService {
return true;
}
public synchronized void curHostSyncEnd() {
if (this.curClient == null) {
logger.warn("当前无正在全量同步数据的客户端");
return;
}
this.migrateClients.add(this.curClient);
this.curClient = null;
this.syncNextHost();
}
/**
* 所有节点全量同步完成
* 注册服务
* 接受后续日志数据
* n秒无日志数据后发送消息可移除数据
* 断开连接
*/
private void allHostSyncFinish() {
logger.info("所有节点数据同步完成:执行注册服务....");
AppContextHolder.doRegistry();
logger.info("注册服务成功");
}
@Deprecated
private void fullSync() {
Map<Long, RouterUser> routerUserMap = RouterUserStorage.getInstance().getRouterUserMap();

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

@ -9,6 +9,7 @@ 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.common.util.netty.Varint32FrameCodec;
import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.datasync.SyncTypes;
import net.sopod.soim.router.datasync.server.codec.SyncCmdCodec;
@ -38,8 +39,11 @@ public class SyncClient {
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel channel) throws Exception {
channel.pipeline()
.addLast(new LoggingHandler(LogLevel.INFO))
// 接收字节可能会分段到达,添加帧编解码器
.addLast(new Varint32FrameCodec())
.addLast(new SyncCmdCodec())
.addLast(new SyncLogEncoder())
.addLast(new SyncCmdClientHandler())

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

@ -10,6 +10,7 @@ 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 net.sopod.soim.common.util.netty.Varint32FrameCodec;
import net.sopod.soim.router.datasync.server.codec.SyncCmdCodec;
import net.sopod.soim.router.datasync.server.codec.SyncLogEncoder;
import net.sopod.soim.router.datasync.server.handler.SyncCmdServerHandler;
@ -49,8 +50,11 @@ public class SyncServer {
: logger.isInfoEnabled() ? LogLevel.INFO
: logger.isWarnEnabled() ? LogLevel.WARN
: logger.isErrorEnabled() ? LogLevel.ERROR : LogLevel.INFO;
ChannelPipeline pipeline = channel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel))
// 接收字节可能会分段到达,添加帧编解码器
.addLast(new Varint32FrameCodec())
.addLast(new SyncCmdCodec())
.addLast(new SyncLogEncoder())
.addLast(new SyncCmdServerHandler())

1
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java

@ -1,6 +1,7 @@
package net.sopod.soim.router.datasync.server.codec;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.PooledByteBufAllocator;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.CombinedChannelDuplexHandler;

2
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java

@ -42,7 +42,7 @@ public class SyncCmd {
/**
* 同步结束命令
*/
public static final Integer SYNC_END = 8;
public static final int SYNC_END = 8;
/**
* SyncLog 推送命令

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

@ -107,7 +107,6 @@ public class SyncLog implements Serializable {
private byte[] toBytes0() {
byte[] dataKeyBytes = dataKey == null ? new byte[0] : dataKey.getBytes();
// byte[] clazzBytes = clazz == null ? new byte[0] : clazz.getBytes();
byte[] methodBytes = method == null ? new byte[0] : method.getBytes();
int argSize = args == null ? 0 : args.length;
@ -144,7 +143,7 @@ public class SyncLog implements Serializable {
+ 4 // logSeq 序列号
+ 4 // dataKey 数据主键标识字节长度
+ dataKeyBytes.length // dataKey 数据字节
+ 4 // clazz 字节长度
// + 4 // clazz 字节长度
// + clazzBytes.length // clazz字节
+ 4 // method 字节长度
+ methodBytes.length // method 字节
@ -157,12 +156,12 @@ public class SyncLog implements Serializable {
// Bytes.short2bytes(MAGIC, bytes, offset);
// offset += 2;
Bytes.int2bytes(bytes.length - 6, bytes, offset);
Bytes.int2bytes(bytes.length - 4, bytes, offset);
offset += 4;
bytes[offset] = (byte) syncType;
offset += 1;
bytes[offset] = (byte)operateType;
bytes[offset] = (byte) operateType;
offset += 1;
Bytes.int2bytes(logSeq, bytes, offset);
@ -237,14 +236,6 @@ public class SyncLog implements Serializable {
int argSize = buf.readInt();
log.args = new String[argSize];
if (argSize > 0) {
// Method method = getClassMethod(log.clazz, log.method);
// if (method == null) {
// throw new IllegalStateException("类" + log.clazz + "方法" + log.method + "未找到");
// }
// Class<?>[] paramTypes = method.getParameterTypes();
// if (paramTypes.length != argSize) {
// throw new IllegalStateException("类" + log.clazz + "方法" + log.method + "指定参数" + argSize + "个,查到参数" + paramTypes.length + "个");
// }
// 复用 bytes
byte[] argBytes = new byte[0];
for (int i = 0; i < argSize; i++) {

11
im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java

@ -2,6 +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.config.AppContextHolder;
import net.sopod.soim.router.datasync.SyncLogMigrateService;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -25,6 +27,9 @@ public class SyncCmdClientHandler extends SimpleChannelInboundHandler<SyncCmd> {
case SyncCmd.PONG:
this.handlePong(ctx, syncCmd);
break;
case SyncCmd.SYNC_END:
this.handleSyncEnd(ctx, syncCmd);
break;
}
}
@ -39,4 +44,10 @@ public class SyncCmdClientHandler extends SimpleChannelInboundHandler<SyncCmd> {
System.out.println("pong: " + ctx.channel());
}
private void handleSyncEnd(ChannelHandlerContext ctx, SyncCmd syncCmd) {
logger.info("sync end: {}", ctx.channel());
SyncLogMigrateService migrateService = AppContextHolder.getBean(SyncLogMigrateService.class);
migrateService.curHostSyncEnd();
}
}

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

@ -2,7 +2,7 @@ package net.sopod.soim.router.datasync.server.handler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import net.sopod.soim.router.datasync.server.SyncLogByHashService;
import net.sopod.soim.router.datasync.SyncLogByHashService;
import net.sopod.soim.router.datasync.server.data.SyncCmd;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

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

@ -2,8 +2,12 @@ package net.sopod.soim.router.datasync.server.handler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import net.sopod.soim.common.util.Collects;
import net.sopod.soim.common.util.Jackson;
import net.sopod.soim.common.util.Reflects;
import net.sopod.soim.common.util.cache.LinkedMapLRUCache;
import net.sopod.soim.router.cache.RouterUserStorage;
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;
@ -11,7 +15,10 @@ import net.sopod.soim.router.datasync.server.data.SyncLog;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.lang.reflect.InvocationTargetException;
import java.lang.reflect.Method;
import java.lang.reflect.Type;
import java.util.*;
/**
* SyncLogHandler
@ -26,17 +33,75 @@ public class SyncLogClientHandler extends SimpleChannelInboundHandler<SyncLog> {
@Override
protected void channelRead0(ChannelHandlerContext ctx, SyncLog syncLog) {
List<byte[]> bytesList = syncLog.getSerializeDataCollect();
switch (syncLog.getOperateType()) {
case SyncLog.OPT_ADD:
this.handleAddSyncLog(ctx, syncLog);
break;
case SyncLog.OPT_REMOVE:
this.handleRemoveSyncLog(ctx, syncLog);
break;
case SyncLog.OPT_UPDATE:
this.handleUpdateSyncLog(ctx, syncLog);
break;
}
}
private void handleAddSyncLog(ChannelHandlerContext ctx, SyncLog syncLog) {
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);
List<byte[]> bytesList = syncLog.getSerializeDataCollect();
if (Collects.isNotEmpty(bytesList)) {
for (byte[] bytes : bytesList) {
DataSync instance = CodecUtil.decode(bytes, syncType.dataType());
syncType.addData(instance);
logger.info("addData: {}", instance);
}
}
// 同步完成响应
SyncCmd syncCmd = new SyncCmd().setCmdType(SyncCmd.SYNC_BY_HASH_ACK);
ctx.writeAndFlush(syncCmd);
logger.info("storage: {}", DataSyncStorage.getStorages());
logger.info("users size: {}", RouterUserStorage.getInstance().getRouterUserMap().size());
}
private void handleRemoveSyncLog(ChannelHandlerContext ctx, SyncLog syncLog) {
SyncTypes.SyncType<DataSync> syncType = SyncTypes.getSyncType(syncLog.getSyncType());
syncType.removeData(syncLog.getDataKey());
}
private static final LinkedMapLRUCache<String, Method> syncTypeMethodCache = new LinkedMapLRUCache<>(16);
private void handleUpdateSyncLog(ChannelHandlerContext ctx, SyncLog syncLog) {
SyncTypes.SyncType<DataSync> syncType = SyncTypes.getSyncType(syncLog.getSyncType());
String methodName = syncLog.getMethod();
String[] args = syncLog.getArgs();
Class<DataSync> dataType = syncType.dataType();
String cacheKey = dataType.getName() + "." + methodName + "(" + args.length + ")";
Method method = syncTypeMethodCache.get(cacheKey);
if (method == null) {
synchronized(syncTypeMethodCache) {
if (null == (method = syncTypeMethodCache.get(cacheKey))) {
List<Method> methods = Reflects.getNonstaticMethods(dataType, methodName, args.length);
if (Collects.isEmpty(methods)) {
logger.warn("{} method {} not found!", dataType, methodName);
return;
}
method = methods.get(0);
syncTypeMethodCache.put(cacheKey, method);
}
}
}
Type[] types = method.getGenericParameterTypes();
Object[] params = new Object[types.length];
for (int i = 0; i < types.length; i++) {
params[i] = Jackson.json().readWithType(args[i], types[i]);
}
DataSync data = syncType.getData(syncLog.getDataKey());
try {
method.setAccessible(true);
method.invoke(data, params);
logger.info("invoke:{}, args:{}", method.getName(), params);
} catch (IllegalAccessException | InvocationTargetException e) {
e.printStackTrace();
}
}
}

Loading…
Cancel
Save