完成入HBase库的测试
This commit is contained in:
parent
603d9ba0f7
commit
89c21a60bf
|
|
@ -14,7 +14,7 @@ import static com.ai.cloud.skywalking.conf.Config.Consumer.*;
|
|||
import static com.ai.cloud.skywalking.conf.Config.Sender.MAX_BUFFER_DATA_SIZE;
|
||||
|
||||
public class BufferGroup {
|
||||
private static Logger logger = Logger.getLogger(BufferGroup.class.getName());
|
||||
// private static Logger logger = Logger.getLogger(BufferGroup.class.getName());
|
||||
private String groupName;
|
||||
private Span[] dataBuffer = new Span[BUFFER_MAX_SIZE];
|
||||
AtomicInteger index = new AtomicInteger(0);
|
||||
|
|
@ -72,14 +72,14 @@ public class BufferGroup {
|
|||
continue;
|
||||
}
|
||||
bool = true;
|
||||
data.append(dataBuffer[i]);
|
||||
data.append(dataBuffer[i] + ";");
|
||||
dataBuffer[i] = null;
|
||||
if (index++ == MAX_BUFFER_DATA_SIZE || data.length() >= Config.Sender.MAX_SEND_LENGTH) {
|
||||
while (!DataSenderFactory.getSender().send(data.toString())) {
|
||||
try {
|
||||
Thread.sleep(CONSUMER_FAIL_RETRY_WAIT_INTERVAL);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep Failure");
|
||||
// logger.log(Level.ALL, "Sleep Failure");
|
||||
}
|
||||
}
|
||||
index = 0;
|
||||
|
|
@ -91,7 +91,7 @@ public class BufferGroup {
|
|||
try {
|
||||
Thread.sleep(MAX_WAIT_TIME);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep Failure");
|
||||
// logger.log(Level.ALL, "Sleep Failure");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -15,6 +15,7 @@ public class Span {
|
|||
private char spanType;
|
||||
private boolean isReceiver = false;
|
||||
private String businessKey;
|
||||
private String processNo;
|
||||
|
||||
public Span(String traceId) {
|
||||
this.traceId = traceId;
|
||||
|
|
@ -80,8 +81,6 @@ public class Span {
|
|||
this.traceId = traceId;
|
||||
}
|
||||
|
||||
private String processNo;
|
||||
|
||||
public String getAddress() {
|
||||
return address;
|
||||
}
|
||||
|
|
@ -91,14 +90,14 @@ public class Span {
|
|||
}
|
||||
|
||||
public byte getStatusCode() {
|
||||
return statusCode;
|
||||
}
|
||||
return statusCode;
|
||||
}
|
||||
|
||||
public void setStatusCode(byte statusCode) {
|
||||
this.statusCode = statusCode;
|
||||
}
|
||||
public void setStatusCode(byte statusCode) {
|
||||
this.statusCode = statusCode;
|
||||
}
|
||||
|
||||
public String getExceptionStack() {
|
||||
public String getExceptionStack() {
|
||||
return exceptionStack;
|
||||
}
|
||||
|
||||
|
|
@ -123,26 +122,65 @@ public class Span {
|
|||
}
|
||||
|
||||
public String getBusinessKey() {
|
||||
return businessKey;
|
||||
}
|
||||
return businessKey;
|
||||
}
|
||||
|
||||
public void setBusinessKey(String businessKey) {
|
||||
this.businessKey = businessKey;
|
||||
}
|
||||
public void setBusinessKey(String businessKey) {
|
||||
this.businessKey = businessKey;
|
||||
}
|
||||
|
||||
@Override
|
||||
@Override
|
||||
public String toString() {
|
||||
StringBuffer stringBuffer = new StringBuffer(traceId + "-" + parentLevel + "-"
|
||||
+ levelId + "-" + viewPointId + "-" + startDate + "-" + cost +
|
||||
"-" + address +
|
||||
"-" + statusCode +
|
||||
"-" + processNo + "-" + spanType + "-" + isReceiver + "-" + businessKey);
|
||||
StringBuilder toStringValue = new StringBuilder();
|
||||
toStringValue.append(traceId + "-");
|
||||
|
||||
|
||||
if (!StringUtil.isEmpty(exceptionStack)) {
|
||||
stringBuffer.append("-" + exceptionStack);
|
||||
if (!StringUtil.isEmpty(parentLevel)) {
|
||||
toStringValue.append(parentLevel + "-");
|
||||
} else {
|
||||
toStringValue.append(" -");
|
||||
}
|
||||
|
||||
return stringBuffer.toString();
|
||||
toStringValue.append(levelId + "-");
|
||||
|
||||
if (!StringUtil.isEmpty(viewPointId)) {
|
||||
toStringValue.append(viewPointId + "-");
|
||||
} else {
|
||||
toStringValue.append(" -");
|
||||
}
|
||||
|
||||
toStringValue.append(startDate + "-");
|
||||
toStringValue.append(cost + "-");
|
||||
|
||||
if (!StringUtil.isEmpty(address)) {
|
||||
toStringValue.append(address + "-");
|
||||
} else {
|
||||
toStringValue.append(" -");
|
||||
}
|
||||
|
||||
toStringValue.append(statusCode + "-");
|
||||
|
||||
if (!StringUtil.isEmpty(exceptionStack)) {
|
||||
toStringValue.append(exceptionStack + "-");
|
||||
} else {
|
||||
toStringValue.append(" -");
|
||||
}
|
||||
|
||||
toStringValue.append(spanType + "-");
|
||||
toStringValue.append(isReceiver + "-");
|
||||
|
||||
|
||||
if (!StringUtil.isEmpty(businessKey)) {
|
||||
toStringValue.append(businessKey + "-");
|
||||
} else {
|
||||
toStringValue.append(" -");
|
||||
}
|
||||
|
||||
if (!StringUtil.isEmpty(processNo)) {
|
||||
toStringValue.append(processNo);
|
||||
} else {
|
||||
toStringValue.append(" -");
|
||||
}
|
||||
|
||||
return toStringValue.toString();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -17,12 +17,11 @@ import static com.ai.cloud.skywalking.conf.Config.SenderChecker.CHECK_POLLING_TI
|
|||
|
||||
public class DataSenderFactory {
|
||||
|
||||
private static Logger logger = Logger.getLogger(DataSenderFactory.getSender().toString());
|
||||
//private static Logger logger = Logger.getLogger(DataSenderFactory.getSender().toString());
|
||||
|
||||
private static Set<SocketAddress> socketAddresses = new HashSet<SocketAddress>();
|
||||
private static Set<SocketAddress> unUsedSocketAddresses = new HashSet<SocketAddress>();
|
||||
private static List<DataSender> availableSenders = new ArrayList<DataSender>();
|
||||
private static DataSenderChecker dataSenderChecker;
|
||||
private static Object lock = new Object();
|
||||
|
||||
static {
|
||||
|
|
@ -38,12 +37,11 @@ public class DataSenderFactory {
|
|||
socketAddresses.add(new InetSocketAddress(server[0], Integer.valueOf(server[1])));
|
||||
}
|
||||
} catch (Exception e) {
|
||||
logger.log(Level.ALL, "Collection service configuration error.");
|
||||
// logger.log(Level.ALL, "Collection service configuration error.");
|
||||
System.exit(-1);
|
||||
}
|
||||
|
||||
dataSenderChecker = new DataSenderChecker();
|
||||
dataSenderChecker.start();
|
||||
new DataSenderChecker().start();
|
||||
}
|
||||
|
||||
public static DataSender getSender() {
|
||||
|
|
@ -51,7 +49,7 @@ public class DataSenderFactory {
|
|||
try {
|
||||
Thread.sleep(RETRY_GET_SENDER_WAIT_INTERVAL);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep failure");
|
||||
// logger.log(Level.ALL, "Sleep failure");
|
||||
}
|
||||
}
|
||||
return availableSenders.get(ThreadLocalRandom.current().nextInt(0, availableSenders.size()));
|
||||
|
|
@ -63,7 +61,7 @@ public class DataSenderFactory {
|
|||
|
||||
public DataSenderChecker() {
|
||||
if (CONNECT_PERCENT <= 0 || CONNECT_PERCENT > 100) {
|
||||
logger.log(Level.ALL, "CONNECT_PERCENT must between 1 and 100");
|
||||
// logger.log(Level.ALL, "CONNECT_PERCENT must between 1 and 100");
|
||||
System.exit(-1);
|
||||
}
|
||||
availableSize = (int) Math.ceil(socketAddresses.size() * ((1.0 * CONNECT_PERCENT / 100) % 100));
|
||||
|
|
@ -112,7 +110,7 @@ public class DataSenderFactory {
|
|||
try {
|
||||
Thread.sleep(CHECK_POLLING_TIME);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep Failure");
|
||||
//logger.log(Level.ALL, "Sleep Failure");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -37,15 +37,15 @@ public final class ContextGenerator {
|
|||
spanData.setParentLevel(context.getParentLevel());
|
||||
}
|
||||
initNewSpanData(spanData, id);
|
||||
|
||||
|
||||
return spanData;
|
||||
}
|
||||
|
||||
private static void initNewSpanData(Span spanData,Identification id){
|
||||
spanData.setSpanType(id.getSpanType());
|
||||
|
||||
private static void initNewSpanData(Span spanData, Identification id) {
|
||||
spanData.setSpanType(id.getSpanType());
|
||||
spanData.setViewPointId(id.getViewPoint());
|
||||
spanData.setBusinessKey(id.getBusinessKey());
|
||||
// 设置基本信息
|
||||
// 设置基本信息
|
||||
spanData.setStartDate(System.currentTimeMillis());
|
||||
spanData.setProcessNo(BuriedPointMachineUtil.getProcessNo());
|
||||
spanData.setAddress(BuriedPointMachineUtil.getHostName() + "/" + BuriedPointMachineUtil.getHostIp());
|
||||
|
|
@ -54,7 +54,7 @@ public final class ContextGenerator {
|
|||
private static Span getSpanFromThreadLocal() {
|
||||
Span span;
|
||||
// 1.获取Context,从ThreadLocal栈中获取中
|
||||
final Span parentSpan = Context.getLastSpan();
|
||||
final Span parentSpan = Context.getLastSpan();
|
||||
// 2 校验Context,Context是否存在
|
||||
if (parentSpan == null) {
|
||||
// 不存在,新创建一个Context
|
||||
|
|
@ -62,7 +62,12 @@ public final class ContextGenerator {
|
|||
} else {
|
||||
// 根据ParentContextData的TraceId和RPCID
|
||||
span = new Span(parentSpan.getTraceId());
|
||||
span.setParentLevel(parentSpan.getParentLevel() + "." + parentSpan.getLevelId());
|
||||
if (!StringUtil.isEmpty(parentSpan.getParentLevel())) {
|
||||
span.setParentLevel(parentSpan.getParentLevel() + "." + parentSpan.getLevelId());
|
||||
} else {
|
||||
span.setParentLevel(String.valueOf(parentSpan.getLevelId()));
|
||||
}
|
||||
|
||||
}
|
||||
return span;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -18,6 +18,10 @@ public class BuriedPointEntry {
|
|||
private String processNo;
|
||||
|
||||
|
||||
private BuriedPointEntry(){
|
||||
|
||||
}
|
||||
|
||||
public String getTraceId() {
|
||||
return traceId;
|
||||
}
|
||||
|
|
@ -78,7 +82,7 @@ public class BuriedPointEntry {
|
|||
result.levelId = Integer.valueOf(fieldValues[2]);
|
||||
result.viewPointId = fieldValues[3];
|
||||
result.startDate = new Date(Long.valueOf(fieldValues[4]));
|
||||
result.cost = Long.getLong(fieldValues[5]);
|
||||
result.cost = Long.parseLong(fieldValues[5]);
|
||||
result.address = fieldValues[6];
|
||||
result.exceptionStack = fieldValues[7];
|
||||
result.spanType = fieldValues[8].charAt(0);
|
||||
|
|
@ -87,4 +91,8 @@ public class BuriedPointEntry {
|
|||
result.processNo = fieldValues[11];
|
||||
return result;
|
||||
}
|
||||
|
||||
public static void main(String[] args){
|
||||
BuriedPointEntry.convert("4d6d38a76c21436e998a81fe798d4ced-.0.0.0-0-com.ai.cloud.skywalking.plugin.spring.common.CallChainE.doBusiness()-1447671990694-0-astraea-PC/192.168.1.108-0- -M-false- -8080;4d6d38a76c21436e998a81fe798d4ced-.0.0-0-com.ai.cloud.skywalking.plugin.spring.common.CallChainC.doBusiness()-1447671990694-11-astraea-PC/192.168.1.108-0- -M-false- -8080");
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@ import com.ai.cloud.skywalking.reciever.hbase.HBaseOperator;
|
|||
import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry;
|
||||
import org.apache.commons.io.FileUtils;
|
||||
import org.apache.commons.io.comparator.NameFileComparator;
|
||||
import org.apache.commons.lang.StringUtils;
|
||||
import org.apache.logging.log4j.LogManager;
|
||||
import org.apache.logging.log4j.Logger;
|
||||
|
||||
|
|
@ -58,12 +59,6 @@ public class PersistenceThread extends Thread {
|
|||
continue;
|
||||
}
|
||||
|
||||
buriedPointData = data.toString().split(";");
|
||||
for (String buriedPoint : buriedPointData) {
|
||||
BuriedPointEntry entry = BuriedPointEntry.convert(buriedPoint);
|
||||
HBaseOperator.insert(entry.getTraceId(), entry.getParentLevel() + "." + entry.getLevelId(), buriedPoint);
|
||||
}
|
||||
|
||||
if ("EOF".equals(data.toString())) {
|
||||
bufferedReader.close();
|
||||
logger.info("Data in file[{}] has been successfully processed", file1.getName());
|
||||
|
|
@ -76,6 +71,17 @@ public class PersistenceThread extends Thread {
|
|||
bool = false;
|
||||
break;
|
||||
}
|
||||
|
||||
buriedPointData = data.toString().split(";");
|
||||
for (String buriedPoint : buriedPointData) {
|
||||
BuriedPointEntry entry = BuriedPointEntry.convert(buriedPoint);
|
||||
if (StringUtils.isEmpty(entry.getParentLevel().trim())) {
|
||||
HBaseOperator.insert(entry.getTraceId(), String.valueOf(entry.getLevelId()), buriedPoint);
|
||||
} else {
|
||||
HBaseOperator.insert(entry.getTraceId(), entry.getParentLevel() + "." + entry.getLevelId(), buriedPoint);
|
||||
}
|
||||
}
|
||||
|
||||
data.delete(0, data.length());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue