修改部分接口,供参考。
This commit is contained in:
parent
f6a38d34f5
commit
b5e86d0330
|
|
@ -5,11 +5,15 @@ import com.ai.cloud.skywalking.conf.Constants;
|
|||
import com.ai.cloud.skywalking.logging.LogManager;
|
||||
import com.ai.cloud.skywalking.logging.Logger;
|
||||
import com.ai.cloud.skywalking.protocol.Span;
|
||||
import com.ai.cloud.skywalking.protocol.common.ISerializable;
|
||||
import com.ai.cloud.skywalking.selfexamination.HeathReading;
|
||||
import com.ai.cloud.skywalking.selfexamination.SDKHealthCollector;
|
||||
import com.ai.cloud.skywalking.sender.DataSenderFactoryWithBalance;
|
||||
import com.ai.cloud.skywalking.util.AtomicRangeInteger;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
import static com.ai.cloud.skywalking.conf.Config.Buffer.BUFFER_MAX_SIZE;
|
||||
import static com.ai.cloud.skywalking.conf.Config.Consumer.*;
|
||||
|
||||
|
|
@ -17,7 +21,7 @@ public class BufferGroup {
|
|||
private static Logger logger = LogManager.getLogger(BufferGroup.class);
|
||||
private String groupName;
|
||||
//注意: 修改这个变量名,需要修改test-api工程的Config类中的SPAN_ARRAY_FIELD_NAME变量
|
||||
private Span[] dataBuffer = new Span[BUFFER_MAX_SIZE];
|
||||
private ISerializable[] dataBuffer = new ISerializable[BUFFER_MAX_SIZE];
|
||||
AtomicRangeInteger index = new AtomicRangeInteger(0, BUFFER_MAX_SIZE);
|
||||
|
||||
public BufferGroup(String groupName) {
|
||||
|
|
@ -65,7 +69,8 @@ public class BufferGroup {
|
|||
|
||||
@Override
|
||||
public void run() {
|
||||
StringBuilder data = new StringBuilder();
|
||||
//TODO: 根据list长度进行限制
|
||||
List<ISerializable> packageData = new ArrayList<ISerializable>();
|
||||
while (true) {
|
||||
boolean bool = false;
|
||||
try {
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@ package com.ai.cloud.skywalking.buffer;
|
|||
import com.ai.cloud.skywalking.logging.LogManager;
|
||||
import com.ai.cloud.skywalking.logging.Logger;
|
||||
import com.ai.cloud.skywalking.protocol.Span;
|
||||
import com.ai.cloud.skywalking.protocol.common.ISerializable;
|
||||
|
||||
import java.util.concurrent.ThreadLocalRandom;
|
||||
|
||||
|
|
@ -18,9 +19,9 @@ public class ContextBuffer {
|
|||
//non
|
||||
}
|
||||
|
||||
public static void save(Span span) {
|
||||
public static void save(ISerializable data) {
|
||||
try{
|
||||
pool.save(span);
|
||||
pool.save(data);
|
||||
}catch(Throwable t){
|
||||
logger.error("save span error.", t);
|
||||
}
|
||||
|
|
@ -36,8 +37,8 @@ public class ContextBuffer {
|
|||
}
|
||||
}
|
||||
|
||||
public void save(Span span) {
|
||||
bufferGroups[ThreadLocalRandom.current().nextInt(0, POOL_SIZE)].save(span);
|
||||
public void save(ISerializable data) {
|
||||
bufferGroups[ThreadLocalRandom.current().nextInt(0, POOL_SIZE)].save(data);
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
package com.ai.cloud.skywalking.sender;
|
||||
|
||||
import com.ai.cloud.skywalking.protocol.common.ISerializable;
|
||||
import io.netty.bootstrap.Bootstrap;
|
||||
import io.netty.channel.Channel;
|
||||
import io.netty.channel.ChannelHandlerContext;
|
||||
|
|
@ -18,6 +19,7 @@ import io.netty.handler.codec.bytes.ByteArrayEncoder;
|
|||
|
||||
import java.io.IOException;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.util.List;
|
||||
|
||||
import com.ai.cloud.skywalking.selfexamination.HeathReading;
|
||||
import com.ai.cloud.skywalking.selfexamination.SDKHealthCollector;
|
||||
|
|
@ -72,7 +74,7 @@ public class DataSender implements IDataSender {
|
|||
* @return
|
||||
*/
|
||||
@Override
|
||||
public boolean send(String data) {
|
||||
public boolean send(List<ISerializable> packageData) {
|
||||
try {
|
||||
if (channel != null && channel.isActive()) {
|
||||
|
||||
|
|
|
|||
|
|
@ -3,8 +3,10 @@ package com.ai.cloud.skywalking.sender;
|
|||
import static com.ai.cloud.skywalking.conf.Config.Sender.MAX_COPY_NUM;
|
||||
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
|
||||
import com.ai.cloud.skywalking.protocol.common.ISerializable;
|
||||
import com.ai.cloud.skywalking.selfexamination.HeathReading;
|
||||
import com.ai.cloud.skywalking.selfexamination.SDKHealthCollector;
|
||||
|
||||
|
|
@ -48,7 +50,7 @@ public class DataSenderWithCopies implements IDataSender {
|
|||
/**
|
||||
* 尝试向所有副本发送
|
||||
*/
|
||||
public boolean send(String data) {
|
||||
public boolean send(List<ISerializable> packageData) {
|
||||
int successNum = 0;
|
||||
for (IDataSender sender : senders) {
|
||||
if (sender.send(data)) {
|
||||
|
|
|
|||
|
|
@ -9,33 +9,33 @@ public class AckSpan extends AbstractDataSerializable {
|
|||
/**
|
||||
* tid,调用链的全局唯一标识
|
||||
*/
|
||||
protected String traceId;
|
||||
private String traceId;
|
||||
/**
|
||||
* 当前调用链的上级描述<br/>
|
||||
* 如当前序号为:0.1.0时,parentLevel=0.1
|
||||
*/
|
||||
protected String parentLevel;
|
||||
private String parentLevel;
|
||||
/**
|
||||
* 当前调用链的本机描述<br/>
|
||||
* 如当前序号为:0.1.0时,levelId=0
|
||||
*/
|
||||
protected int levelId = 0;
|
||||
private int levelId = 0;
|
||||
/**
|
||||
* 节点调用花费时间
|
||||
*/
|
||||
protected long cost = 0L;
|
||||
private long cost = 0L;
|
||||
/**
|
||||
* 节点调用的状态<br/>
|
||||
* 0:成功<br/>
|
||||
* 1:异常<br/>
|
||||
* 异常判断原则:代码产生exception,并且此exception不在忽略列表中
|
||||
*/
|
||||
protected byte statusCode = 0;
|
||||
private byte statusCode = 0;
|
||||
/**
|
||||
* 节点调用的错误堆栈<br/>
|
||||
* 堆栈以JAVA的exception为主要判断依据
|
||||
*/
|
||||
protected String exceptionStack;
|
||||
private String exceptionStack;
|
||||
|
||||
public String getTraceId() {
|
||||
return traceId;
|
||||
|
|
|
|||
|
|
@ -4,6 +4,9 @@ import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable;
|
|||
import com.ai.cloud.skywalking.protocol.common.CallType;
|
||||
import com.ai.cloud.skywalking.protocol.common.SpanType;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* Created by wusheng on 16/7/4.
|
||||
*/
|
||||
|
|
@ -11,66 +14,69 @@ public class RequestSpan extends AbstractDataSerializable {
|
|||
/**
|
||||
* tid,调用链的全局唯一标识
|
||||
*/
|
||||
protected String traceId;
|
||||
private String traceId;
|
||||
/**
|
||||
* 当前调用链的上级描述<br/>
|
||||
* 如当前序号为:0.1.0时,parentLevel=0.1
|
||||
*/
|
||||
protected String parentLevel;
|
||||
private String parentLevel;
|
||||
/**
|
||||
* 当前调用链的本机描述<br/>
|
||||
* 如当前序号为:0.1.0时,levelId=0
|
||||
*/
|
||||
protected int levelId = 0;
|
||||
private int levelId = 0;
|
||||
/**
|
||||
* 调用链中单个节点的入口描述<br/>
|
||||
* 如:java方法名,调用的RPC地址等等
|
||||
*/
|
||||
protected String viewPointId = "";
|
||||
private String viewPointId = "";
|
||||
/**
|
||||
* 节点调用开始时间
|
||||
*/
|
||||
protected long startDate = System.currentTimeMillis();
|
||||
private long startDate = System.currentTimeMillis();
|
||||
/**
|
||||
* 节点调用的发生机器描述<br/>
|
||||
* 包含机器名 + IP地址
|
||||
*/
|
||||
protected String address = "";
|
||||
private String address = "";
|
||||
|
||||
/**
|
||||
* 节点类型描述<br/>
|
||||
* 已字符串的形式描述<br/>
|
||||
* 如:java,dubbo等
|
||||
*/
|
||||
protected String spanTypeDesc = "";
|
||||
private String spanTypeDesc = "";
|
||||
|
||||
/**
|
||||
* 节点调用类型描述<br/>
|
||||
*
|
||||
* @see CallType
|
||||
*/
|
||||
protected String callType = "";
|
||||
private String callType = "";
|
||||
|
||||
/**
|
||||
* 节点分布式类型<br/>
|
||||
* 本地调用 / RPC服务端 / RPC客户端
|
||||
*/
|
||||
protected SpanType spanType = SpanType.LOCAL;
|
||||
private SpanType spanType = SpanType.LOCAL;
|
||||
|
||||
/**
|
||||
* 节点调用的所在进程号
|
||||
*/
|
||||
protected String processNo = "";
|
||||
private String processNo = "";
|
||||
/**
|
||||
* 节点调用所在的系统逻辑名称<br/>
|
||||
* 由授权文件指定
|
||||
*/
|
||||
protected String applicationId = "";
|
||||
private String applicationId = "";
|
||||
|
||||
/**
|
||||
* 用户id<br/>
|
||||
* 由授权文件指定
|
||||
*/
|
||||
protected String userId;
|
||||
private String userId;
|
||||
|
||||
private Map<String, String> paramters = new HashMap<String, String>();
|
||||
|
||||
public String getTraceId() {
|
||||
return traceId;
|
||||
|
|
@ -168,6 +174,14 @@ public class RequestSpan extends AbstractDataSerializable {
|
|||
this.userId = userId;
|
||||
}
|
||||
|
||||
public Map<String, String> getParamters() {
|
||||
return paramters;
|
||||
}
|
||||
|
||||
public void setParamters(Map<String, String> paramters) {
|
||||
this.paramters = paramters;
|
||||
}
|
||||
|
||||
@Override
|
||||
public int getDataType() {
|
||||
return 1;
|
||||
|
|
|
|||
|
|
@ -21,11 +21,18 @@ public abstract class AbstractDataSerializable implements ISerializable, Nullabl
|
|||
|
||||
public abstract byte[] getData();
|
||||
|
||||
/**
|
||||
* 消息包结构:
|
||||
* 4位消息体类型
|
||||
* n位数据正文
|
||||
*
|
||||
* @return
|
||||
*/
|
||||
@Override
|
||||
public byte[] convert2Bytes() {
|
||||
byte[] type = IntegerAssist.intToBytes(SerializableDataTypeRegister.getType(this.getClass()));
|
||||
|
||||
//TODO:类型+ data = 消息包
|
||||
//TODO:消息包 = 4位
|
||||
return getData();
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,9 +1,12 @@
|
|||
package com.ai.cloud.skywalking.util;
|
||||
|
||||
import com.ai.cloud.skywalking.protocol.common.ISerializable;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.List;
|
||||
|
||||
public class TransportPackager {
|
||||
public static byte[] pack(byte[] data) {
|
||||
public static byte[] pack(List<ISerializable> packageData) {
|
||||
// 对协议格式进行修改
|
||||
// | check sum(4 byte) | data
|
||||
byte[] dataPackage = new byte[data.length + 4];
|
||||
|
|
@ -24,7 +27,7 @@ public class TransportPackager {
|
|||
}
|
||||
|
||||
|
||||
public static byte[] unpack(byte[] dataPackage) {
|
||||
public static List<ISerializable> unpack(byte[] dataPackage) {
|
||||
if (validateCheckSum(dataPackage)) {
|
||||
return unpackDataText(dataPackage);
|
||||
} else {
|
||||
|
|
|
|||
Loading…
Reference in New Issue