diff --git a/src/main/java/com/profiler/Agent.java b/src/main/java/com/profiler/Agent.java index 2c213e8a1..20f209c7f 100644 --- a/src/main/java/com/profiler/Agent.java +++ b/src/main/java/com/profiler/Agent.java @@ -63,7 +63,7 @@ public class Agent { private void initializeTraceContext() { this.traceContext = TraceContext.getTraceContext(); - this.traceContext.setDataSender(this.dataSender); +// this.traceContext.setDataSender(this.dataSender); this.traceContext.setAgentId(this.agentId); this.traceContext.setApplicationId(this.applicationName); diff --git a/src/main/java/com/profiler/context/AsyncTrace.java b/src/main/java/com/profiler/context/AsyncTrace.java index 37c8692d3..93e1d07b8 100644 --- a/src/main/java/com/profiler/context/AsyncTrace.java +++ b/src/main/java/com/profiler/context/AsyncTrace.java @@ -27,15 +27,16 @@ public class AsyncTrace { private int asyncId = NON_REGIST; private SubSpan subSpan; - private DataSender dataSender; + + private Storage storage; private TimerTask timeoutTask; public AsyncTrace(SubSpan subspan) { this.subSpan = subspan; } - public void setDataSender(DataSender dataSender) { - this.dataSender = dataSender; + public void setStorage(Storage storage) { + this.storage = storage; } public void setTimeoutTask(TimerTask timeoutTask) { @@ -77,6 +78,11 @@ public class AsyncTrace { public void traceBlockEnd() { logSpan(this.subSpan); +// clearReference(); + } + + private void clearReference() { + // 관련 reference를 null로 하는게 좋지 않을까 함? } public void markAfterTime() { @@ -100,41 +106,28 @@ public class AsyncTrace { } public void recordRpcName(final ServiceType serviceType, final String service, final String rpc) { - try { - this.subSpan.setServiceType(serviceType); - this.subSpan.setServiceName(service); - this.subSpan.setRpc(rpc); - } catch (Exception e) { - logger.log(Level.SEVERE, e.getMessage(), e); - } + this.subSpan.setServiceType(serviceType); + this.subSpan.setServiceName(service); + this.subSpan.setRpc(rpc); } // TODO: final String... endPoint로 받으면 합치는데 비용이 들어가 그냥 한번에 받는게 나을것 같음. public void recordEndPoint(final String endPoint) { - try { - this.subSpan.setEndPoint(endPoint); - } catch (Exception e) { - logger.log(Level.SEVERE, e.getMessage(), e); - } + this.subSpan.setEndPoint(endPoint); } private void annotate(final String key) { + this.subSpan.addAnnotation(new HippoAnnotation(key)); - try { - this.subSpan.addAnnotation(new HippoAnnotation(key)); - } catch (Exception e) { - logger.log(Level.SEVERE, e.getMessage(), e); - } } - void logSpan(SubSpan span) { + void logSpan(SubSpan subSpan) { try { if (logger.isLoggable(Level.INFO)) { Thread thread = Thread.currentThread(); - logger.info("[WRITE SubSPAN]" + span + " CurrentThreadID=" + thread.getId() + ",\n\t CurrentThreadName=" + thread.getName()); + logger.info("[WRITE SubSPAN]" + subSpan + " CurrentThreadID=" + thread.getId() + ",\n\t CurrentThreadName=" + thread.getName()); } - - this.dataSender.send(span.toThrift()); + this.storage.store(subSpan); } catch (Exception e) { logger.log(Level.SEVERE, e.getMessage(), e); } diff --git a/src/main/java/com/profiler/context/GlobalCallTrace.java b/src/main/java/com/profiler/context/GlobalCallTrace.java index eff6395a3..51405ecf8 100644 --- a/src/main/java/com/profiler/context/GlobalCallTrace.java +++ b/src/main/java/com/profiler/context/GlobalCallTrace.java @@ -20,13 +20,11 @@ public class GlobalCallTrace { private ConcurrentMap trace = new ConcurrentHashMap(32); private AtomicInteger idGenerator = new AtomicInteger(0); - private DataSender dataSender; private Timer timer = new Timer("GlobalCallTrace-Timer-" + timerId.getAndIncrement(), true); public int registerTraceObject(AsyncTrace asyncTrace) { // TODO 연관관계가 전달부분이 영 별로임. - asyncTrace.setDataSender(this.dataSender); TimeoutTask timeoutTask = new TimeoutTask(trace, asyncTrace.getAsyncId()); asyncTrace.setTimeoutTask(timeoutTask); @@ -59,9 +57,6 @@ public class GlobalCallTrace { return asyncTrace; } - public void setDataSender(DataSender dataSender) { - this.dataSender = dataSender; - } private final class TimeoutTask extends TimerTask { private ConcurrentMap trace; diff --git a/src/main/java/com/profiler/context/Span.java b/src/main/java/com/profiler/context/Span.java index bac2129fe..fcdc06849 100644 --- a/src/main/java/com/profiler/context/Span.java +++ b/src/main/java/com/profiler/context/Span.java @@ -146,15 +146,10 @@ public class Span implements Thriftable { List subSpanList = this.getSubSpanList(); if (subSpanList != null && subSpanList.size() != 0) { - SubSpan first = null; + List tSubSpanList = new ArrayList(subSpanList.size()); for (SubSpan subSpan : subSpanList) { com.profiler.common.dto.thrift.SubSpan tSubSpan = subSpan.toThrift(true); - if (first == null) { - // 첫번째 subSpan에는 sequence를 마크한다. - tSubSpan.setSequence(subSpan.getSequence()); - first = subSpan; - } tSubSpanList.add(tSubSpan); } span.setSubSpanList(tSubSpanList); diff --git a/src/main/java/com/profiler/context/Storage.java b/src/main/java/com/profiler/context/Storage.java index ffb13df9a..aa5f8d759 100644 --- a/src/main/java/com/profiler/context/Storage.java +++ b/src/main/java/com/profiler/context/Storage.java @@ -11,7 +11,17 @@ public interface Storage { DataSender getDataSender(); + /** + * store(SubSpan subSpan)와 store(Span span)간 동기화가 구현되어 있어야 한다. + * + * @param subSpan + */ void store(SubSpan subSpan); + /** + * store(SubSpan subSpan)와 store(Span span)간 동기화가 구현되어 있어야 한다. + * + * @param span + */ void store(Span span); } diff --git a/src/main/java/com/profiler/context/SubSpan.java b/src/main/java/com/profiler/context/SubSpan.java index 386773359..43644002e 100644 --- a/src/main/java/com/profiler/context/SubSpan.java +++ b/src/main/java/com/profiler/context/SubSpan.java @@ -131,37 +131,38 @@ public class SubSpan implements Thriftable { } public com.profiler.common.dto.thrift.SubSpan toThrift(boolean child) { - com.profiler.common.dto.thrift.SubSpan span = new com.profiler.common.dto.thrift.SubSpan(); + com.profiler.common.dto.thrift.SubSpan subSpan = new com.profiler.common.dto.thrift.SubSpan(); long parentSpanStartTime = parentSpan.getStartTime(); - span.setStartElapsed((int) (startTime - parentSpanStartTime)); - span.setEndElapsed((int) (endTime - startTime)); + subSpan.setStartElapsed((int) (startTime - parentSpanStartTime)); + subSpan.setEndElapsed((int) (endTime - startTime)); + subSpan.setSequence(sequence); // 다른 span의 sub로 들어가지 않을 경우 if (!child) { - span.setAgentId(Agent.getInstance().getAgentId()); + subSpan.setAgentId(Agent.getInstance().getAgentId()); TraceID parentSpanTraceID = parentSpan.getTraceID(); - span.setMostTraceId(parentSpanTraceID.getId().getMostSignificantBits()); - span.setLeastTraceId(parentSpanTraceID.getId().getLeastSignificantBits()); - span.setSpanId(parentSpanTraceID.getSpanId()); - span.setSequence(sequence); + subSpan.setMostTraceId(parentSpanTraceID.getId().getMostSignificantBits()); + subSpan.setLeastTraceId(parentSpanTraceID.getId().getLeastSignificantBits()); + subSpan.setSpanId(parentSpanTraceID.getSpanId()); } - span.setRpc(rpc); - span.setServiceName(serviceName); - span.setServiceType(serviceType.getCode()); + + subSpan.setRpc(rpc); + subSpan.setServiceName(serviceName); + subSpan.setServiceType(serviceType.getCode()); - span.setEndPoint(endPoint); + subSpan.setEndPoint(endPoint); // 여기서 데이터 인코딩을 하자. List annotationList = new ArrayList(annotations.size()); for (HippoAnnotation a : annotations) { annotationList.add(a.toThrift()); } - span.setAnnotations(annotationList); + subSpan.setAnnotations(annotationList); - return span; + return subSpan; } } diff --git a/src/main/java/com/profiler/context/SubSpanList.java b/src/main/java/com/profiler/context/SubSpanList.java index 80f27c8d0..18ea25bdc 100644 --- a/src/main/java/com/profiler/context/SubSpanList.java +++ b/src/main/java/com/profiler/context/SubSpanList.java @@ -29,7 +29,6 @@ public class SubSpanList implements Thriftable { tSubSpanList.setMostTraceId(id.getMostSignificantBits()); tSubSpanList.setLeastTraceId(id.getLeastSignificantBits()); tSubSpanList.setSpanId(parentSpan.getTraceID().getSpanId()); - tSubSpanList.setStartSequence(first.getSequence()); List tSubSpan = createSubSpan(subSpanList); @@ -49,6 +48,8 @@ public class SubSpanList implements Thriftable { tSubSpan.setStartElapsed((int) (subSpan.getStartTime() - parentSpanStartTime)); tSubSpan.setEndElapsed((int) (subSpan.getEndTime() - subSpan.getStartTime())); + tSubSpan.setSequence(subSpan.getSequence()); + tSubSpan.setRpc(subSpan.getRpc()); tSubSpan.setServiceName(subSpan.getServiceName()); tSubSpan.setServiceType(subSpan.getServiceType().getCode()); diff --git a/src/main/java/com/profiler/context/TimeBaseStorage.java b/src/main/java/com/profiler/context/TimeBaseStorage.java index b86781034..11fcba967 100644 --- a/src/main/java/com/profiler/context/TimeBaseStorage.java +++ b/src/main/java/com/profiler/context/TimeBaseStorage.java @@ -10,13 +10,16 @@ import java.util.logging.Logger; * */ public class TimeBaseStorage implements Storage { + + private static final int RESERVE_BUFFER_SIZE = 2; + private boolean discard = true; private boolean limit; private long limitTime = 1000; private int bufferSize = 20; - private List storage = new ArrayList(bufferSize); + private List storage = new ArrayList(bufferSize + RESERVE_BUFFER_SIZE); private DataSender dataSender; public TimeBaseStorage() { @@ -46,29 +49,42 @@ public class TimeBaseStorage implements Storage { @Override public void store(SubSpan subSpan) { - addSubSpan(subSpan); // flush유무 확인 if (!limit) { // 절대 시간만 체크한다. 1초 이내 라서 절대 데이터를 flush하지 않는다. + synchronized (this) { + addSubSpan(subSpan); + } limit = checkLimit(subSpan); } else { // 1초가 지났다면. // 데이터가 flushCount이상일 경우 먼저 flush한다. - if (storage.size() >= bufferSize) { - SubSpanList subSpanList = new SubSpanList(storage); - storage = new ArrayList(bufferSize); - dataSender.send(subSpanList); + List flushData = null; + synchronized (this) { + if (!addSubSpan(subSpan)) { + dataSender.send(subSpan); + return; + } + if (storage.size() >= bufferSize) { + flushData = storage; + storage = new ArrayList(bufferSize + RESERVE_BUFFER_SIZE); + } + } + if (flushData != null) { + dataSender.send(new SubSpanList(flushData)); } } } - private void addSubSpan(SubSpan subSpan) { + private boolean addSubSpan(SubSpan subSpan) { if (storage == null) { Logger logger = Logger.getLogger(this.getClass().getName()); - logger.warning("storage is null."); - return; + logger.fine("storage is null. direct send"); + // 이미 span이 와서 flush된 상황임. + return false; } storage.add(subSpan); + return true; } private boolean checkLimit(SubSpan subSpan) { @@ -89,7 +105,9 @@ public class TimeBaseStorage implements Storage { limit = checkLimit(span); if (!limit) { // 제한시간내 빨리 끝난 경우는 subspan을 버린다. - this.storage = null; + synchronized (this) { + this.storage = null; + } dataSender.send(span); } else { @@ -102,11 +120,14 @@ public class TimeBaseStorage implements Storage { } private void flushAll(Span span) { - List subSpanList = storage; + List subSpanList; + synchronized (this) { + subSpanList = storage; + this.storage = null; + } if (subSpanList != null && subSpanList.size() != 0) { span.setSubSpanList(subSpanList); } - this.storage = null; dataSender.send(span); } diff --git a/src/main/java/com/profiler/context/Trace.java b/src/main/java/com/profiler/context/Trace.java index 625254a98..0923f6880 100644 --- a/src/main/java/com/profiler/context/Trace.java +++ b/src/main/java/com/profiler/context/Trace.java @@ -64,13 +64,17 @@ public final class Trace { return storage.getDataSender(); } + public short getSequence() { + return sequence++; + } public AsyncTrace createAsyncTrace() { // 경우에 따라 별도 timeout 처리가 있어야 될수도 있음. SubSpan subSpan = new SubSpan(callStack.getSpan()); - subSpan.setSequence(sequence++); + subSpan.setSequence(getSequence()); AsyncTrace asyncTrace = new AsyncTrace(subSpan); - asyncTrace.setDataSender(this.getDataSender()); +// asyncTrace.setDataSender(this.getDataSender()); + asyncTrace.setStorage(this.storage); return asyncTrace; } @@ -78,7 +82,7 @@ public final class Trace { SubSpan subSpan = new SubSpan(callStack.getSpan()); SubStackFrame stackFrame = new SubStackFrame(subSpan); stackFrame.setStackFrameId(stackId); - stackFrame.setSequence(sequence++); + stackFrame.setSequence(getSequence()); return stackFrame; } diff --git a/src/main/java/com/profiler/context/TraceContext.java b/src/main/java/com/profiler/context/TraceContext.java index f9df33998..9bc962ec3 100644 --- a/src/main/java/com/profiler/context/TraceContext.java +++ b/src/main/java/com/profiler/context/TraceContext.java @@ -29,10 +29,6 @@ public class TraceContext { // internal stacktrace 추적때 필요한 unique 아이디, activethreadcount의 slow 타임 계산의 위해서도 필요할듯 함. private final AtomicInteger transactionId = new AtomicInteger(0); - private static final DataSender DEFAULT_DATA_SENDER = new LoggingDataSender(); - - private DataSender dataSender = DEFAULT_DATA_SENDER; - private GlobalCallTrace globalCallTrace = new GlobalCallTrace(); private String agentId; @@ -75,11 +71,6 @@ public class TraceContext { return activeThreadCounter; } - public void setDataSender(DataSender dataSender) { - this.dataSender = dataSender; - this.globalCallTrace.setDataSender(dataSender); - } - public void setAgentId(String agentId) { this.agentId = agentId; }