diff --git a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java index e265bdd77..79c5128de 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java @@ -21,7 +21,7 @@ import java.net.InetSocketAddress; import com.ai.cloud.skywalking.selfexamination.HeathReading; import com.ai.cloud.skywalking.selfexamination.SDKHealthCollector; -import com.ai.cloud.skywalking.util.ProtocolPackager; +import com.ai.cloud.skywalking.util.TransportPackager; public class DataSender implements IDataSender { private EventLoopGroup group; @@ -76,7 +76,7 @@ public class DataSender implements IDataSender { try { if (channel != null && channel.isActive()) { - byte[] dataPackage = ProtocolPackager.pack(data.getBytes()); + byte[] dataPackage = TransportPackager.pack(data.getBytes()); channel.writeAndFlush(dataPackage); SDKHealthCollector.getCurrentHeathReading("sender").updateData(HeathReading.INFO, "DataSender[" + socketAddress + "] send data successfully."); diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/exception/SerializableDataTypeRegisterException.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/exception/SerializableDataTypeRegisterException.java new file mode 100644 index 000000000..82a40e368 --- /dev/null +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/exception/SerializableDataTypeRegisterException.java @@ -0,0 +1,11 @@ +package com.ai.cloud.skywalking.exception; + +/** + * Created by wusheng on 16/7/4. + */ +public class SerializableDataTypeRegisterException extends RuntimeException { + + public SerializableDataTypeRegisterException(String message) { + super(message); + } +} diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AbstractDataSerializable.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AbstractDataSerializable.java new file mode 100644 index 000000000..ef2110c34 --- /dev/null +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AbstractDataSerializable.java @@ -0,0 +1,42 @@ +package com.ai.cloud.skywalking.protocol; + +import com.ai.cloud.skywalking.util.IntegerAssist; + +import java.util.HashSet; +import java.util.Set; + +/** + * Created by wusheng on 16/7/4. + */ +public abstract class AbstractDataSerializable implements ISerializable, NullableClass{ + private static Set DATA_TYPE_SCOPE = new HashSet(); + + public AbstractDataSerializable(){ + SerializableDataTypeRegister.init(getDataType(), this.getClass()); + } + + public abstract int getDataType(); + + public abstract byte[] getData(); + + @Override + public byte[] convert2Bytes() { + byte[] type = IntegerAssist.intToBytes(SerializableDataTypeRegister.getType(this.getClass())); + + //TODO:类型+ data = 消息包 + return getData(); + } + + @Override + public Object convert2Object(byte[] data) { + // TODO:data的前4位转成type; + int dataType = 1; + if(!SerializableDataTypeRegister.isTypeAndClassMatch(dataType, this.getClass())){ + return new NullClass(); + } + // TODO: 反序列化 + return null; + } + + +} diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/ISerializable.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/ISerializable.java new file mode 100644 index 000000000..58c532b84 --- /dev/null +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/ISerializable.java @@ -0,0 +1,11 @@ +package com.ai.cloud.skywalking.protocol; + +/** + * Created by wusheng on 16/7/4. + */ +public interface ISerializable { + byte[] convert2Bytes(); + + Object convert2Object(byte[] data); + +} diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/NullClass.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/NullClass.java new file mode 100644 index 000000000..d9a8984f4 --- /dev/null +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/NullClass.java @@ -0,0 +1,11 @@ +package com.ai.cloud.skywalking.protocol; + +/** + * Created by wusheng on 16/7/4. + */ +public class NullClass implements NullableClass { + @Override + public boolean isNull() { + return true; + } +} diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/NullableClass.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/NullableClass.java new file mode 100644 index 000000000..f50f9c216 --- /dev/null +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/NullableClass.java @@ -0,0 +1,9 @@ +package com.ai.cloud.skywalking.protocol; + +/** + * Created by wusheng on 16/7/4. + */ +public interface NullableClass { + boolean isNull(); + +} diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializableDataTypeRegister.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializableDataTypeRegister.java new file mode 100644 index 000000000..9fefeddd4 --- /dev/null +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializableDataTypeRegister.java @@ -0,0 +1,43 @@ +package com.ai.cloud.skywalking.protocol; + +import com.ai.cloud.skywalking.exception.SerializableDataTypeRegisterException; + +import java.util.HashMap; +import java.util.Map; + +/** + * Created by wusheng on 16/7/4. + */ +public class SerializableDataTypeRegister { + private static Map> DATA_TYPE_MAPPING_CLASS = new HashMap>(); + + private static Map, Integer> CLASS_MAPPING_DATA_TYPE = new HashMap, Integer>(); + + public static void init(Integer dataType, Class clazz) { + if (DATA_TYPE_MAPPING_CLASS.containsKey(dataType)) { + if (!DATA_TYPE_MAPPING_CLASS.get(dataType).equals(clazz)) { + throw new SerializableDataTypeRegisterException("dataType=" + dataType + " has been registered to " + DATA_TYPE_MAPPING_CLASS.get(dataType)); + } else { + return; + } + } + DATA_TYPE_MAPPING_CLASS.put(dataType, clazz); + CLASS_MAPPING_DATA_TYPE.put(clazz, dataType); + } + + public static boolean isTypeAndClassMatch(Integer dataType, Class clazz) { + if (DATA_TYPE_MAPPING_CLASS.containsKey(dataType)) { + return true; + } else { + return false; + } + } + + public static Integer getType(Class clazz) { + if (CLASS_MAPPING_DATA_TYPE.containsKey(clazz)) { + return CLASS_MAPPING_DATA_TYPE.get(clazz); + } else { + throw new SerializableDataTypeRegisterException("class " + clazz + " not found in SerializableDataTypeRegister."); + } + } +} diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/IntegerAssist.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/IntegerAssist.java new file mode 100644 index 000000000..76c18155e --- /dev/null +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/IntegerAssist.java @@ -0,0 +1,15 @@ +package com.ai.cloud.skywalking.util; + +/** + * Created by wusheng on 16/7/4. + */ +public class IntegerAssist { + public static byte[] intToBytes(int value) { + byte[] src = new byte[4]; + src[0] = (byte) ((value >> 24) & 0xFF); + src[1] = (byte) ((value >> 16) & 0xFF); + src[2] = (byte) ((value >> 8) & 0xFF); + src[3] = (byte) (value & 0xFF); + return src; + } +} diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/ProtocolPackager.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/TransportPackager.java similarity index 82% rename from skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/ProtocolPackager.java rename to skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/TransportPackager.java index e162612c0..2eb326411 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/ProtocolPackager.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/TransportPackager.java @@ -2,7 +2,7 @@ package com.ai.cloud.skywalking.util; import java.util.Arrays; -public class ProtocolPackager { +public class TransportPackager { public static byte[] pack(byte[] data) { // 对协议格式进行修改 // | check sum(4 byte) | data @@ -58,16 +58,6 @@ public class ProtocolPackager { result ^= data[i]; } - return intToBytes(result); + return IntegerAssist.intToBytes(result); } - - private static byte[] intToBytes(int value) { - byte[] src = new byte[4]; - src[0] = (byte) ((value >> 24) & 0xFF); - src[1] = (byte) ((value >> 16) & 0xFF); - src[2] = (byte) ((value >> 8) & 0xFF); - src[3] = (byte) (value & 0xFF); - return src; - } - } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java index a63966637..0a0938757 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java @@ -3,7 +3,7 @@ package com.ai.cloud.skywalking.reciever.handler; import com.ai.cloud.skywalking.reciever.buffer.DataBufferThreadContainer; import com.ai.cloud.skywalking.reciever.conf.Config; import com.ai.cloud.skywalking.reciever.util.RedisConnector; -import com.ai.cloud.skywalking.util.ProtocolPackager; +import com.ai.cloud.skywalking.util.TransportPackager; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.SimpleChannelInboundHandler; import redis.clients.jedis.Jedis; @@ -18,7 +18,7 @@ public class CollectionServerDataHandler extends SimpleChannelInboundHandler= 0 && msg.length < Config.DataPackage.MAX_DATA_PACKAGE) { - byte[] data = ProtocolPackager.unpack(msg); + byte[] data = TransportPackager.unpack(msg); if (data != null) { DataBufferThreadContainer.getDataBufferThread().saveTemporarily(data);