diff --git a/src/main/java/com/profiler/context/Annotation.java b/src/main/java/com/profiler/context/Annotation.java index 08dc76d03..8231dec08 100644 --- a/src/main/java/com/profiler/context/Annotation.java +++ b/src/main/java/com/profiler/context/Annotation.java @@ -2,31 +2,36 @@ package com.profiler.context; public interface Annotation { + public static final String CLIENT_SEND = "@CLIENT_SEND"; + public static final String CLIENT_RECV = "@CLIENT_RECV"; + public static final String SERVER_SEND = "@SERVER_SEND"; + public static final String SERVER_RECV = "@SERVER_RECV"; + public static class ClientSend implements Annotation { @Override public String toString() { - return "@CLIENT_SEND"; + return CLIENT_SEND; } } public static class ClientRecv implements Annotation { @Override public String toString() { - return "@CLIENT_RECV"; + return CLIENT_RECV; } } public static class ServerSend implements Annotation { @Override public String toString() { - return "@SERVER_SEND"; + return SERVER_SEND; } } public static class ServerRecv implements Annotation { @Override public String toString() { - return "@SERVER_RECV"; + return SERVER_RECV; } } diff --git a/src/main/java/com/profiler/context/DeadlineSpanMap.java b/src/main/java/com/profiler/context/DeadlineSpanMap.java new file mode 100644 index 000000000..8a0b6aaf1 --- /dev/null +++ b/src/main/java/com/profiler/context/DeadlineSpanMap.java @@ -0,0 +1,32 @@ +package com.profiler.context; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +import com.profiler.context.tracer.Tracer; + +public class DeadlineSpanMap { + + private final Map map = new ConcurrentHashMap(); + + private final Tracer tracer; + + public DeadlineSpanMap(Tracer tracer) { + this.tracer = tracer; + } + + public Span update(TraceID traceId, SpanUpdater spanUpdater) { + Span span = map.get(traceId.toString()); + + if (span == null) { + span = new Span(traceId, null, null); + map.put(traceId.toString(), span); + } + + return spanUpdater.updateSpan(span); + } + + public Span remove(TraceID traceId) { + return map.remove(traceId.toString()); + } +} diff --git a/src/main/java/com/profiler/context/Record.java b/src/main/java/com/profiler/context/Record.java index e4dd65208..483ef2e5d 100644 --- a/src/main/java/com/profiler/context/Record.java +++ b/src/main/java/com/profiler/context/Record.java @@ -14,6 +14,14 @@ public class Record { this.duration = duration; } + public TraceID getTraceId() { + return this.traceId; + } + + public Annotation getAnnotation() { + return this.annotation; + } + @Override public String toString() { StringBuilder sb = new StringBuilder(); diff --git a/src/main/java/com/profiler/context/Span.java b/src/main/java/com/profiler/context/Span.java index eb22de774..2b7239f1d 100644 --- a/src/main/java/com/profiler/context/Span.java +++ b/src/main/java/com/profiler/context/Span.java @@ -1,11 +1,11 @@ package com.profiler.context; -import java.util.Comparator; -import java.util.SortedSet; -import java.util.TreeSet; +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; /** - * A span represents one RPC request. A trace is made up of many spans. * * @author netspider * @@ -15,26 +15,20 @@ public class Span { private final TraceID traceID; private final String name; private final EndPoint endPoint; - private final long createTime; - private final SortedSet annotations = new TreeSet(new Comparator() { - @Override - public int compare(Annotation a1, Annotation a2) { - throw new RuntimeException("Comparator not implemented"); - // return (int) (a1.getTimestamp() - a2.getTimestamp()); - } - }); + private final List annotations = new ArrayList(); + private final Set annotationDesc = new HashSet(); public Span(TraceID traceId, String name, EndPoint endPoint) { this.traceID = traceId; this.name = name; this.endPoint = endPoint; - this.createTime = System.nanoTime(); } public boolean addAnnotation(Annotation annotation) { + annotationDesc.add(annotation.toString()); return annotations.add(annotation); } @@ -42,6 +36,10 @@ public class Span { return annotations.size(); } + public boolean isExistsAnnotation(String annotation) { + return annotationDesc.contains(annotation); + } + public String toString() { StringBuilder sb = new StringBuilder(); diff --git a/src/main/java/com/profiler/context/SpanUpdater.java b/src/main/java/com/profiler/context/SpanUpdater.java new file mode 100644 index 000000000..bf81e943b --- /dev/null +++ b/src/main/java/com/profiler/context/SpanUpdater.java @@ -0,0 +1,5 @@ +package com.profiler.context; + +public interface SpanUpdater { + Span updateSpan(Span span); +} diff --git a/src/main/java/com/profiler/context/tracer/DefaultTracer.java b/src/main/java/com/profiler/context/tracer/DefaultTracer.java index 1d4327942..52cc057ab 100644 --- a/src/main/java/com/profiler/context/tracer/DefaultTracer.java +++ b/src/main/java/com/profiler/context/tracer/DefaultTracer.java @@ -1,10 +1,38 @@ package com.profiler.context.tracer; +import com.profiler.context.Annotation; +import com.profiler.context.DeadlineSpanMap; import com.profiler.context.Record; +import com.profiler.context.Span; +import com.profiler.context.SpanUpdater; +import com.profiler.context.TraceID; public class DefaultTracer implements Tracer { + + private final DeadlineSpanMap spanMap = new DeadlineSpanMap(this); + + private void mutate(TraceID traceId, SpanUpdater spanUpdater) { + Span span = spanMap.update(traceId, spanUpdater); + + if (span.isExistsAnnotation(Annotation.CLIENT_RECV) || span.isExistsAnnotation(Annotation.SERVER_SEND)) { + spanMap.remove(traceId); + logSpan(span); + } + } + + private void logSpan(Span span) { + System.out.println("Write span=" + span); + } + @Override - public void record(Record record) { - System.out.println("record=" + record); + public void record(final Record record) { + mutate(record.getTraceId(), new SpanUpdater() { + @Override + public Span updateSpan(Span span) { + span.addAnnotation(record.getAnnotation()); + return span; + } + }); + } } diff --git a/src/test/java/com/profiler/context/SpanTest.java b/src/test/java/com/profiler/context/SpanTest.java index 4e4afefc6..565cfbf14 100644 --- a/src/test/java/com/profiler/context/SpanTest.java +++ b/src/test/java/com/profiler/context/SpanTest.java @@ -8,20 +8,16 @@ import junit.framework.Assert; import org.junit.Test; +import com.profiler.context.tracer.DefaultTracer; import com.profiler.context.tracer.Tracer; public class SpanTest { @Test public void span() { - Trace.addTracer(new Tracer() { - @Override - public void record(Record record) { - System.out.printf("[%s] Record=%s\n", Thread.currentThread().getId(), record); - } - }); + Trace.addTracer(new DefaultTracer()); - int testSize = 3; + int testSize = 1; final CountDownLatch startLatch = new CountDownLatch(1); final CountDownLatch endLatch = new CountDownLatch(testSize); @@ -67,7 +63,7 @@ public class SpanTest { Trace.record("msg:server send"); Trace.record(new Annotation.ServerSend()); - Trace.record("msg:client recg"); + Trace.record("msg:client recv"); Trace.record(new Annotation.ClientRecv()); endLatch.countDown();