Add exportable to Span and use it to determine if data are exported

If the flag is set then annotations are collected and the data are exported
in zipkin or stream. If not you still get the correlation ids, so a purely
log-oriented solution will always have useful data on all requests.
This commit is contained in:
Dave Syer
2015-12-17 16:35:06 +00:00
parent b48f00f677
commit c11e530560
20 changed files with 237 additions and 203 deletions

View File

@@ -30,22 +30,27 @@ import lombok.Singular;
* @author Spencer Gibb
*/
@Data
@Builder
@Builder(toBuilder=true)
public class MilliSpan implements Span {
private final long begin;
private long end = 0;
private final String name;
private final String traceId;
@Singular
private List<String> parents;
private List<String> parents = new ArrayList<>();
private final String spanId;
private boolean remote = false;
private boolean exportable = true;
private final Map<String, String> annotations = new LinkedHashMap<>();
private final String processId;
@Singular
private final List<TimelineAnnotation> timelineAnnotations = new ArrayList<>();
public MilliSpan(long begin, long end, String name, String traceId, List<String> parents, String spanId, boolean remote, String processId) {
public static MilliSpan.MilliSpanBuilder builder() {
return new MilliSpan().toBuilder();
}
public MilliSpan(long begin, long end, String name, String traceId, List<String> parents, String spanId, boolean remote, boolean exportable, String processId) {
this.begin = begin<=0 ? System.currentTimeMillis() : begin;
this.end = end;
this.name = name;
@@ -53,16 +58,15 @@ public class MilliSpan implements Span {
this.parents = parents;
this.spanId = spanId;
this.remote = remote;
this.exportable = exportable;
this.processId = processId;
}
//for serialization
@SuppressWarnings("unused")
private MilliSpan() {
this.begin = 0;
this.name = null;
this.traceId = null;
this.parents = null;
this.spanId = null;
this.processId = null;
}

View File

@@ -96,6 +96,12 @@ public interface Span {
*/
boolean isRunning();
/**
* Is the span eligible for export? If not then we may not need accumulate annotations
* (for instance).
*/
boolean isExportable();
/**
* Add a data annotation associated with this span
*/

View File

@@ -26,6 +26,8 @@ import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
/**
* Utility for manipulating message headers related to span data.
*
* @author Dave Syer
*
*/
@@ -40,24 +42,28 @@ public class SpanMessageHeaders {
return message;
}
addAnnotations(message, span);
Map<String, String> headers = new HashMap<>();
addHeader(headers, Trace.TRACE_ID_NAME, span.getTraceId());
addHeader(headers, Trace.SPAN_ID_NAME, span.getSpanId());
addHeader(headers, Trace.PARENT_ID_NAME, getFirst(span.getParents()));
addHeader(headers, Trace.SPAN_NAME_NAME, span.getName());
addHeader(headers, Trace.PROCESS_ID_NAME, span.getProcessId());
if (span.isExportable()) {
addAnnotations(message, span);
addHeader(headers, Trace.PARENT_ID_NAME, getFirst(span.getParents()));
addHeader(headers, Trace.SPAN_NAME_NAME, span.getName());
addHeader(headers, Trace.PROCESS_ID_NAME, span.getProcessId());
} else {
addHeader(headers, Trace.NOT_SAMPLED_NAME, "");
}
return MessageBuilder.fromMessage(message).copyHeaders(headers).build();
}
public static void addAnnotations(Message<?> message, Span span) {
for ( Map.Entry<String, Object> entry : message.getHeaders().entrySet()) {
if (!Trace.HEADERS.contains(entry.getKey())) { //filter out trace headers
for (Map.Entry<String, Object> entry : message.getHeaders().entrySet()) {
if (!Trace.HEADERS.contains(entry.getKey())) { // filter out trace headers
String key = "/messaging/headers/" + entry.getKey().toLowerCase();
String value = null;
if (entry.getValue() != null) {
value = entry.getValue().toString(); //TODO: better way to serialize?
value = entry.getValue().toString(); // TODO: better way to serialize?
}
span.addAnnotation(key, value);
}
@@ -69,15 +75,17 @@ public class SpanMessageHeaders {
payload.getClass().getCanonicalName());
if (payload instanceof String) {
span.addAnnotation("/messaging/payload/size",
String.valueOf(((String)payload).length()));
} else if (payload instanceof byte[]) {
String.valueOf(((String) payload).length()));
}
else if (payload instanceof byte[]) {
span.addAnnotation("/messaging/payload/size",
String.valueOf(((byte[])payload).length));
String.valueOf(((byte[]) payload).length));
}
}
}
private static void addHeader(Map<String, String> headers, String name, String value) {
private static void addHeader(Map<String, String> headers, String name,
String value) {
if (value != null) {
headers.put(name, value);
}

View File

@@ -16,8 +16,6 @@
package org.springframework.cloud.sleuth.instrument.integration;
import static org.springframework.util.StringUtils.hasText;
import org.springframework.cloud.sleuth.MilliSpan;
import org.springframework.cloud.sleuth.MilliSpan.MilliSpanBuilder;
import org.springframework.cloud.sleuth.Trace;
@@ -27,6 +25,7 @@ import org.springframework.integration.context.IntegrationObjectSupport;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.support.ChannelInterceptorAdapter;
import org.springframework.util.StringUtils;
/**
* @author Dave Syer
@@ -51,8 +50,7 @@ public class TraceChannelInterceptor extends ChannelInterceptorAdapter {
@Override
public Message<?> preSend(Message<?> message, MessageChannel channel) {
if (this.traceManager.isTracing()
|| message.getHeaders().containsKey(Trace.NOT_SAMPLED_NAME)) {
if (this.traceManager.isTracing()) {
return SpanMessageHeaders.addSpanHeaders(message,
this.traceManager.getCurrentSpan());
}
@@ -60,10 +58,13 @@ public class TraceChannelInterceptor extends ChannelInterceptorAdapter {
String traceId = getHeader(message, Trace.TRACE_ID_NAME);
String name = "message/" + getChannelName(channel);
Trace trace;
if (hasText(spanId) && hasText(traceId)) {
if (StringUtils.hasText(traceId)) {
MilliSpanBuilder span = MilliSpan.builder().traceId(traceId).spanId(spanId);
String parentId = getHeader(message, Trace.PARENT_ID_NAME);
if (message.getHeaders().containsKey(Trace.NOT_SAMPLED_NAME)) {
span.exportable(false);
}
String processId = getHeader(message, Trace.PROCESS_ID_NAME);
String spanName = getHeader(message, Trace.SPAN_NAME_NAME);
if (spanName != null) {
@@ -77,7 +78,6 @@ public class TraceChannelInterceptor extends ChannelInterceptorAdapter {
}
span.remote(true);
// TODO: traceManager description?
trace = this.traceManager.startSpan(name, span.build());
}
else {

View File

@@ -101,50 +101,52 @@ public class TraceFilter extends OncePerRequestFilter
else if (skip) {
addToResponseIfNotPresent(response, Trace.NOT_SAMPLED_NAME, "");
}
else {
String spanId = getHeader(request, response, Trace.SPAN_ID_NAME);
String traceId = getHeader(request, response, Trace.TRACE_ID_NAME);
String name = "http" + uri;
if (hasText(spanId) && hasText(traceId)) {
MilliSpanBuilder span = MilliSpan.builder().traceId(traceId)
.spanId(spanId);
String parentId = getHeader(request, response, Trace.PARENT_ID_NAME);
String processId = getHeader(request, response, Trace.PROCESS_ID_NAME);
String parentName = getHeader(request, response, Trace.SPAN_NAME_NAME);
if (parentName != null) {
span.name(parentName);
}
if (processId != null) {
span.processId(processId);
}
if (parentId != null) {
span.parent(parentId);
}
span.remote(true);
String spanId = getHeader(request, response, Trace.SPAN_ID_NAME);
String traceId = getHeader(request, response, Trace.TRACE_ID_NAME);
String name = "http" + uri;
if (hasText(traceId)) {
// TODO: trace description?
Span parent = span.build();
trace = this.traceManager.startSpan(name, parent);
publish(new ServerReceivedEvent(this, parent, trace.getSpan()));
request.setAttribute(TRACE_REQUEST_ATTR, trace);
// Send new span id back
addToResponseIfNotPresent(response, Trace.TRACE_ID_NAME,
trace.getSpan().getTraceId());
addToResponseIfNotPresent(response, Trace.SPAN_ID_NAME,
trace.getSpan().getSpanId());
MilliSpanBuilder span = MilliSpan.builder().traceId(traceId).spanId(spanId);
if (skip) {
span.exportable(false);
}
else {
trace = this.traceManager.startSpan(name);
request.setAttribute(TRACE_REQUEST_ATTR, trace);
String parentId = getHeader(request, response, Trace.PARENT_ID_NAME);
String processId = getHeader(request, response, Trace.PROCESS_ID_NAME);
String parentName = getHeader(request, response, Trace.SPAN_NAME_NAME);
if (parentName != null) {
span.name(parentName);
}
if (processId != null) {
span.processId(processId);
}
if (parentId != null) {
span.parent(parentId);
}
span.remote(true);
Span parent = span.build();
trace = this.traceManager.startSpan(name, parent);
publish(new ServerReceivedEvent(this, parent, trace.getSpan()));
request.setAttribute(TRACE_REQUEST_ATTR, trace);
}
else {
trace = this.traceManager.startSpan(name);
request.setAttribute(TRACE_REQUEST_ATTR, trace);
}
// Send new trace id back to the caller
addToResponseIfNotPresent(response, Trace.TRACE_ID_NAME,
trace.getSpan().getTraceId());
addToResponseIfNotPresent(response, Trace.SPAN_ID_NAME,
trace.getSpan().getSpanId());
try {
addRequestAnnotations(request);
filterChain.doFilter(request, response);
}
finally {
if (isAsyncStarted(request) || request.isAsyncStarted()) {

View File

@@ -107,10 +107,9 @@ public class TraceWebAspect {
}
}
@SuppressWarnings("unchecked")
@Around("anyControllerOrRestControllerWithPublicWebAsyncTaskMethod()")
public Object wrapWebAsyncTaskWithCorrelationId(ProceedingJoinPoint pjp) throws Throwable {
final WebAsyncTask webAsyncTask = (WebAsyncTask) pjp.proceed();
final WebAsyncTask<?> webAsyncTask = (WebAsyncTask<?>) pjp.proceed();
if (this.accessor.isTracing()) {
try {
log.debug("Wrapping callable with span ["

View File

@@ -17,7 +17,6 @@
package org.springframework.cloud.sleuth.instrument.web.client;
import static java.util.Collections.singletonList;
import static org.springframework.cloud.sleuth.trace.TraceContextHolder.isTracing;
import java.io.IOException;
import java.lang.reflect.Type;
@@ -75,11 +74,14 @@ public class TraceFeignClientAutoConfiguration {
public Decoder feignDecoder() {
return new ResponseEntityDecoder(new SpringDecoder(this.messageConverters)) {
@Override
public Object decode(Response response, Type type) throws IOException, FeignException {
public Object decode(Response response, Type type)
throws IOException, FeignException {
try {
return super.decode(Response.create(response.status(), response.reason(),
headersWithTraceId(response.headers()), response.body()), type);
} finally {
return super.decode(Response.create(response.status(),
response.reason(), headersWithTraceId(response.headers()),
response.body()), type);
}
finally {
Span span = getCurrentSpan();
if (span != null) {
publish(new ClientReceivedEvent(this, span));
@@ -95,16 +97,21 @@ public class TraceFeignClientAutoConfiguration {
@Override
public void apply(RequestTemplate template) {
Span span = getCurrentSpan();
if (span != null) {
template.header(Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(template, Trace.SPAN_NAME_NAME, span.getName());
setHeader(template, Trace.SPAN_ID_NAME, span.getSpanId());
setHeader(template, Trace.PARENT_ID_NAME, getParentId(span));
setHeader(template, Trace.PROCESS_ID_NAME, span.getProcessId());
publish(new ClientSentEvent(this, span));
} else {
if (span == null) {
setHeader(template, Trace.NOT_SAMPLED_NAME, "");
return;
}
if (span.getSpanId() == null) {
setHeader(template, Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(template, Trace.NOT_SAMPLED_NAME, "");
return;
}
template.header(Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(template, Trace.SPAN_NAME_NAME, span.getName());
setHeader(template, Trace.SPAN_ID_NAME, span.getSpanId());
setHeader(template, Trace.PARENT_ID_NAME, getParentId(span));
setHeader(template, Trace.PROCESS_ID_NAME, span.getProcessId());
publish(new ClientSentEvent(this, span));
}
};
}
@@ -116,30 +123,40 @@ public class TraceFeignClientAutoConfiguration {
}
private String getParentId(Span span) {
return span.getParents() != null && !span.getParents().isEmpty() ? span.getParents().get(0) : null;
return span.getParents() != null && !span.getParents().isEmpty()
? span.getParents().get(0) : null;
}
public void setHeader(RequestTemplate request, String name, String value) {
if (value != null && !request.headers().containsKey(name) && isTracing()) {
if (value != null && !request.headers().containsKey(name)
&& this.accessor.isTracing()) {
request.header(name, value);
}
}
private Map<String, Collection<String>> headersWithTraceId(Map<String, Collection<String>> headers) {
private Map<String, Collection<String>> headersWithTraceId(
Map<String, Collection<String>> headers) {
Map<String, Collection<String>> newHeaders = new HashMap<>();
newHeaders.putAll(headers);
if (getCurrentSpan() == null) {
Span span = getCurrentSpan();
if (span == null) {
setHeader(newHeaders, Trace.NOT_SAMPLED_NAME, "");
return newHeaders;
}
setHeader(newHeaders, Trace.TRACE_ID_NAME, getCurrentSpan().getTraceId());
setHeader(newHeaders, Trace.SPAN_ID_NAME, getCurrentSpan().getSpanId());
setHeader(newHeaders, Trace.PARENT_ID_NAME, getParentId(getCurrentSpan()));
if (span.getSpanId() == null) {
setHeader(newHeaders, Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(newHeaders, Trace.NOT_SAMPLED_NAME, "");
return newHeaders;
}
setHeader(newHeaders, Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(newHeaders, Trace.SPAN_ID_NAME, span.getSpanId());
setHeader(newHeaders, Trace.PARENT_ID_NAME, getParentId(span));
return newHeaders;
}
public void setHeader(Map<String, Collection<String>> headers, String name, String value) {
if (value != null && !headers.containsKey(name) && isTracing()) {
public void setHeader(Map<String, Collection<String>> headers, String name,
String value) {
if (value != null && !headers.containsKey(name) && this.accessor.isTracing()) {
headers.put(name, singletonList(value));
}
}

View File

@@ -15,8 +15,6 @@
*/
package org.springframework.cloud.sleuth.instrument.web.client;
import static org.springframework.cloud.sleuth.trace.TraceContextHolder.isTracing;
import java.io.IOException;
import org.springframework.cloud.sleuth.Span;
@@ -61,16 +59,22 @@ ApplicationEventPublisherAware {
@Override
public ClientHttpResponse intercept(HttpRequest request, byte[] body,
ClientHttpRequestExecution execution) throws IOException {
if (getCurrentSpan() == null) {
Span span = getCurrentSpan();
if (span == null) {
setHeader(request, Trace.NOT_SAMPLED_NAME, "");
return execution.execute(request, body);
}
setHeader(request, Trace.SPAN_ID_NAME, getCurrentSpan().getSpanId());
setHeader(request, Trace.TRACE_ID_NAME, getCurrentSpan().getTraceId());
setHeader(request, Trace.SPAN_NAME_NAME, getCurrentSpan().getName());
setHeader(request, Trace.PARENT_ID_NAME, getParentId(getCurrentSpan()));
setHeader(request, Trace.PROCESS_ID_NAME, getCurrentSpan().getProcessId());
publish(new ClientSentEvent(this, getCurrentSpan()));
if (span.getSpanId()==null) {
setHeader(request, Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(request, Trace.NOT_SAMPLED_NAME, "");
return execution.execute(request, body);
}
setHeader(request, Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(request, Trace.SPAN_ID_NAME, span.getSpanId());
setHeader(request, Trace.SPAN_NAME_NAME, span.getName());
setHeader(request, Trace.PARENT_ID_NAME, getParentId(span));
setHeader(request, Trace.PROCESS_ID_NAME, span.getProcessId());
publish(new ClientSentEvent(this, span));
return new TraceHttpResponse(this, execution.execute(request, body));
}
@@ -93,7 +97,7 @@ ApplicationEventPublisherAware {
}
public void setHeader(HttpRequest request, String name, String value) {
if (value != null && !request.getHeaders().containsKey(name) && isTracing()) {
if (value != null && !request.getHeaders().containsKey(name) && this.accessor.isTracing()) {
request.getHeaders().add(name, value);
}
}

View File

@@ -16,8 +16,6 @@
package org.springframework.cloud.sleuth.instrument.zuul;
import static org.springframework.cloud.sleuth.trace.TraceContextHolder.isTracing;
import java.util.Map;
import org.springframework.cloud.sleuth.Span;
@@ -62,18 +60,24 @@ ApplicationEventPublisherAware {
RequestContext ctx = RequestContext.getCurrentContext();
Map<String, String> response = ctx.getZuulRequestHeaders();
// N.B. this will only work with the simple host filter (not ribbon) unless you set hystrix.execution.isolation.strategy=SEMAPHORE
if (getCurrentSpan() == null) {
Span span = getCurrentSpan();
if (span == null) {
setHeader(response, Trace.NOT_SAMPLED_NAME, "");
return null;
}
if (span.getSpanId()==null) {
setHeader(response, Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(response, Trace.NOT_SAMPLED_NAME, "");
return null;
}
try {
setHeader(response, Trace.SPAN_ID_NAME, getCurrentSpan().getSpanId());
setHeader(response, Trace.TRACE_ID_NAME, getCurrentSpan().getTraceId());
setHeader(response, Trace.SPAN_NAME_NAME, getCurrentSpan().getName());
setHeader(response, Trace.PARENT_ID_NAME, getParentId(getCurrentSpan()));
setHeader(response, Trace.PROCESS_ID_NAME, getCurrentSpan().getProcessId());
setHeader(response, Trace.SPAN_ID_NAME, span.getSpanId());
setHeader(response, Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(response, Trace.SPAN_NAME_NAME, span.getName());
setHeader(response, Trace.PARENT_ID_NAME, getParentId(span));
setHeader(response, Trace.PROCESS_ID_NAME, span.getProcessId());
// TODO: the client sent event should come from the client not the filter!
publish(new ClientSentEvent(this, getCurrentSpan()));
publish(new ClientSentEvent(this, span));
}
catch (Exception ex) {
ReflectionUtils.rethrowRuntimeException(ex);
@@ -91,7 +95,7 @@ ApplicationEventPublisherAware {
}
public void setHeader(Map<String, String> request, String name, String value) {
if (value != null && !request.containsKey(name) && isTracing()) {
if (value != null && !request.containsKey(name) && this.accessor.isTracing()) {
request.put(name, value);
}
}

View File

@@ -16,8 +16,6 @@
package org.springframework.cloud.sleuth.instrument.zuul;
import static org.springframework.cloud.sleuth.trace.TraceContextHolder.isTracing;
import java.io.InputStream;
import java.net.URISyntaxException;
@@ -93,19 +91,25 @@ public class TraceRestClientRibbonCommandFactory extends RestClientRibbonCommand
@Override
protected void customizeRequest(HttpRequest.Builder requestBuilder) {
if (getCurrentSpan() == null) {
Span span = getCurrentSpan();
if (span == null) {
setHeader(requestBuilder, Trace.NOT_SAMPLED_NAME, "");
return;
}
if (span.getSpanId()==null) {
setHeader(requestBuilder, Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(requestBuilder, Trace.NOT_SAMPLED_NAME, "");
return;
}
setHeader(requestBuilder, Trace.SPAN_ID_NAME, getCurrentSpan().getSpanId());
setHeader(requestBuilder, Trace.TRACE_ID_NAME, getCurrentSpan().getTraceId());
setHeader(requestBuilder, Trace.SPAN_NAME_NAME, getCurrentSpan().getName());
setHeader(requestBuilder, Trace.TRACE_ID_NAME, span.getTraceId());
setHeader(requestBuilder, Trace.SPAN_ID_NAME, span.getSpanId());
setHeader(requestBuilder, Trace.SPAN_NAME_NAME, span.getName());
setHeader(requestBuilder, Trace.PARENT_ID_NAME,
getParentId(getCurrentSpan()));
getParentId(span));
setHeader(requestBuilder, Trace.PROCESS_ID_NAME,
getCurrentSpan().getProcessId());
publish(new ClientSentEvent(this, getCurrentSpan()));
span.getProcessId());
publish(new ClientSentEvent(this, span));
}
private void publish(ApplicationEvent event) {
@@ -120,7 +124,7 @@ public class TraceRestClientRibbonCommandFactory extends RestClientRibbonCommand
}
public void setHeader(HttpRequest.Builder builder, String name, String value) {
if (value != null && isTracing()) {
if (value != null && this.accessor.isTracing()) {
builder.header(name, value);
}
}

View File

@@ -19,7 +19,6 @@ package org.springframework.cloud.sleuth.template;
import org.springframework.cloud.sleuth.Trace;
import org.springframework.cloud.sleuth.TraceManager;
import org.springframework.cloud.sleuth.instrument.TraceDelegate;
import org.springframework.cloud.sleuth.trace.TraceContextHolder;
/**
* @author Spencer Gibb
@@ -34,7 +33,7 @@ public class TraceTemplate implements TraceOperations {
@Override
public <T> T trace(final TraceCallback<T> callback) {
if (TraceContextHolder.isTracing()) {
if (this.traceManager.isTracing()) {
DelegateCallback<T> delegate = new DelegateCallback<>(this.traceManager);
Trace trace = delegate.startSpan();
try {

View File

@@ -74,9 +74,16 @@ public class DefaultTraceManager implements TraceManager {
@Override
public <T> Trace startSpan(String name, Sampler<T> s, T info) {
Span span = null;
if (TraceContextHolder.isTracing() || s.next(info)) {
if (isTracing() || s.next(info)) {
span = createChild(getCurrentSpan(), name);
}
else {
// Non-exportable so we keep the trace but not other data
String id = createId();
span = MilliSpan.builder().begin(System.currentTimeMillis()).name(name)
.traceId(id).spanId(id).exportable(false).build();
this.publisher.publishEvent(new SpanAcquiredEvent(this, span));
}
return continueSpan(span);
}
@@ -94,7 +101,7 @@ public class DefaultTraceManager implements TraceManager {
+ ". You have " + "probably forgotten to close or detach " + cur);
}
else {
if (span != NullTrace.INSTANCE) {
if (trace.getSavedTrace() != null) {
TraceContextHolder.setCurrentTrace(trace.getSavedTrace());
}
else {
@@ -119,7 +126,7 @@ public class DefaultTraceManager implements TraceManager {
+ ". You have " + "probably forgotten to close or detach " + cur);
}
else {
if (span != NullTrace.INSTANCE && span != null) {
if (span != null) {
span.stop();
if (savedTrace != null
&& span.getParents().contains(savedTrace.getSpan().getSpanId())) {
@@ -140,9 +147,10 @@ public class DefaultTraceManager implements TraceManager {
}
protected Span createChild(Span parent, String name) {
String id = createId();
if (parent == null) {
MilliSpan span = MilliSpan.builder().begin(System.currentTimeMillis())
.name(name).traceId(createId()).spanId(createId()).build();
.name(name).traceId(id).spanId(id).build();
this.publisher.publishEvent(new SpanAcquiredEvent(this, span));
return span;
}
@@ -153,7 +161,7 @@ public class DefaultTraceManager implements TraceManager {
}
MilliSpan span = MilliSpan.builder().begin(System.currentTimeMillis())
.name(name).traceId(parent.getTraceId()).parent(parent.getSpanId())
.spanId(createId()).processId(parent.getProcessId()).build();
.spanId(id).processId(parent.getProcessId()).build();
this.publisher.publishEvent(new SpanAcquiredEvent(this, parent, span));
return span;
}
@@ -165,11 +173,9 @@ public class DefaultTraceManager implements TraceManager {
@Override
public Trace continueSpan(Span span) {
// Return an empty Trace that does nothing on close
if (span == null) {
return NullTrace.INSTANCE;
if (span != null) {
this.publisher.publishEvent(new SpanContinuedEvent(this, span));
}
this.publisher.publishEvent(new SpanContinuedEvent(this, span));
Trace trace = createTrace(TraceContextHolder.getCurrentTrace(), span);
TraceContextHolder.setCurrentTrace(trace);
return trace;
@@ -192,7 +198,7 @@ public class DefaultTraceManager implements TraceManager {
@Override
public void addAnnotation(String key, String value) {
Span s = getCurrentSpan();
if (s != null) {
if (s != null && s.isExportable()) {
s.addAnnotation(key, value);
}
}
@@ -204,7 +210,7 @@ public class DefaultTraceManager implements TraceManager {
*/
@Override
public <V> Callable<V> wrap(Callable<V> callable) {
if (TraceContextHolder.isTracing()) {
if (isTracing()) {
return new TraceCallable<>(this, callable);
}
return callable;
@@ -217,7 +223,7 @@ public class DefaultTraceManager implements TraceManager {
*/
@Override
public Runnable wrap(Runnable runnable) {
if (TraceContextHolder.isTracing()) {
if (isTracing()) {
return new TraceRunnable(this, runnable);
}
return runnable;

View File

@@ -1,39 +0,0 @@
/*
* Copyright 2013-2015 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.cloud.sleuth.trace;
import org.springframework.cloud.sleuth.Trace;
/**
* @author Spencer Gibb
*/
public final class NullTrace extends Trace {
/**
* Singleton instance representing an empty {@link Trace}.
*/
public static final Trace INSTANCE = new NullTrace();
private NullTrace() {
super(null);
}
@Override
public String toString() {
return "NullTrace";
}
}

View File

@@ -28,33 +28,33 @@ import lombok.extern.apachecommons.CommonsLog;
@CommonsLog
public class TraceContextHolder {
private static final ThreadLocal<Trace> currentSpan = new NamedThreadLocal<>("Trace Context");
private static final ThreadLocal<Trace> currentTrace = new NamedThreadLocal<>("Trace Context");
public static Trace getCurrentTrace() {
return currentSpan.get();
return currentTrace.get();
}
public static Span getCurrentSpan() {
return isTracing() ? currentSpan.get().getSpan() : null;
return isTracing() ? currentTrace.get().getSpan() : null;
}
public static void setCurrentTrace(Trace span) {
public static void setCurrentTrace(Trace trace) {
// backwards compatibility
if (span == null) {
currentSpan.remove();
if (trace == null) {
currentTrace.remove();
return;
}
if (log.isTraceEnabled()) {
log.trace("Setting current span " + span);
log.trace("Setting current trace " + trace);
}
currentSpan.set(span);
currentTrace.set(trace);
}
public static void removeCurrentTrace() {
currentSpan.remove();
currentTrace.remove();
}
public static boolean isTracing() {
return currentSpan.get() != null;
return currentTrace.get() != null;
}
}

View File

@@ -28,7 +28,7 @@ public class MilliSpanTests {
@Test(expected = UnsupportedOperationException.class)
public void getAnnotationsReadOnly() {
MilliSpan span = new MilliSpan(1, 2, "name", "traceId", Collections.<String>emptyList(), "spanId", true, "processId");
MilliSpan span = new MilliSpan(1, 2, "name", "traceId", Collections.<String>emptyList(), "spanId", true, true, "processId");
span.getAnnotations().put("a", "b");
}
@@ -36,7 +36,7 @@ public class MilliSpanTests {
@Test(expected = UnsupportedOperationException.class)
public void getTimelineAnnotationsReadOnly() {
MilliSpan span = new MilliSpan(1, 2, "name", "traceId", Collections.<String>emptyList(), "spanId", true, "processId");
MilliSpan span = new MilliSpan(1, 2, "name", "traceId", Collections.<String>emptyList(), "spanId", true, true, "processId");
span.getTimelineAnnotations().add(new TimelineAnnotation(1, "1"));
}

View File

@@ -5,6 +5,7 @@ import static org.assertj.core.api.BDDAssertions.then;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.Mockito;
@@ -24,6 +25,11 @@ public class TraceRunnableTest {
TraceManager traceManager = new DefaultTraceManager(new AlwaysSampler(),
new JdkIdGenerator(), Mockito.mock(ApplicationEventPublisher.class));
@Before
public void init() {
TraceContextHolder.removeCurrentTrace();
}
@Test
public void should_remove_span_from_thread_local_after_finishing_work()
throws Exception {

View File

@@ -17,7 +17,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@RunWith(SpringJUnit4ClassRunner.class)
@SpringApplicationConfiguration(classes = {ScheduledTestConfiguration.class})
public class TracingOnScheduledITest {
public class TracingOnScheduledTests {
@Autowired TestBeanWithScheduledMethod beanWithScheduledMethod;
@@ -36,7 +36,7 @@ public class TracingOnScheduledITest {
return new Runnable() {
@Override
public void run() {
Span storedSpan = TracingOnScheduledITest.this.beanWithScheduledMethod.getSpan();
Span storedSpan = TracingOnScheduledTests.this.beanWithScheduledMethod.getSpan();
then(storedSpan).isNotNull();
then(storedSpan.getTraceId()).isNotNull();
}
@@ -47,7 +47,7 @@ public class TracingOnScheduledITest {
return new Runnable() {
@Override
public void run() {
then(TracingOnScheduledITest.this.beanWithScheduledMethod.getSpan()).isNotEqualTo(spanToCompare);
then(TracingOnScheduledTests.this.beanWithScheduledMethod.getSpan()).isNotEqualTo(spanToCompare);
}
};
}

View File

@@ -17,22 +17,20 @@
package org.springframework.cloud.sleuth.instrument.web;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNull;
import static org.mockito.Matchers.any;
import static org.mockito.Matchers.anyString;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.MockitoAnnotations.initMocks;
import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.get;
import org.junit.Before;
import org.junit.Test;
import org.mockito.Mock;
import org.mockito.Mockito;
import org.springframework.cloud.sleuth.Sampler;
import org.springframework.cloud.sleuth.Span;
import org.springframework.cloud.sleuth.Trace;
import org.springframework.cloud.sleuth.TraceManager;
import org.springframework.cloud.sleuth.sampler.AlwaysSampler;
import org.springframework.cloud.sleuth.sampler.IsTracingSampler;
import org.springframework.cloud.sleuth.trace.DefaultTraceManager;
import org.springframework.cloud.sleuth.trace.TraceContextHolder;
import org.springframework.context.ApplicationEventPublisher;
@@ -62,12 +60,13 @@ public class TraceFilterTests {
private MockHttpServletRequest request;
private MockHttpServletResponse response;
private MockFilterChain filterChain;
private Sampler<Void> sampler = new AlwaysSampler();
@Before
@SneakyThrows
public void init() {
initMocks(this);
this.traceManager = new DefaultTraceManager(new AlwaysSampler(),
this.traceManager = new DefaultTraceManager(new DelegateSampler(),
new JdkIdGenerator(), this.publisher) {
@Override
protected Trace createTrace(Trace trace, Span span) {
@@ -88,16 +87,16 @@ public class TraceFilterTests {
@Test
public void notTraced() throws Exception {
TraceManager mockTraceManager = Mockito.mock(TraceManager.class);
TraceFilter filter = new TraceFilter(mockTraceManager);
this.sampler = new IsTracingSampler();
TraceFilter filter = new TraceFilter(this.traceManager);
this.request = get("/favicon.ico").accept(MediaType.ALL)
.buildRequest(new MockServletContext());
filter.doFilter(this.request, this.response, this.filterChain);
verify(mockTraceManager, never()).startSpan(anyString());
verify(mockTraceManager, never()).close(any(Trace.class));
assertFalse(this.span.isExportable());
assertNull(TraceContextHolder.getCurrentTrace());
}
@Test
@@ -153,4 +152,11 @@ public class TraceFilterTests {
private void hasAnnotation(Span span, String name, String value) {
assertEquals(value, span.getAnnotations().get(name));
}
private class DelegateSampler implements Sampler<Void> {
@Override
public boolean next(Void info) {
return TraceFilterTests.this.sampler.next(info);
}
}
}

View File

@@ -95,7 +95,9 @@ public class StreamSpanListener {
@Order(0)
public void release(SpanReleasedEvent event) {
event.getSpan().addTimelineAnnotation("release");
this.queue.add(event.getSpan());
if (event.getSpan().isExportable()) {
this.queue.add(event.getSpan());
}
}
@InboundChannelAdapter(value = SleuthSource.OUTPUT)

View File

@@ -50,8 +50,8 @@ public class ZipkinSpanListener {
private SpanCollector spanCollector;
/**
* Endpoint is the visible IP address of this service, the port it is
* listening on and the service name from discovery.
* Endpoint is the visible IP address of this service, the port it is listening on and
* the service name from discovery.
*/
private Endpoint localEndpoint;
@@ -63,7 +63,8 @@ public class ZipkinSpanListener {
@EventListener
@Order(0)
public void start(SpanAcquiredEvent event) {
// Starting a span in zipkin means adding: traceId, id, parentId(optional), and timestamp
// Starting a span in zipkin means adding: traceId, id, parentId(optional), and
// timestamp
event.getSpan().addTimelineAnnotation("acquire");
}
@@ -107,7 +108,9 @@ public class ZipkinSpanListener {
public void release(SpanReleasedEvent event) {
// Ending a span in zipkin means adding duration and sending it out
event.getSpan().addTimelineAnnotation("release");
this.spanCollector.collect(convert(event.getSpan()));
if (event.getSpan().isExportable()) {
this.spanCollector.collect(convert(event.getSpan()));
}
}
/**
@@ -120,12 +123,14 @@ public class ZipkinSpanListener {
*/
public com.twitter.zipkin.gen.Span convert(Span span) {
com.twitter.zipkin.gen.Span zipkinSpan = new com.twitter.zipkin.gen.Span();
addZipkinAnnotations(zipkinSpan, span, localEndpoint);
List<BinaryAnnotation> binaryAnnotationList = createZipkinBinaryAnnotations(span, localEndpoint);
addZipkinAnnotations(zipkinSpan, span, this.localEndpoint);
List<BinaryAnnotation> binaryAnnotationList = createZipkinBinaryAnnotations(span,
this.localEndpoint);
zipkinSpan.setDuration(span.getEnd() - span.getBegin());
zipkinSpan.setTrace_id(hash(span.getTraceId()));
if (span.getParents().size() > 0) {
if (span.getParents().size() > 1) {
log.error("zipkin doesn't support spans with multiple parents. Omitting "
log.error("Zipkin doesn't support spans with multiple parents. Omitting "
+ "other parents for " + span);
}
zipkinSpan.setParent_id(hash(span.getParents().get(0)));
@@ -138,12 +143,11 @@ public class ZipkinSpanListener {
return zipkinSpan;
}
/**
* Add annotations from the sleuth Span.
*/
private void addZipkinAnnotations(com.twitter.zipkin.gen.Span zipkinSpan,
Span span, Endpoint endpoint) {
private void addZipkinAnnotations(com.twitter.zipkin.gen.Span zipkinSpan, Span span,
Endpoint endpoint) {
Long startTs = null;
Long endTs = null;
for (TimelineAnnotation ta : span.getTimelineAnnotations()) {
@@ -151,9 +155,11 @@ public class ZipkinSpanListener {
ta.getTime(), endpoint);
if (zipkinAnnotation.getValue().equals("acquire")) {
startTs = zipkinAnnotation.getTimestamp();
} else if (zipkinAnnotation.getValue().equals("release")) {
}
else if (zipkinAnnotation.getValue().equals("release")) {
endTs = zipkinAnnotation.getTimestamp();
} else {
}
else {
zipkinSpan.addToAnnotations(zipkinAnnotation);
}
}
@@ -209,7 +215,7 @@ public class ZipkinSpanListener {
private static long hash(String string) {
long h = 1125899906842597L;
if (string==null) {
if (string == null) {
return h;
}
int len = string.length();