diff --git a/benchmarks/pom.xml b/benchmarks/pom.xml
index 338ad62be..09c4438dc 100644
--- a/benchmarks/pom.xml
+++ b/benchmarks/pom.xml
@@ -163,6 +163,9 @@
true
+
+ false
+
spring-milestones
@@ -172,14 +175,6 @@
false
-
- spring-staging
- Spring Staging
- https://repo.spring.io/libs-staging-local/
-
- false
-
-
spring-releases
Spring Releases
diff --git a/pom.xml b/pom.xml
index ec376b7e3..2332a11b8 100644
--- a/pom.xml
+++ b/pom.xml
@@ -222,6 +222,9 @@
true
+
+ false
+
spring-milestones
@@ -234,14 +237,6 @@
false
-
- spring-staging
- Spring Staging
- https://repo.spring.io/libs-staging-local/
-
- false
-
-
spring-releases
Spring Releases
diff --git a/spring-cloud-sleuth-dependencies/pom.xml b/spring-cloud-sleuth-dependencies/pom.xml
index 5175c6347..9c6786413 100644
--- a/spring-cloud-sleuth-dependencies/pom.xml
+++ b/spring-cloud-sleuth-dependencies/pom.xml
@@ -16,7 +16,8 @@
1.2.0.BUILD-SNAPSHOT
1.8.4
- 1.8.4
+ 1.11.1
+ 0.4.3
@@ -92,6 +93,11 @@
zipkin-junit
${zipkin.version}
+
+ io.zipkin.reporter
+ zipkin-reporter
+ ${zipkin-reporter.version}
+
diff --git a/spring-cloud-sleuth-samples/pom.xml b/spring-cloud-sleuth-samples/pom.xml
index f6471e8d6..002c1ee70 100644
--- a/spring-cloud-sleuth-samples/pom.xml
+++ b/spring-cloud-sleuth-samples/pom.xml
@@ -59,12 +59,12 @@
io.zipkin.java
zipkin
- 1.8.4
+ 1.11.1
io.zipkin.java
zipkin-server
- 1.8.4
+ 1.11.1
diff --git a/spring-cloud-sleuth-zipkin/pom.xml b/spring-cloud-sleuth-zipkin/pom.xml
index b6a82e16b..ba690c42b 100644
--- a/spring-cloud-sleuth-zipkin/pom.xml
+++ b/spring-cloud-sleuth-zipkin/pom.xml
@@ -65,6 +65,10 @@
io.zipkin.java
zipkin
+
+ io.zipkin.reporter
+ zipkin-reporter
+
org.springframework
spring-messaging
diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporter.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporter.java
index 9857e00d3..72284aa76 100644
--- a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporter.java
+++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporter.java
@@ -2,30 +2,13 @@ package org.springframework.cloud.sleuth.zipkin;
import java.io.Closeable;
import java.io.Flushable;
-import java.io.IOException;
-import java.net.URI;
-import java.nio.charset.Charset;
-import java.util.ArrayList;
-import java.util.LinkedList;
-import java.util.List;
-import java.util.concurrent.BlockingQueue;
-import java.util.concurrent.Executors;
-import java.util.concurrent.LinkedBlockingQueue;
-import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
-import org.apache.commons.logging.Log;
import org.springframework.cloud.sleuth.metric.SpanMetricReporter;
-import org.springframework.http.HttpHeaders;
-import org.springframework.http.HttpMethod;
-import org.springframework.http.MediaType;
-import org.springframework.http.RequestEntity;
-import org.springframework.web.client.RestClientException;
import org.springframework.web.client.RestTemplate;
-import zipkin.Codec;
import zipkin.Span;
-
-import static java.util.concurrent.TimeUnit.SECONDS;
+import zipkin.reporter.AsyncReporter;
/**
* Submits spans using Zipkin's {@code POST /spans} endpoint.
@@ -33,17 +16,9 @@ import static java.util.concurrent.TimeUnit.SECONDS;
* @author Adrian Cole
* @since 1.0.0
*/
-public final class HttpZipkinSpanReporter
- implements ZipkinSpanReporter, Flushable, Closeable {
- private static final Log log = org.apache.commons.logging.LogFactory
- .getLog(HttpZipkinSpanReporter.class);
- private static final Charset UTF_8 = Charset.forName("UTF-8");
-
- private final RestTemplate restTemplate;
- private final String url;
- private final BlockingQueue pending = new LinkedBlockingQueue<>(1000);
- private final Flusher flusher; // Nullable for testing
- private final SpanMetricReporter spanMetricReporter;
+public final class HttpZipkinSpanReporter implements ZipkinSpanReporter, Flushable, Closeable {
+ private final RestTemplateSender sender;
+ private final AsyncReporter delegate;
/**
* @param restTemplate {@link RestTemplate} used for sending requests to Zipkin
@@ -53,10 +28,12 @@ public final class HttpZipkinSpanReporter
*/
public HttpZipkinSpanReporter(RestTemplate restTemplate, String baseUrl, int flushInterval,
SpanMetricReporter spanMetricReporter) {
- this.restTemplate = restTemplate;
- this.url = baseUrl + (baseUrl.endsWith("/") ? "" : "/") + "api/v1/spans";
- this.flusher = flushInterval > 0 ? new Flusher(this, flushInterval) : null;
- this.spanMetricReporter = spanMetricReporter;
+ this.sender = new RestTemplateSender(restTemplate, baseUrl);
+ this.delegate = AsyncReporter.builder(this.sender)
+ .queuedMaxSpans(1000) // historical constraint. Note: AsyncReporter supports memory bounds
+ .messageTimeout(flushInterval, TimeUnit.SECONDS)
+ .metrics(new ReporterMetricsAdapter(spanMetricReporter))
+ .build();
}
/**
@@ -66,10 +43,7 @@ public final class HttpZipkinSpanReporter
*/
@Override
public void report(Span span) {
- this.spanMetricReporter.incrementAcceptedSpans(1);
- if (!this.pending.offer(span)) {
- this.spanMetricReporter.incrementDroppedSpans(1);
- }
+ this.delegate.report(span);
}
/**
@@ -77,78 +51,15 @@ public final class HttpZipkinSpanReporter
*/
@Override
public void flush() {
- if (this.pending.isEmpty())
- return;
- List drained = new ArrayList<>(this.pending.size());
- this.pending.drainTo(drained);
- if (drained.isEmpty())
- return;
-
- // json-encode the spans for transport
- byte[] json = Codec.JSON.writeSpans(drained);
- // NOTE: https://github.com/openzipkin/zipkin-java/issues/66 will throw instead of return null.
- if (json == null) {
- if (log.isDebugEnabled()) {
- log.debug("failed to encode spans, dropping them: " + drained);
- }
- this.spanMetricReporter.incrementDroppedSpans(drained.size());
- return;
- }
-
- // Send the json to the zipkin endpoint
- try {
- postSpans(json);
- }
- catch (RestClientException e) {
- if (log.isDebugEnabled()) { // don't pollute logs unless debug is on.
- // TODO: logger test
- log.debug(
- "error POSTing spans to " + this.url + ": as json: " + new String(json,
- UTF_8), e);
- }
- this.spanMetricReporter.incrementDroppedSpans(drained.size());
- }
+ this.delegate.flush();
}
/**
- * Calls flush on a fixed interval
- */
- static final class Flusher implements Runnable {
- final Flushable flushable;
- final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
-
- Flusher(Flushable flushable, int flushInterval) {
- this.flushable = flushable;
- this.scheduler.scheduleWithFixedDelay(this, 0, flushInterval, SECONDS);
- }
-
- @Override
- public void run() {
- try {
- this.flushable.flush();
- }
- catch (IOException ignored) {
- }
- }
- }
-
- void postSpans(byte[] json) {
- HttpHeaders httpHeaders = new HttpHeaders();
- httpHeaders.setContentType(MediaType.APPLICATION_JSON);
- RequestEntity requestEntity = new RequestEntity<>(json, httpHeaders, HttpMethod.POST, URI.create(this.url));
- this.restTemplate.exchange(requestEntity, String.class);
- }
-
- /**
- * Requests a cease of delivery. There will be at most one in-flight request processing after this
- * call returns.
+ * Blocks until in-flight spans are sent and drops any that are left pending.
*/
@Override
public void close() {
- if (this.flusher != null)
- this.flusher.scheduler.shutdown();
- // throw any outstanding spans on the floor
- int dropped = this.pending.drainTo(new LinkedList<>());
- this.spanMetricReporter.incrementDroppedSpans(dropped);
+ this.delegate.close();
+ this.sender.close();
}
}
diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ReporterMetricsAdapter.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ReporterMetricsAdapter.java
new file mode 100644
index 000000000..b1772ff64
--- /dev/null
+++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ReporterMetricsAdapter.java
@@ -0,0 +1,43 @@
+package org.springframework.cloud.sleuth.zipkin;
+
+import org.springframework.cloud.sleuth.metric.SpanMetricReporter;
+
+import zipkin.reporter.ReporterMetrics;
+
+final class ReporterMetricsAdapter implements ReporterMetrics {
+ private final SpanMetricReporter spanMetricReporter;
+
+ public ReporterMetricsAdapter(SpanMetricReporter spanMetricReporter) {
+ this.spanMetricReporter = spanMetricReporter;
+ }
+
+ @Override
+ public ReporterMetrics forTransport(String transport) {
+ return this;
+ }
+
+ @Override
+ public void incrementMessages() {
+ }
+
+ @Override
+ public void incrementMessagesDropped() {
+ }
+
+ @Override
+ public void incrementSpans(int i) {
+ this.spanMetricReporter.incrementAcceptedSpans(i);
+ }
+
+ @Override
+ public void incrementSpanBytes(int i) {
+ }
+
+ @Override
+ public void incrementMessageBytes(int i) {
+ }
+
+ @Override public void incrementSpansDropped(int i) {
+ this.spanMetricReporter.incrementDroppedSpans(i);
+ }
+}
diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/RestTemplateSender.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/RestTemplateSender.java
new file mode 100644
index 000000000..90c27c7d2
--- /dev/null
+++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/RestTemplateSender.java
@@ -0,0 +1,75 @@
+package org.springframework.cloud.sleuth.zipkin;
+
+import java.net.URI;
+import java.util.List;
+
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.HttpMethod;
+import org.springframework.http.MediaType;
+import org.springframework.http.RequestEntity;
+import org.springframework.web.client.RestTemplate;
+
+import zipkin.reporter.BytesMessageEncoder;
+import zipkin.reporter.Callback;
+import zipkin.reporter.Encoding;
+import zipkin.reporter.Sender;
+
+final class RestTemplateSender implements Sender {
+ final RestTemplate restTemplate;
+ final String url;
+
+ RestTemplateSender(RestTemplate restTemplate, String baseUrl) {
+ this.restTemplate = restTemplate;
+ this.url = baseUrl + (baseUrl.endsWith("/") ? "" : "/") + "api/v1/spans";
+ }
+
+ @Override public Encoding encoding() {
+ return Encoding.JSON;
+ }
+
+ @Override public int messageMaxBytes() {
+ // This will drop a span larger than 5MiB. Note: values like 512KiB benchmark better.
+ return 5 * 1024 * 1024;
+ }
+
+ @Override public int messageSizeInBytes(List spans) {
+ return encoding().listSizeInBytes(spans);
+ }
+
+ /** close is typically called from a different thread */
+ transient boolean closeCalled;
+
+ @Override public void sendSpans(List encodedSpans, Callback callback) {
+ if (this.closeCalled) throw new IllegalStateException("close");
+ try {
+ byte[] message = BytesMessageEncoder.JSON.encode(encodedSpans);
+ post(message);
+ callback.onComplete();
+ } catch (Throwable e) {
+ callback.onError(e);
+ if (e instanceof Error) throw (Error) e;
+ }
+ }
+
+ /** Sends an empty json message to the configured endpoint. */
+ @Override public CheckResult check() {
+ try {
+ post(new byte[] {'[', ']'});
+ return CheckResult.OK;
+ } catch (Exception e) {
+ return CheckResult.failed(e);
+ }
+ }
+
+ @Override public void close() {
+ this.closeCalled = true;
+ }
+
+ void post(byte[] json) {
+ HttpHeaders httpHeaders = new HttpHeaders();
+ httpHeaders.setContentType(MediaType.APPLICATION_JSON);
+ RequestEntity requestEntity =
+ new RequestEntity<>(json, httpHeaders, HttpMethod.POST, URI.create(this.url));
+ this.restTemplate.exchange(requestEntity, String.class);
+ }
+}