diff --git a/benchmarks/pom.xml b/benchmarks/pom.xml
index 45295d16b..4dbbf5e39 100644
--- a/benchmarks/pom.xml
+++ b/benchmarks/pom.xml
@@ -25,16 +25,23 @@
2.2.8.BUILD-SNAPSHOT
benchmarks
+
+ org.springframework.boot
+ spring-boot-starter-parent
+ 2.3.9.RELEASE
+
+
+
${project.basedir}/..
- 1.22
3.2.1
true
1.8
1.8
- 2.3.8.RELEASE
- 5.12.7
- 3.14.6
+ 4.9.0
+ 0.2.0.RELEASE
+ 1.21
+ Horsham.SR11
@@ -47,10 +54,9 @@
import
-
- org.springframework.boot
- spring-boot-dependencies
- ${spring-boot.version}
+ org.springframework.cloud
+ spring-cloud-stream-dependencies
+ ${spring-cloud-stream.version}
pom
import
@@ -91,20 +97,35 @@
org.assertj
assertj-core
- 3.14.0
compile
- org.hamcrest
- hamcrest-core
- 1.3
+ org.springframework.cloud
+ spring-cloud-starter-stream-kafka
+
+
+ org.springframework.cloud
+ spring-cloud-stream
+ test-jar
+ compile
+ test-binder
+
+
+ org.springframework.boot
+ spring-boot-starter-test
compile
-
- org.openjdk.jmh
- jmh-core
- ${jmh.version}
+ com.github.mp911de.microbenchmark-runner
+ microbenchmark-runner-junit5
+ ${microbenchmark-runner.version}
+ test
+
+
+ com.github.mp911de.microbenchmark-runner
+ microbenchmark-runner-extras
+ ${microbenchmark-runner.version}
+ test
@@ -122,12 +143,16 @@
io.zipkin.brave
brave-instrumentation-httpclient
- ${brave.version}
org.apache.httpcomponents
httpclient
+
+ org.awaitility
+ awaitility
+ test
+
@@ -140,28 +165,19 @@
${maven.compiler.target}
-
-
- maven-deploy-plugin
-
- true
-
-
-
- maven-install-plugin
-
- true
-
-
+
+ jitpack.io
+ https://jitpack.io
+
spring-snapshots
Spring Snapshots
- https://repo.spring.io/libs-snapshot-local
+ https://repo.spring.io/snapshot
true
@@ -184,7 +200,7 @@
spring-milestones
Spring Milestones
- https://repo.spring.io/libs-milestone-local
+ https://repo.spring.io/milestone
false
@@ -202,7 +218,7 @@
spring-snapshots
Spring Snapshots
- https://repo.spring.io/libs-snapshot-local
+ https://repo.spring.io/snapshot
true
@@ -213,7 +229,7 @@
spring-milestones
Spring Milestones
- https://repo.spring.io/libs-milestone-local
+ https://repo.spring.io/milestone
false
@@ -221,7 +237,7 @@
spring-releases
Spring Releases
- https://repo.spring.io/libs-release-local
+ https://repo.spring.io/release
false
@@ -229,75 +245,6 @@
-
- jmh
-
- false
-
-
-
-
- maven-shade-plugin
- ${maven-shade-plugin.version}
-
-
- org.springframework.boot
- spring-boot-maven-plugin
- ${spring-boot.version}
-
-
-
- true
-
- true
-
-
- *:*
-
- META-INF/*.SF
- META-INF/*.DSA
- META-INF/*.RSA
-
-
-
-
-
-
- package
-
- shade
-
-
- benchmarks
-
-
- META-INF/spring.handlers
-
-
- META-INF/spring.factories
-
-
- META-INF/spring.schemas
-
-
-
- org.openjdk.jmh.Main
-
-
- false
-
-
-
-
-
-
-
-
jmeter
@@ -380,7 +327,7 @@
com.lazerycode.jmeter
jmeter-maven-plugin
- 1.10.1
+ 3.1.1
false
false
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/mvc/SleuthBenchmarkingSpringApp.java b/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/mvc/SleuthBenchmarkingSpringApp.java
index a47396692..99efe676c 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/mvc/SleuthBenchmarkingSpringApp.java
+++ b/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/mvc/SleuthBenchmarkingSpringApp.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2019 the original author or authors.
+ * Copyright 2013-2020 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.
@@ -16,10 +16,6 @@
package org.springframework.cloud.sleuth.benchmarks.app.mvc;
-import java.util.concurrent.Callable;
-import java.util.concurrent.ExecutionException;
-import java.util.concurrent.ExecutorService;
-import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.regex.Pattern;
@@ -27,7 +23,6 @@ import javax.annotation.PreDestroy;
import brave.Span;
import brave.Tracer;
-import brave.sampler.Sampler;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
@@ -41,28 +36,26 @@ import org.springframework.boot.web.servlet.server.ServletWebServerFactory;
import org.springframework.cloud.sleuth.annotation.ContinueSpan;
import org.springframework.cloud.sleuth.annotation.NewSpan;
import org.springframework.cloud.sleuth.annotation.SpanTag;
+import org.springframework.cloud.sleuth.benchmarks.app.mvc.controller.AsyncSimulationController;
import org.springframework.cloud.sleuth.instrument.web.SkipPatternProvider;
import org.springframework.context.ApplicationListener;
import org.springframework.context.annotation.Bean;
-import org.springframework.scheduling.annotation.Async;
+import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.EnableAsync;
import org.springframework.util.SocketUtils;
-import org.springframework.web.bind.annotation.RequestMapping;
-import org.springframework.web.bind.annotation.RestController;
/**
* @author Marcin Grzejszczak
*/
@SpringBootApplication
-@RestController
@EnableAsync
-public class SleuthBenchmarkingSpringApp
- implements ApplicationListener {
+public class SleuthBenchmarkingSpringApp implements ApplicationListener {
private static final Log log = LogFactory.getLog(SleuthBenchmarkingSpringApp.class);
- public final ExecutorService pool = Executors.newWorkStealingPool();
-
+ /**
+ * Port of the app.
+ */
public int port;
@Autowired(required = false)
@@ -71,28 +64,16 @@ public class SleuthBenchmarkingSpringApp
@Autowired
AClass aClass;
+ @Autowired
+ AsyncSimulationController controller;
+
public static void main(String... args) {
SpringApplication.run(SleuthBenchmarkingSpringApp.class, args);
}
- @RequestMapping("/foo")
- public String foo() {
- return "foo";
- }
-
- @RequestMapping("/bar")
- public Callable bar() {
- return () -> "bar";
- }
-
- @RequestMapping("/async")
- public String asyncHttp() throws ExecutionException, InterruptedException {
- return this.async().get();
- }
-
- @Async
- public Future async() {
- return this.pool.submit(() -> "async");
+ @PreDestroy
+ public void clean() {
+ this.controller.clean();
}
public String manualSpan() {
@@ -108,48 +89,43 @@ public class SleuthBenchmarkingSpringApp
this.port = event.getSource().getPort();
}
- @Bean
- public ServletWebServerFactory servletContainer(
- @Value("${server.port:0}") int serverPort) {
- log.info("Starting container at port [" + serverPort + "]");
- return new TomcatServletWebServerFactory(
- serverPort == 0 ? SocketUtils.findAvailableTcpPort() : serverPort);
+ public Future async() {
+ return this.controller.async();
}
- @PreDestroy
- public void clean() {
- this.pool.shutdownNow();
- }
+ @Configuration
+ static class Config {
+ @Autowired(required = false)
+ Tracer tracer;
- @Bean
- Sampler alwaysSampler() {
- return Sampler.ALWAYS_SAMPLE;
- }
+ @Bean
+ AnotherClass anotherClass() {
+ return new AnotherClass(this.tracer);
+ }
- @Bean
- AnotherClass anotherClass() {
- return new AnotherClass(this.tracer);
- }
+ @Bean
+ AClass aClass() {
+ return new AClass(this.tracer, anotherClass());
+ }
- @Bean
- AClass aClass() {
- return new AClass(this.tracer, anotherClass());
- }
+ @Bean
+ SkipPatternProvider patternProvider() {
+ return new SkipPatternProvider() {
+ @Override
+ public Pattern skipPattern() {
+ return Pattern.compile("");
+ }
+ };
+ }
- @Bean
- SkipPatternProvider patternProvider() {
- return new SkipPatternProvider() {
- @Override
- public Pattern skipPattern() {
- return Pattern.compile("");
- }
- };
- }
- public ExecutorService getPool() {
- return this.pool;
- }
+ @Bean
+ public ServletWebServerFactory servletContainer(@Value("${server.port:0}") int serverPort) {
+ log.info("Starting container at port [" + serverPort + "]");
+ return new TomcatServletWebServerFactory(serverPort == 0 ? SocketUtils.findAvailableTcpPort() : serverPort);
+ }
+ }
}
class AClass {
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/mvc/controller/AsyncSimulationController.java b/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/mvc/controller/AsyncSimulationController.java
new file mode 100644
index 000000000..94b17fdd8
--- /dev/null
+++ b/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/mvc/controller/AsyncSimulationController.java
@@ -0,0 +1,50 @@
+package org.springframework.cloud.sleuth.benchmarks.app.mvc.controller;
+
+import java.util.concurrent.Callable;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+
+import javax.annotation.PreDestroy;
+
+import org.springframework.scheduling.annotation.Async;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+/**
+ * @author Marcin Grzejszczak
+ */
+@RestController
+public class AsyncSimulationController {
+ private final ExecutorService pool = Executors.newWorkStealingPool();
+
+ @RequestMapping("/foo")
+ public String foo() {
+ return "foo";
+ }
+
+ @RequestMapping("/bar")
+ public Callable bar() {
+ return () -> "bar";
+ }
+
+ @RequestMapping("/async")
+ public String asyncHttp() throws ExecutionException, InterruptedException {
+ return this.async().get();
+ }
+
+ @Async
+ public Future async() {
+ return this.pool.submit(() -> "async");
+ }
+
+ @PreDestroy
+ public void clean() {
+ this.pool.shutdownNow();
+ }
+
+ public ExecutorService getPool() {
+ return this.pool;
+ }
+}
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/stream/SleuthBenchmarkingStreamApplication.java b/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/stream/SleuthBenchmarkingStreamApplication.java
new file mode 100644
index 000000000..a793ec2bd
--- /dev/null
+++ b/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/stream/SleuthBenchmarkingStreamApplication.java
@@ -0,0 +1,209 @@
+/*
+ * Copyright 2013-2020 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
+ *
+ * https://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.benchmarks.app.stream;
+
+import java.io.IOException;
+import java.time.Duration;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.function.Function;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+import reactor.core.scheduler.Scheduler;
+import reactor.core.scheduler.Schedulers;
+
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.cloud.stream.binder.test.InputDestination;
+import org.springframework.cloud.stream.binder.test.OutputDestination;
+import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
+import org.springframework.context.ConfigurableApplicationContext;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Import;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.support.MessageBuilder;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+@SpringBootApplication
+@Import(TestChannelBinderConfiguration.class)
+public class SleuthBenchmarkingStreamApplication {
+
+ private static final Logger log = LoggerFactory.getLogger(SleuthBenchmarkingStreamApplication.class);
+
+ public static void main(String[] args) throws InterruptedException, IOException {
+ // System.setProperty("spring.sleuth.enabled", "false");
+ // System.setProperty("spring.sleuth.reactor.instrumentation-type",
+ // "DECORATE_ON_EACH");
+ // System.setProperty("spring.sleuth.reactor.instrumentation-type",
+ // "DECORATE_ON_LAST");
+ // System.setProperty("spring.sleuth.reactor.instrumentation-type", "MANUAL");
+ System.setProperty("spring.sleuth.reactor.instrumentation-type", "MANUAL");
+ System.setProperty("spring.sleuth.function.type", "simple");
+ ConfigurableApplicationContext context = SpringApplication.run(SleuthBenchmarkingStreamApplication.class, args);
+ for (int i = 0; i < 1; i++) {
+ InputDestination input = context.getBean(InputDestination.class);
+ input.send(MessageBuilder.withPayload("hello".getBytes())
+ .setHeader("b3", "4883117762eb9420-4883117762eb9420-1").build());
+ log.info("Retrieving the message for tests");
+ OutputDestination output = context.getBean(OutputDestination.class);
+ Message message = output.receive(200L);
+ log.info("Got the message from output");
+ assertThat(message).isNotNull();
+ log.info("Message is not null");
+ assertThat(message.getPayload()).isEqualTo("HELLO".getBytes());
+ log.info("Payload is HELLO");
+ String b3 = message.getHeaders().get("b3", String.class);
+ log.info("Checking the b3 header [" + b3 + "]");
+ assertThat(b3).startsWith("4883117762eb9420");
+ }
+ }
+
+ @Bean
+ ExecutorService sleuthExecutorService() {
+ return Executors.newCachedThreadPool();
+ }
+
+ @Bean(name = "myFlux")
+ @ConditionalOnProperty(value = "spring.sleuth.function.type", havingValue = "simple")
+ public Function simple() {
+ log.info("simple_function");
+ return new SimpleFunction();
+ }
+
+ @Bean(name = "myFlux")
+ @ConditionalOnProperty(value = "spring.sleuth.function.type", havingValue = "reactive_simple")
+ public Function, Flux> reactiveSimple() {
+ log.info("simple_reactive_function");
+ return new SimpleReactiveFunction();
+ }
+
+ @Bean(name = "myFlux")
+ @ConditionalOnProperty(value = "spring.sleuth.function.type", havingValue = "simple_function_with_around")
+ public Function, Message> simpleFunctionWithAround() {
+ log.info("simple_function_with_around");
+ return new SimpleMessageFunction();
+ }
+
+ @Bean(name = "myFlux")
+ @ConditionalOnProperty(value = "spring.sleuth.nonreactive.function.enabled", havingValue = "true")
+ public Function nonReactiveFunction(ExecutorService executorService) {
+ log.info("no sleuth non reactive function");
+ return new SleuthNonReactiveFunction(executorService);
+ }
+
+ @Bean(name = "myFlux")
+ @ConditionalOnProperty(value = "spring.sleuth.function.type", havingValue = "DECORATE_ON_EACH",
+ matchIfMissing = true)
+ public Function, Flux> onEachFunction() {
+ log.info("on each function");
+ return new SleuthFunction();
+ }
+
+ @Bean(name = "myFlux")
+ @ConditionalOnProperty(value = "spring.sleuth.function.type", havingValue = "DECORATE_ON_LAST")
+ public Function, Flux> onLastFunction() {
+ log.info("on last function");
+ return new SleuthFunction();
+ }
+
+}
+
+class SimpleFunction implements Function {
+
+ private static final Logger log = LoggerFactory.getLogger(SimpleFunction.class);
+
+ @Override
+ public String apply(String input) {
+ // tracing works cause headers from the input message get propagated to the output
+ // message
+ log.info("Hello from simple [{}]", input);
+ return input.toUpperCase();
+ }
+
+}
+
+class SimpleReactiveFunction implements Function, Flux> {
+
+ private static final Logger log = LoggerFactory.getLogger(SimpleReactiveFunction.class);
+
+ @Override
+ public Flux apply(Flux input) {
+ return input.doOnNext(s -> log.info("Hello from simple [{}]", s)).map(String::toUpperCase);
+ }
+
+}
+
+class SimpleMessageFunction implements Function, Message> {
+
+ private static final Logger log = LoggerFactory.getLogger(SimpleFunction.class);
+
+ @Override
+ public Message apply(Message input) {
+ log.info("Hello from message simple [{}]", input.getPayload());
+ return MessageBuilder.withPayload(input.getPayload().toUpperCase()).build();
+ }
+
+}
+
+class SleuthNonReactiveFunction implements Function {
+
+ private static final Logger log = LoggerFactory.getLogger(SleuthNonReactiveFunction.class);
+
+ private final ExecutorService executorService;
+
+ SleuthNonReactiveFunction(ExecutorService executorService) {
+ this.executorService = executorService;
+ }
+
+ @Override
+ public String apply(String input) {
+ log.info("Got a message");
+ try {
+ return this.executorService.submit(() -> {
+ log.info("Logging [{}] from a new thread", input);
+ return input.toUpperCase();
+ }).get(20, TimeUnit.MILLISECONDS);
+ }
+ catch (Exception e) {
+ throw new IllegalStateException(e);
+ }
+ }
+
+}
+
+class SleuthFunction implements Function, Flux> {
+
+ private static final Logger log = LoggerFactory.getLogger(SleuthFunction.class);
+
+ static final Scheduler SCHEDULER = Schedulers.newParallel("sleuthFunction");
+
+ @Override
+ public Flux apply(Flux input) {
+ return input.doOnEach(signal -> log.info("Got a message"))
+ .flatMap(s -> Mono.delay(Duration.ofMillis(1), SCHEDULER).map(aLong -> {
+ log.info("Logging [{}] from flat map", s);
+ return s.toUpperCase();
+ }));
+ }
+
+}
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/webflux/SleuthBenchmarkingSpringWebFluxApp.java b/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/webflux/SleuthBenchmarkingSpringWebFluxApp.java
index 631919c6b..121b9ceb1 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/webflux/SleuthBenchmarkingSpringWebFluxApp.java
+++ b/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/webflux/SleuthBenchmarkingSpringWebFluxApp.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2019 the original author or authors.
+ * Copyright 2013-2020 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.
@@ -16,13 +16,17 @@
package org.springframework.cloud.sleuth.benchmarks.app.webflux;
+import java.time.Duration;
import java.util.regex.Pattern;
+import java.util.stream.Collectors;
-import brave.sampler.Sampler;
-import brave.handler.SpanHandler;
-import org.apache.commons.logging.Log;
-import org.apache.commons.logging.LogFactory;
+import brave.propagation.TraceContext;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
+import reactor.core.scheduler.Scheduler;
+import reactor.core.scheduler.Schedulers;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.WebApplicationType;
@@ -33,7 +37,9 @@ import org.springframework.boot.web.reactive.context.ReactiveWebServerInitialize
import org.springframework.cloud.sleuth.instrument.web.SkipPatternProvider;
import org.springframework.context.ApplicationListener;
import org.springframework.context.annotation.Bean;
+import org.springframework.util.Assert;
import org.springframework.util.SocketUtils;
+import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@@ -42,17 +48,20 @@ import org.springframework.web.bind.annotation.RestController;
*/
@SpringBootApplication
@RestController
-public class SleuthBenchmarkingSpringWebFluxApp
- implements ApplicationListener {
+public class SleuthBenchmarkingSpringWebFluxApp implements ApplicationListener {
- private static final Log log = LogFactory
- .getLog(SleuthBenchmarkingSpringWebFluxApp.class);
+ static final Scheduler FOO_SCHEDULER = Schedulers.newParallel("foo");
+ private static final Logger log = LoggerFactory.getLogger(SleuthBenchmarkingSpringWebFluxApp.class);
+
+ /**
+ * Port to set.
+ */
public int port;
public static void main(String... args) {
- new SpringApplicationBuilder(SleuthBenchmarkingSpringWebFluxApp.class)
- .web(WebApplicationType.REACTIVE).application().run(args);
+ new SpringApplicationBuilder(SleuthBenchmarkingSpringWebFluxApp.class).web(WebApplicationType.REACTIVE)
+ .application().run(args);
}
@RequestMapping("/foo")
@@ -60,29 +69,15 @@ public class SleuthBenchmarkingSpringWebFluxApp
return Mono.just("foo");
}
- @Bean
- Sampler alwaysSampler() {
- return Sampler.ALWAYS_SAMPLE;
- }
-
@Bean
SkipPatternProvider patternProvider() {
return () -> Pattern.compile("");
}
@Bean
- NettyReactiveWebServerFactory nettyReactiveWebServerFactory(
- @Value("${server.port:0}") int serverPort) {
+ NettyReactiveWebServerFactory nettyReactiveWebServerFactory(@Value("${server.port:0}") int serverPort) {
log.info("Starting container at port [" + serverPort + "]");
- return new NettyReactiveWebServerFactory(
- serverPort == 0 ? SocketUtils.findAvailableTcpPort() : serverPort);
- }
-
- @Bean
- public SpanHandler spanHandler() {
- return new SpanHandler() {
- // intentionally anonymous to prevent logging fallback on NOOP
- };
+ return new NettyReactiveWebServerFactory(serverPort == 0 ? SocketUtils.findAvailableTcpPort() : serverPort);
}
@Override
@@ -90,4 +85,40 @@ public class SleuthBenchmarkingSpringWebFluxApp
this.port = event.getWebServer().getPort();
}
+ @GetMapping("/simple")
+ public Mono simple() {
+ return Mono.just("hello").map(String::toUpperCase).doOnNext(s -> log.info("Hello from simple [{}]", s));
+ }
+
+
+ @GetMapping("/complexNoSleuth")
+ public Mono complexNoSleuth() {
+ return Flux.range(1, 10).map(String::valueOf).collect(Collectors.toList())
+ .doOnEach(signal -> log.info("Got a request"))
+ .flatMap(s -> Mono.delay(Duration.ofMillis(1), FOO_SCHEDULER).map(aLong -> {
+ log.info("Logging [{}] from flat map", s);
+ return "";
+ }));
+ }
+
+ @GetMapping("/complex")
+ public Mono complex() {
+ return Flux.range(1, 10).map(String::valueOf).collect(Collectors.toList())
+ .doOnEach(signal -> log.info("Got a request"))
+ .flatMap(s -> Mono.delay(Duration.ofMillis(1), FOO_SCHEDULER).map(aLong -> {
+ log.info("Logging [{}] from flat map", s);
+ return "";
+ })).doOnEach(signal -> {
+ log.info("Doing assertions");
+ TraceContext traceContext = signal.getContext().get(TraceContext.class);
+ Assert.notNull(traceContext, "Context must be set by Sleuth instrumentation");
+ if (traceContext.traceIdString().startsWith("0000000000000000")) {
+ Assert.state(traceContext.traceIdString().equals("00000000000000004883117762eb9420"), "TraceId must be propagated");
+ } else {
+ Assert.state(traceContext.traceIdString().equals("4883117762eb9420"), "TraceId must be propagated");
+ }
+ log.info("Assertions passed");
+ });
+ }
+
}
diff --git a/benchmarks/src/main/resources/application.yml b/benchmarks/src/main/resources/application.yml
index 246f23e1c..8ffac9305 100644
--- a/benchmarks/src/main/resources/application.yml
+++ b/benchmarks/src/main/resources/application.yml
@@ -1,3 +1,5 @@
logging.level:
org.springframework: ERROR
- org.springframework.cloud.sleuth.benchmarks: INFO
+ org.springframework.sleuth: ERROR
+ org.springframework.sleuth.benchmarks: INFO
+ brave: ERROR
diff --git a/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/Pair.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/Pair.java
new file mode 100644
index 000000000..bad8a8d5a
--- /dev/null
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/Pair.java
@@ -0,0 +1,51 @@
+/*
+ * Copyright 2016-2019 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
+ *
+ * https://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.benchmarks.jmh;
+
+public class Pair {
+ final String key;
+ final String value;
+
+ public Pair(String key, String value) {
+ this.key = key;
+ this.value = value;
+ }
+
+ public String asProp() {
+ return this.key + "=" + this.value;
+ }
+
+ public static Pair of(String key, String value) {
+ return new Pair(key, value);
+ }
+
+ public static Pair noHook() {
+ return new Pair("spring.sleuth.reactor.decorate-hooks", "false");
+ }
+
+ public static Pair noSleuth() {
+ return new Pair("spring.sleuth.enabled", "false");
+ }
+
+ public static Pair onEach() {
+ return new Pair("spring.sleuth.reactor.decorate-on-each", "true");
+ }
+
+ public static Pair onLast() {
+ return new Pair("spring.sleuth.reactor.decorate-on-each", "false");
+ }
+}
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/ProcessLauncherState.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/ProcessLauncherState.java
similarity index 93%
rename from benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/ProcessLauncherState.java
rename to benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/ProcessLauncherState.java
index 9478d6955..7a57bfc13 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/ProcessLauncherState.java
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/ProcessLauncherState.java
@@ -14,7 +14,7 @@
* limitations under the License.
*/
-package org.springframework.cloud.sleuth.benchmarks.jmh.benchmarks;
+package org.springframework.cloud.sleuth.benchmarks.jmh;
import java.io.BufferedReader;
import java.io.File;
@@ -57,15 +57,13 @@ public class ProcessLauncherState {
this.args.add(count++, "-Djava.security.egd=file:/dev/./urandom");
this.args.add(count++, "-XX:TieredStopAtLevel=1"); // zoom
if (System.getProperty("bench.args") != null) {
- this.args.addAll(count++,
- Arrays.asList(System.getProperty("bench.args").split(" ")));
+ this.args.addAll(count++, Arrays.asList(System.getProperty("bench.args").split(" ")));
}
this.length = args.length;
this.home = new File(dir);
}
- protected static String output(InputStream inputStream, String marker)
- throws IOException {
+ protected static String output(InputStream inputStream, String marker) throws IOException {
StringBuilder sb = new StringBuilder();
BufferedReader br = null;
br = new BufferedReader(new InputStreamReader(inputStream));
@@ -100,8 +98,7 @@ public class ProcessLauncherState {
public void after() throws Exception {
if (started != null && started.isAlive()) {
- System.err.println(
- "Stopped " + mainClass + ": " + started.destroyForcibly().waitFor());
+ System.err.println("Stopped " + mainClass + ": " + started.destroyForcibly().waitFor());
}
}
@@ -120,6 +117,7 @@ public class ProcessLauncherState {
}
public void run() throws Exception {
+ System.out.println("Running process");
List args = new ArrayList<>(this.args);
args.add(args.size() - this.length, this.mainClass);
if (extraArgs != null) {
diff --git a/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/RunSleuthJmhBenchmarksFromIde.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/RunSleuthJmhBenchmarksFromIde.java
deleted file mode 100644
index 1c6e7d03e..000000000
--- a/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/RunSleuthJmhBenchmarksFromIde.java
+++ /dev/null
@@ -1,36 +0,0 @@
-/*
- * Copyright 2013-2019 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
- *
- * https://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.benchmarks.jmh;
-
-import org.openjdk.jmh.runner.Runner;
-import org.openjdk.jmh.runner.RunnerException;
-import org.openjdk.jmh.runner.options.Options;
-import org.openjdk.jmh.runner.options.OptionsBuilder;
-
-public class RunSleuthJmhBenchmarksFromIde {
-
- // Convenience main entry-point for testing from IDE
- public static void main(String[] args) throws RunnerException {
- Options opt = new OptionsBuilder()
- .include(RunSleuthJmhBenchmarksFromIde.class.getPackage().getName()
- + ".benchmarks.*")
- .build();
-
- new Runner(opt).run();
- }
-
-}
diff --git a/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/SampleTests.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/SampleTests.java
new file mode 100644
index 000000000..3cadc3077
--- /dev/null
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/SampleTests.java
@@ -0,0 +1,163 @@
+/*
+ * Copyright 2013-2020 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
+ *
+ * https://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.benchmarks.jmh;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import java.util.stream.Collectors;
+
+import brave.Tracing;
+import org.junit.jupiter.api.Disabled;
+import org.junit.jupiter.api.Test;
+import org.openjdk.jmh.annotations.Param;
+import org.openjdk.jmh.annotations.Setup;
+import org.openjdk.jmh.annotations.TearDown;
+
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.WebApplicationType;
+import org.springframework.boot.builder.SpringApplicationBuilder;
+import org.springframework.cloud.sleuth.benchmarks.app.stream.SleuthBenchmarkingStreamApplication;
+import org.springframework.cloud.stream.binder.test.InputDestination;
+import org.springframework.cloud.stream.binder.test.OutputDestination;
+import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
+import org.springframework.context.ConfigurableApplicationContext;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.context.annotation.Import;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.support.MessageBuilder;
+import org.springframework.util.StringUtils;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+@Disabled
+public class SampleTests {
+
+ @Test
+ public void testStream() throws Exception {
+ for (BenchmarkContext.Instrumentation value : BenchmarkContext.Instrumentation.values()) {
+ run(value);
+ }
+ // run(BenchmarkContext.Instrumentation.sleuthReactiveSimpleManual);
+ }
+
+ private void run(BenchmarkContext.Instrumentation value) throws Exception {
+ BenchmarkContext context = new BenchmarkContext();
+ System.out.println("\n\n\n\n WILL WORK WITH [" + value + "]\n\n\n\n");
+ context.instrumentation = value;
+ context.setup();
+
+ try {
+ context.run(value);
+ }
+ finally {
+ context.clean();
+ }
+ System.out.println("\n\n FINISHED WITH [" + value + "]\n\n\n\n");
+ }
+
+ public static class BenchmarkContext {
+
+ volatile ConfigurableApplicationContext applicationContext;
+
+ volatile InputDestination input;
+
+ volatile OutputDestination output;
+
+ @Param
+ private Instrumentation instrumentation;
+
+ @Setup
+ public void setup() {
+ this.applicationContext = initContext();
+ this.input = this.applicationContext.getBean(InputDestination.class);
+ this.output = this.applicationContext.getBean(OutputDestination.class);
+ }
+
+ protected ConfigurableApplicationContext initContext() {
+ SpringApplication application = new SpringApplicationBuilder(SleuthBenchmarkingStreamApplication.class)
+ .web(WebApplicationType.REACTIVE).application();
+ return application.run(runArgs());
+ }
+
+ protected String[] runArgs() {
+ List strings = new ArrayList<>();
+ strings.addAll(Arrays.asList("--spring.jmx.enabled=false",
+ "--spring.application.name=defaultTraceContextForStream" + instrumentation.name()));
+ strings.addAll(instrumentation.entires.stream().map(s -> "--" + s).collect(Collectors.toList()));
+ return strings.toArray(new String[0]);
+ }
+
+ void run(Instrumentation value) {
+ System.out.println("Sending the message to input");
+ input.send(MessageBuilder.withPayload("hello".getBytes())
+ .setHeader("b3", "4883117762eb9420-4883117762eb9420-1").build());
+ System.out.println("Retrieving the message for tests");
+ Message message = output.receive(200L);
+ System.out.println("Got the message from output");
+ assertThat(message).isNotNull();
+ System.out.println("Message is not null");
+ assertThat(message.getPayload()).isEqualTo("HELLO".getBytes());
+ System.out.println("Payload is HELLO");
+ if (!value.toString().toLowerCase().contains("nosleuth")) {
+ String b3 = message.getHeaders().get("b3", String.class);
+ System.out.println("Checking the b3 header [" + b3 + "]");
+ assertThat(b3).startsWith("4883117762eb9420");
+ }
+ }
+
+ @TearDown
+ public void clean() throws Exception {
+ Tracing current = Tracing.current();
+ if (current != null) {
+ current.close();
+ }
+ try {
+ this.applicationContext.close();
+ }
+ catch (Exception ig) {
+
+ }
+ }
+
+ public enum Instrumentation {
+
+ noSleuthSimple("spring.sleuth.enabled=false,spring.sleuth.function.type=simple");
+
+ private Set entires = new HashSet<>();
+
+ Instrumentation(String key, String value) {
+ this.entires.add(key + "=" + value);
+ }
+
+ Instrumentation(String commaSeparated) {
+ this.entires.addAll(StringUtils.commaDelimitedListToSet(commaSeparated));
+ }
+
+ }
+
+ }
+
+ @Configuration(proxyBeanMethods = false)
+ @Import(TestChannelBinderConfiguration.class)
+ static class TestConfiguration {
+
+ }
+
+}
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/SpringWebFluxOnLastBenchmark.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/TracerImplementation.java
similarity index 56%
rename from benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/SpringWebFluxOnLastBenchmark.java
rename to benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/TracerImplementation.java
index b050885a6..75fa5373b 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/SpringWebFluxOnLastBenchmark.java
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/TracerImplementation.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2019 the original author or authors.
+ * Copyright 2016-2019 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.
@@ -14,16 +14,15 @@
* limitations under the License.
*/
-package org.springframework.cloud.sleuth.benchmarks.jmh.benchmarks;
+package org.springframework.cloud.sleuth.benchmarks.jmh;
-public class SpringWebFluxOnLastBenchmark extends SpringWebFluxBenchmarks {
+public enum TracerImplementation {
+
+ brave;
@Override
- protected String[] runArgs() {
- return new String[] { "--spring.jmx.enabled=false",
- "--spring.application.name=defaultTraceContextWithOnLastOperator",
- "--spring.sleuth.enabled=true",
- "--spring.sleuth.reactor.on-each-operator=false" };
+ public String toString() {
+ return this.name();
}
-}
+}
\ No newline at end of file
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/AnnotationBenchmarks.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/AnnotationBenchmarksTests.java
similarity index 78%
rename from benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/AnnotationBenchmarks.java
rename to benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/AnnotationBenchmarksTests.java
index c1bfb8c0f..8a4e5eeda 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/AnnotationBenchmarks.java
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/AnnotationBenchmarksTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2019 the original author or authors.
+ * Copyright 2013-2020 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.
@@ -14,16 +14,18 @@
* limitations under the License.
*/
-package org.springframework.cloud.sleuth.benchmarks.jmh.benchmarks;
+package org.springframework.cloud.sleuth.benchmarks.jmh.mvc;
import java.util.concurrent.TimeUnit;
+import jmh.mbr.junit5.Microbenchmark;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.BenchmarkMode;
import org.openjdk.jmh.annotations.Fork;
import org.openjdk.jmh.annotations.Measurement;
import org.openjdk.jmh.annotations.Mode;
import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
import org.openjdk.jmh.annotations.Scope;
import org.openjdk.jmh.annotations.Setup;
import org.openjdk.jmh.annotations.State;
@@ -33,17 +35,19 @@ import org.openjdk.jmh.annotations.Warmup;
import org.springframework.boot.SpringApplication;
import org.springframework.cloud.sleuth.benchmarks.app.mvc.SleuthBenchmarkingSpringApp;
+import org.springframework.cloud.sleuth.benchmarks.jmh.TracerImplementation;
import org.springframework.context.ConfigurableApplicationContext;
import static org.assertj.core.api.BDDAssertions.then;
-@Measurement(iterations = 5)
-@Warmup(iterations = 10)
-@Fork(3)
+@Measurement(iterations = 10, time = 1)
+@Warmup(iterations = 10, time = 1)
+@Fork(4)
@BenchmarkMode(Mode.SampleTime)
@OutputTimeUnit(TimeUnit.MICROSECONDS)
@Threads(Threads.MAX)
-public class AnnotationBenchmarks {
+@Microbenchmark
+public class AnnotationBenchmarksTests {
@Benchmark
public void manuallyCreatedSpans(BenchmarkContext context) throws Exception {
@@ -62,11 +66,14 @@ public class AnnotationBenchmarks {
volatile SleuthBenchmarkingSpringApp sleuth;
+ @Param
+ private TracerImplementation tracerImplementation;
+
@Setup
public void setup() {
- this.withSleuth = new SpringApplication(SleuthBenchmarkingSpringApp.class)
- .run("--spring.jmx.enabled=false",
- "--spring.application.name=withSleuth");
+ this.withSleuth = new SpringApplication(SleuthBenchmarkingSpringApp.class).run("--spring.jmx.enabled=false",
+
+ "--spring.application.name=withSleuth_" + this.tracerImplementation.name());
this.sleuth = this.withSleuth.getBean(SleuthBenchmarkingSpringApp.class);
}
diff --git a/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/AsyncWithSleuthBenchmarksTests.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/AsyncWithSleuthBenchmarksTests.java
new file mode 100644
index 000000000..d82fd607c
--- /dev/null
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/AsyncWithSleuthBenchmarksTests.java
@@ -0,0 +1,81 @@
+/*
+ * Copyright 2013-2020 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
+ *
+ * https://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.benchmarks.jmh.mvc;
+
+import java.util.concurrent.TimeUnit;
+
+import jmh.mbr.junit5.Microbenchmark;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.Setup;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.TearDown;
+import org.openjdk.jmh.annotations.Threads;
+import org.openjdk.jmh.annotations.Warmup;
+
+import org.springframework.boot.SpringApplication;
+import org.springframework.cloud.sleuth.benchmarks.app.mvc.SleuthBenchmarkingSpringApp;
+import org.springframework.cloud.sleuth.benchmarks.jmh.TracerImplementation;
+import org.springframework.context.ConfigurableApplicationContext;
+
+import static org.assertj.core.api.BDDAssertions.then;
+
+@Measurement(iterations = 10, time = 1)
+@Warmup(iterations = 10, time = 1)
+@Fork(4)
+@BenchmarkMode(Mode.SampleTime)
+@OutputTimeUnit(TimeUnit.MICROSECONDS)
+@Threads(Threads.MAX)
+@Microbenchmark
+public class AsyncWithSleuthBenchmarksTests {
+ @Benchmark
+ public void asyncMethodWithSleuth(BenchmarkContext context) throws Exception {
+ then(context.app.async().get()).isEqualTo("async");
+ }
+
+ @State(Scope.Benchmark)
+ public static class BenchmarkContext {
+ volatile ConfigurableApplicationContext context;
+ volatile SleuthBenchmarkingSpringApp app;
+
+ @Param
+ private TracerImplementation tracerImplementation;
+
+ @Setup
+ public void setup() {
+ this.context = new SpringApplication(SleuthBenchmarkingSpringApp.class).run(
+ "--spring.jmx.enabled=false",
+ "--spring.application.name=withSleuth_" + this.tracerImplementation.name()
+ );
+ this.app = this.context.getBean(SleuthBenchmarkingSpringApp.class);
+ }
+
+ @TearDown
+ public void clean() {
+ this.app.clean();
+ this.context.close();
+ }
+
+ }
+
+}
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/AsyncBenchmarks.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/AsyncWithoutSleuthBenchmarksTests.java
similarity index 53%
rename from benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/AsyncBenchmarks.java
rename to benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/AsyncWithoutSleuthBenchmarksTests.java
index 61fe76cc7..93ea2e7ce 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/AsyncBenchmarks.java
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/AsyncWithoutSleuthBenchmarksTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2019 the original author or authors.
+ * Copyright 2013-2020 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.
@@ -14,10 +14,11 @@
* limitations under the License.
*/
-package org.springframework.cloud.sleuth.benchmarks.jmh.benchmarks;
+package org.springframework.cloud.sleuth.benchmarks.jmh.mvc;
import java.util.concurrent.TimeUnit;
+import jmh.mbr.junit5.Microbenchmark;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.BenchmarkMode;
import org.openjdk.jmh.annotations.Fork;
@@ -37,57 +38,39 @@ import org.springframework.context.ConfigurableApplicationContext;
import static org.assertj.core.api.BDDAssertions.then;
-@Measurement(iterations = 5)
-@Warmup(iterations = 10)
-@Fork(3)
+@Measurement(iterations = 10, time = 1)
+@Warmup(iterations = 10, time = 1)
+@Fork(4)
@BenchmarkMode(Mode.SampleTime)
@OutputTimeUnit(TimeUnit.MICROSECONDS)
@Threads(Threads.MAX)
-public class AsyncBenchmarks {
-
+@Microbenchmark
+public class AsyncWithoutSleuthBenchmarksTests {
@Benchmark
public void asyncMethodWithoutSleuth(BenchmarkContext context) throws Exception {
- then(context.untracedAsyncMethodHavingBean.async().get()).isEqualTo("async");
- }
-
- @Benchmark
- public void asyncMethodWithSleuth(BenchmarkContext context) throws Exception {
- then(context.tracedAsyncMethodHavingBean.async().get()).isEqualTo("async");
+ then(context.app.async().get()).isEqualTo("async");
}
@State(Scope.Benchmark)
public static class BenchmarkContext {
-
- volatile ConfigurableApplicationContext withSleuth;
-
- volatile ConfigurableApplicationContext withoutSleuth;
-
- volatile SleuthBenchmarkingSpringApp tracedAsyncMethodHavingBean;
-
- volatile SleuthBenchmarkingSpringApp untracedAsyncMethodHavingBean;
+ volatile ConfigurableApplicationContext context;
+ volatile SleuthBenchmarkingSpringApp app;
@Setup
public void setup() {
- this.withSleuth = new SpringApplication(SleuthBenchmarkingSpringApp.class)
- .run("--spring.jmx.enabled=false",
- "--spring.application.name=withSleuth");
- this.withoutSleuth = new SpringApplication(SleuthBenchmarkingSpringApp.class)
- .run("--spring.jmx.enabled=false",
- "--spring.application.name=withoutSleuth",
- "--spring.sleuth.enabled=false",
- "--spring.sleuth.async.enabled=false");
- this.tracedAsyncMethodHavingBean = this.withSleuth
- .getBean(SleuthBenchmarkingSpringApp.class);
- this.untracedAsyncMethodHavingBean = this.withoutSleuth
- .getBean(SleuthBenchmarkingSpringApp.class);
+ this.context = new SpringApplication(SleuthBenchmarkingSpringApp.class).run(
+ "--spring.jmx.enabled=false",
+ "--spring.application.name=withoutSleuth",
+ "--spring.sleuth.enabled=false",
+ "--spring.sleuth.async.enabled=false"
+ );
+ this.app = this.context.getBean(SleuthBenchmarkingSpringApp.class);
}
@TearDown
public void clean() {
- this.tracedAsyncMethodHavingBean.clean();
- this.untracedAsyncMethodHavingBean.clean();
- this.withSleuth.close();
- this.withoutSleuth.close();
+ this.app.clean();
+ this.context.close();
}
}
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/HttpFilterBenchmarks.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/HttpFilterBenchmarksTests.java
similarity index 82%
rename from benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/HttpFilterBenchmarks.java
rename to benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/HttpFilterBenchmarksTests.java
index 1ad2ab604..4ea058739 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/HttpFilterBenchmarks.java
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/HttpFilterBenchmarksTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2019 the original author or authors.
+ * Copyright 2013-2020 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.
@@ -14,7 +14,7 @@
* limitations under the License.
*/
-package org.springframework.cloud.sleuth.benchmarks.jmh.benchmarks;
+package org.springframework.cloud.sleuth.benchmarks.jmh.mvc;
import java.io.IOException;
import java.util.concurrent.Callable;
@@ -28,12 +28,14 @@ import javax.servlet.ServletRequest;
import javax.servlet.ServletResponse;
import brave.servlet.TracingFilter;
+import jmh.mbr.junit5.Microbenchmark;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.BenchmarkMode;
import org.openjdk.jmh.annotations.Fork;
import org.openjdk.jmh.annotations.Measurement;
import org.openjdk.jmh.annotations.Mode;
import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
import org.openjdk.jmh.annotations.Scope;
import org.openjdk.jmh.annotations.Setup;
import org.openjdk.jmh.annotations.State;
@@ -43,6 +45,8 @@ import org.openjdk.jmh.annotations.Warmup;
import org.springframework.boot.SpringApplication;
import org.springframework.cloud.sleuth.benchmarks.app.mvc.SleuthBenchmarkingSpringApp;
+import org.springframework.cloud.sleuth.benchmarks.app.mvc.controller.AsyncSimulationController;
+import org.springframework.cloud.sleuth.benchmarks.jmh.TracerImplementation;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.http.MediaType;
import org.springframework.mock.web.MockFilterChain;
@@ -62,17 +66,17 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.request;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
-@Warmup(iterations = 10)
+@Warmup(iterations = 5)
@BenchmarkMode(Mode.SampleTime)
@OutputTimeUnit(TimeUnit.MICROSECONDS)
@Threads(Threads.MAX)
-public class HttpFilterBenchmarks {
+@Microbenchmark
+public class HttpFilterBenchmarksTests {
@Benchmark
@Measurement(iterations = 5, time = 1)
- @Fork(3)
- public void filterWithoutSleuth(BenchmarkContext context)
- throws IOException, ServletException {
+ @Fork(2)
+ public void filterWithoutSleuth(BenchmarkContext context) throws IOException, ServletException {
MockHttpServletRequest request = builder().buildRequest(new MockServletContext());
MockHttpServletResponse response = new MockHttpServletResponse();
response.setContentType(MediaType.APPLICATION_JSON_VALUE);
@@ -82,9 +86,8 @@ public class HttpFilterBenchmarks {
@Benchmark
@Measurement(iterations = 5, time = 1)
- @Fork(3)
- public void filterWithSleuth(BenchmarkContext context)
- throws ServletException, IOException {
+ @Fork(2)
+ public void filterWithSleuth(BenchmarkContext context) throws ServletException, IOException {
MockHttpServletRequest request = builder().buildRequest(new MockServletContext());
MockHttpServletResponse response = new MockHttpServletResponse();
response.setContentType(MediaType.APPLICATION_JSON_VALUE);
@@ -107,12 +110,10 @@ public class HttpFilterBenchmarks {
}
private MockHttpServletRequestBuilder builder() {
- return get("/").accept(MediaType.APPLICATION_JSON).header("User-Agent",
- "MockMvc");
+ return get("/").accept(MediaType.APPLICATION_JSON).header("User-Agent", "MockMvc");
}
- private void performRequest(MockMvc mockMvc, String url, String expectedResult)
- throws Exception {
+ private void performRequest(MockMvc mockMvc, String url, String expectedResult) throws Exception {
MvcResult mvcResult = mockMvc.perform(get("/" + url)).andExpect(status().isOk())
.andExpect(request().asyncStarted()).andReturn();
@@ -133,18 +134,18 @@ public class HttpFilterBenchmarks {
volatile MockMvc mockMvcForUntracedController;
+ @Param
+ private TracerImplementation tracerImplementation;
+
@Setup
public void setup() {
- this.withSleuth = new SpringApplication(SleuthBenchmarkingSpringApp.class)
- .run("--spring.jmx.enabled=false",
- "--spring.application.name=withSleuth");
+ this.withSleuth = new SpringApplication(SleuthBenchmarkingSpringApp.class).run("--spring.jmx.enabled=false",
+
+ "--spring.application.name=withSleuth_" + this.tracerImplementation.name());
this.tracingFilter = this.withSleuth.getBean(TracingFilter.class);
this.mockMvcForTracedController = MockMvcBuilders
- .standaloneSetup(
- this.withSleuth.getBean(SleuthBenchmarkingSpringApp.class))
- .build();
- this.mockMvcForUntracedController = MockMvcBuilders
- .standaloneSetup(new VanillaController()).build();
+ .standaloneSetup(this.withSleuth.getBean(AsyncSimulationController.class)).build();
+ this.mockMvcForUntracedController = MockMvcBuilders.standaloneSetup(new VanillaController()).build();
}
@TearDown
@@ -162,8 +163,8 @@ public class HttpFilterBenchmarks {
}
@Override
- public void doFilter(ServletRequest request, ServletResponse response,
- FilterChain chain) throws IOException, ServletException {
+ public void doFilter(ServletRequest request, ServletResponse response, FilterChain chain)
+ throws IOException, ServletException {
chain.doFilter(request, response);
}
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/RestTemplateBenchmark.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/RestTemplateBenchmarkTests.java
similarity index 66%
rename from benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/RestTemplateBenchmark.java
rename to benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/RestTemplateBenchmarkTests.java
index 8823ff1e6..c18673216 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/RestTemplateBenchmark.java
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/RestTemplateBenchmarkTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2019 the original author or authors.
+ * Copyright 2013-2020 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.
@@ -14,7 +14,7 @@
* limitations under the License.
*/
-package org.springframework.cloud.sleuth.benchmarks.jmh.benchmarks;
+package org.springframework.cloud.sleuth.benchmarks.jmh.mvc;
import java.io.IOException;
import java.util.Collections;
@@ -23,12 +23,14 @@ import java.util.concurrent.TimeUnit;
import javax.servlet.ServletException;
import brave.spring.web.TracingClientHttpRequestInterceptor;
+import jmh.mbr.junit5.Microbenchmark;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.BenchmarkMode;
import org.openjdk.jmh.annotations.Fork;
import org.openjdk.jmh.annotations.Measurement;
import org.openjdk.jmh.annotations.Mode;
import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
import org.openjdk.jmh.annotations.Scope;
import org.openjdk.jmh.annotations.Setup;
import org.openjdk.jmh.annotations.State;
@@ -38,6 +40,8 @@ import org.openjdk.jmh.annotations.Warmup;
import org.springframework.boot.SpringApplication;
import org.springframework.cloud.sleuth.benchmarks.app.mvc.SleuthBenchmarkingSpringApp;
+import org.springframework.cloud.sleuth.benchmarks.app.mvc.controller.AsyncSimulationController;
+import org.springframework.cloud.sleuth.benchmarks.jmh.TracerImplementation;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.test.web.client.MockMvcClientHttpRequestFactory;
import org.springframework.test.web.servlet.MockMvc;
@@ -49,24 +53,22 @@ import static org.assertj.core.api.BDDAssertions.then;
/**
* We're checking how much overhead does the instrumentation of the RestTemplate take
*/
-@Measurement(iterations = 5)
-@Warmup(iterations = 10)
-@Fork(3)
+@Measurement(iterations = 10, time = 1)
+@Warmup(iterations = 10, time = 1)
+@Fork(4)
@BenchmarkMode(Mode.SampleTime)
@OutputTimeUnit(TimeUnit.MICROSECONDS)
@Threads(Threads.MAX)
-public class RestTemplateBenchmark {
+@Microbenchmark
+public class RestTemplateBenchmarkTests {
@Benchmark
- public void syncEndpointWithoutSleuth(BenchmarkContext context)
- throws IOException, ServletException {
- then(context.untracedTemplate.getForObject("/foo", String.class))
- .isEqualTo("foo");
+ public void syncEndpointWithoutSleuth(BenchmarkContext context) throws IOException, ServletException {
+ then(context.untracedTemplate.getForObject("/foo", String.class)).isEqualTo("foo");
}
@Benchmark
- public void syncEndpointWithSleuth(BenchmarkContext context)
- throws ServletException, IOException {
+ public void syncEndpointWithSleuth(BenchmarkContext context) throws ServletException, IOException {
then(context.tracedTemplate.getForObject("/foo", String.class)).isEqualTo("foo");
}
@@ -81,20 +83,21 @@ public class RestTemplateBenchmark {
volatile RestTemplate untracedTemplate;
+ @Param
+ private TracerImplementation tracerImplementation;
+
@Setup
public void setup() {
- new SpringApplication(SleuthBenchmarkingSpringApp.class).run(
- "--spring.jmx.enabled=false", "--spring.application.name=withSleuth");
- this.mockMvc = MockMvcBuilders
- .standaloneSetup(
- this.withSleuth.getBean(SleuthBenchmarkingSpringApp.class))
+ this.withSleuth = new SpringApplication(SleuthBenchmarkingSpringApp.class).run(
+ "--spring.jmx.enabled=false",
+ "--spring.application.name=withSleuth_" + this.tracerImplementation.name()
+ );
+ this.mockMvc = MockMvcBuilders.standaloneSetup(this.withSleuth.getBean(AsyncSimulationController.class))
.build();
- this.tracedTemplate = new RestTemplate(
- new MockMvcClientHttpRequestFactory(this.mockMvc));
- this.tracedTemplate.setInterceptors(Collections.singletonList(
- this.withSleuth.getBean(TracingClientHttpRequestInterceptor.class)));
- this.untracedTemplate = new RestTemplate(
- new MockMvcClientHttpRequestFactory(this.mockMvc));
+ this.tracedTemplate = new RestTemplate(new MockMvcClientHttpRequestFactory(this.mockMvc));
+ this.tracedTemplate.setInterceptors(
+ Collections.singletonList(this.withSleuth.getBean(TracingClientHttpRequestInterceptor.class)));
+ this.untracedTemplate = new RestTemplate(new MockMvcClientHttpRequestFactory(this.mockMvc));
}
@TearDown
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/StartupBenchmark.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/StartupBenchmarkTests.java
similarity index 70%
rename from benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/StartupBenchmark.java
rename to benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/StartupBenchmarkTests.java
index 0510186ea..63a6adc65 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/StartupBenchmark.java
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/mvc/StartupBenchmarkTests.java
@@ -14,24 +14,32 @@
* limitations under the License.
*/
-package org.springframework.cloud.sleuth.benchmarks.jmh.benchmarks;
+package org.springframework.cloud.sleuth.benchmarks.jmh.mvc;
+import jmh.mbr.junit5.Microbenchmark;
+import org.junit.jupiter.api.Disabled;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.BenchmarkMode;
import org.openjdk.jmh.annotations.Fork;
import org.openjdk.jmh.annotations.Level;
import org.openjdk.jmh.annotations.Measurement;
import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.Param;
import org.openjdk.jmh.annotations.Scope;
import org.openjdk.jmh.annotations.State;
import org.openjdk.jmh.annotations.TearDown;
import org.openjdk.jmh.annotations.Warmup;
+import org.springframework.cloud.sleuth.benchmarks.jmh.ProcessLauncherState;
+import org.springframework.cloud.sleuth.benchmarks.jmh.TracerImplementation;
+
@Measurement(iterations = 5)
@Warmup(iterations = 1)
@Fork(value = 2, warmups = 0)
@BenchmarkMode(Mode.AverageTime)
-public class StartupBenchmark {
+@Microbenchmark
+@Disabled("Process doesn't stop")
+public class StartupBenchmarkTests {
@Benchmark
public void withAnnotations(ApplicationState state) throws Exception {
@@ -46,31 +54,30 @@ public class StartupBenchmark {
@Benchmark
public void withoutAsync(ApplicationState state) throws Exception {
- state.setExtraArgs("--spring.sleuth.async.enabled=false",
- "--spring.sleuth.annotation.enabled=false");
+ state.setExtraArgs("--spring.sleuth.async.enabled=false", "--spring.sleuth.annotation.enabled=false");
state.run();
}
@Benchmark
public void withoutScheduled(ApplicationState state) throws Exception {
- state.setExtraArgs("--spring.sleuth.scheduled.enabled=false",
- "--spring.sleuth.async.enabled=false",
+ state.setExtraArgs("--spring.sleuth.scheduled.enabled=false", "--spring.sleuth.async.enabled=false",
"--spring.sleuth.annotation.enabled=false");
state.run();
}
@Benchmark
public void withoutWeb(ApplicationState state) throws Exception {
- state.setExtraArgs("--spring.sleuth.web.enabled=false",
- "--spring.sleuth.scheduled.enabled=false",
- "--spring.sleuth.async.enabled=false",
- "--spring.sleuth.annotation.enabled=false");
+ state.setExtraArgs("--spring.sleuth.web.enabled=false", "--spring.sleuth.scheduled.enabled=false",
+ "--spring.sleuth.async.enabled=false", "--spring.sleuth.annotation.enabled=false");
state.run();
}
@State(Scope.Benchmark)
public static class ApplicationState extends ProcessLauncherState {
+ @Param
+ private TracerImplementation tracerImplementation;
+
public ApplicationState() {
super("target", "--server.port=0");
}
diff --git a/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/stream/MicroBenchmarkStreamTests.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/stream/MicroBenchmarkStreamTests.java
new file mode 100644
index 000000000..d58dee940
--- /dev/null
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/stream/MicroBenchmarkStreamTests.java
@@ -0,0 +1,196 @@
+/*
+ * Copyright 2013-2020 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
+ *
+ * https://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.benchmarks.jmh.stream;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+import brave.Tracing;
+import jmh.mbr.junit5.Microbenchmark;
+import org.junit.platform.commons.annotation.Testable;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.Setup;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.TearDown;
+import org.openjdk.jmh.annotations.Warmup;
+
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.WebApplicationType;
+import org.springframework.boot.builder.SpringApplicationBuilder;
+import org.springframework.cloud.sleuth.benchmarks.app.stream.SleuthBenchmarkingStreamApplication;
+import org.springframework.cloud.sleuth.benchmarks.jmh.Pair;
+import org.springframework.cloud.sleuth.benchmarks.jmh.TracerImplementation;
+import org.springframework.cloud.stream.binder.test.InputDestination;
+import org.springframework.cloud.stream.binder.test.OutputDestination;
+import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
+import org.springframework.context.ConfigurableApplicationContext;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.context.annotation.Import;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.support.MessageBuilder;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+@Measurement(iterations = 10, time = 1)
+@Warmup(iterations = 10, time = 1)
+@Fork(4)
+@BenchmarkMode(Mode.SampleTime)
+@OutputTimeUnit(TimeUnit.MILLISECONDS)
+@Microbenchmark
+public class MicroBenchmarkStreamTests {
+
+ @Benchmark
+ @Testable
+ public void testStream(BenchmarkContext context) throws Exception {
+ context.run();
+ }
+
+ @State(Scope.Benchmark)
+ public static class BenchmarkContext {
+
+ volatile ConfigurableApplicationContext applicationContext;
+
+ volatile InputDestination input;
+
+ volatile OutputDestination output;
+
+ @Param
+ private Instrumentation instrumentation;
+
+ @Param
+ private TracerImplementation tracerImplementation;
+
+ @Setup
+ public void setup() {
+ this.applicationContext = initContext();
+ this.input = this.applicationContext.getBean(InputDestination.class);
+ this.output = this.applicationContext.getBean(OutputDestination.class);
+ }
+
+ private void sendInputMessage() {
+ // System.out.println("Sending the message to input");
+ input.send(MessageBuilder.withPayload("hello".getBytes())
+ .setHeader("b3", "4883117762eb9420-4883117762eb9420-1").build());
+ }
+
+ protected ConfigurableApplicationContext initContext() {
+ SpringApplication application = new SpringApplicationBuilder(SleuthBenchmarkingStreamApplication.class)
+ .web(WebApplicationType.NONE).application();
+ return application.run(runArgs());
+ }
+
+ protected String[] runArgs() {
+ List strings = new ArrayList<>();
+ strings.addAll(Arrays.asList("--spring.jmx.enabled=false",
+ "--spring.application.name=defaultTraceContextForStream" + instrumentation.name() + "_"
+ + tracerImplementation.name()));
+ strings.addAll(Arrays.asList(instrumentation.asParams()));
+ return strings.toArray(new String[0]);
+ }
+
+ void run() {
+ sendInputMessage();
+ assertThatOutputMessageGotReceived();
+ }
+
+ private void assertThatOutputMessageGotReceived() {
+ // System.out.println("Retrieving the message for tests");
+ Message message = output.receive(200L);
+ // System.out.println("Got the message from output");
+ assertThat(message).isNotNull();
+ // System.out.println("Message is not null");
+ assertThat(message.getPayload()).isEqualTo("HELLO".getBytes());
+ // System.out.println("Payload is HELLO");
+ if (!instrumentation.toString().toLowerCase().contains("nosleuth")) {
+ String b3 = message.getHeaders().get("b3", String.class);
+ // System.out.println("Checking the b3 header [" + b3 + "]");
+ assertThat(b3).isNotEmpty();
+ if (b3.startsWith("0000000000000000")) {
+ assertThat(b3).startsWith("00000000000000004883117762eb9420");
+ } else {
+ assertThat(b3).startsWith("4883117762eb9420");
+ }
+ }
+ }
+
+ @TearDown
+ public void clean() throws Exception {
+ Tracing current = Tracing.current();
+ if (current != null) {
+ current.close();
+ }
+ try {
+ this.applicationContext.close();
+ }
+ catch (Exception ig) {
+
+ }
+ }
+
+ public enum Instrumentation {
+
+ // @formatter:off
+ noSleuthSimple(Pair.noSleuth(), function("simple")),
+ sleuthSimpleOnHooks(function("simple")),
+ sleuthSimpleOnEach(function("simple"), Pair.noHook(), Pair.onEach()),
+ sleuthSimpleOnLast(function("simple"), Pair.noHook(), Pair.onLast()),
+ sleuthSimpleWithAroundOnHooks(function("simple_function_with_around")),
+ sleuthSimpleWithAroundOnEach(function("simple_function_with_around"), Pair.noHook(), Pair.onEach()),
+ sleuthSimpleWithAroundOnLast(function("simple_function_with_around"), Pair.noHook(), Pair.onLast()),
+ noSleuthReactiveSimple(function("reactive_simple"), Pair.noSleuth()),
+ sleuthReactiveSimpleOnHooks(function("DECORATE_ON_EACH")),
+ sleuthReactiveSimpleOnEach(function("DECORATE_ON_EACH"), Pair.noHook(), Pair.onEach(), integrationEnabled());
+ // @formatter:on
+
+ private List pairs;
+
+ Instrumentation(Pair... pairs) {
+ this.pairs = Arrays.asList(pairs);
+ }
+
+ String[] asParams() {
+ return this.pairs.stream().map(p -> "--" + p.asProp()).collect(Collectors.toList()).toArray(new String[0]);
+ }
+
+ static Pair function(String type) {
+ return Pair.of("spring.sleuth.function.type", type);
+ }
+
+ static Pair integrationEnabled() {
+ return Pair.of("spring.sleuth.integration.enabled", "true");
+ }
+ }
+
+ }
+
+ @Configuration(proxyBeanMethods = false)
+ @Import(TestChannelBinderConfiguration.class)
+ static class TestConfiguration {
+
+ }
+
+}
diff --git a/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/MicroBenchmarkHttpTests.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/MicroBenchmarkHttpTests.java
new file mode 100644
index 000000000..2b622e136
--- /dev/null
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/MicroBenchmarkHttpTests.java
@@ -0,0 +1,146 @@
+/*
+ * Copyright 2013-2020 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
+ *
+ * https://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.benchmarks.jmh.webflux;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+import brave.Tracing;
+import jmh.mbr.junit5.Microbenchmark;
+import org.junit.platform.commons.annotation.Testable;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.Setup;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.TearDown;
+import org.openjdk.jmh.annotations.Warmup;
+
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.WebApplicationType;
+import org.springframework.boot.builder.SpringApplicationBuilder;
+import org.springframework.cloud.sleuth.benchmarks.app.webflux.SleuthBenchmarkingSpringWebFluxApp;
+import org.springframework.cloud.sleuth.benchmarks.jmh.Pair;
+import org.springframework.cloud.sleuth.benchmarks.jmh.TracerImplementation;
+import org.springframework.context.ConfigurableApplicationContext;
+import org.springframework.test.web.reactive.server.WebTestClient;
+
+@Measurement(iterations = 10, time = 1)
+@Warmup(iterations = 10, time = 1)
+@Fork(4)
+@BenchmarkMode(Mode.SampleTime)
+@OutputTimeUnit(TimeUnit.MILLISECONDS)
+@Microbenchmark
+public class MicroBenchmarkHttpTests {
+
+ @Benchmark
+ @Testable
+ public void test(BenchmarkContext context) throws Exception {
+ context.run();
+ }
+
+ @State(Scope.Benchmark)
+ public static class BenchmarkContext {
+
+ volatile ConfigurableApplicationContext applicationContext;
+
+ volatile WebTestClient webTestClient;
+
+ @Param
+ private Instrumentation instrumentation;
+
+ @Param
+ private TracerImplementation tracerImplementation;
+
+ @Setup
+ public void setup() {
+ this.applicationContext = initContext();
+ this.webTestClient = WebTestClient.bindToApplicationContext(applicationContext).build();
+ }
+
+ protected ConfigurableApplicationContext initContext() {
+ SpringApplication application = new SpringApplicationBuilder(SleuthBenchmarkingSpringWebFluxApp.class)
+ .web(WebApplicationType.REACTIVE).application();
+ return application.run(runArgs());
+ }
+
+ protected String[] runArgs() {
+ String[] defaultArgs = new String[] { "--spring.jmx.enabled=false",
+ "--spring.application.name=defaultTraceContext" + instrumentation.name() + "_"
+ + tracerImplementation.name() };
+ List list = new ArrayList<>(Arrays.asList(defaultArgs));
+ list.addAll(Arrays.asList(instrumentation.asParams()));
+ return list.toArray(new String[0]);
+ }
+
+ void run() {
+ this.webTestClient.get().uri(instrumentation.url).header("X-B3-TraceId", "4883117762eb9420")
+ .header("X-B3-SpanId", "4883117762eb9420").exchange().expectStatus().isOk();
+ }
+
+ @TearDown
+ public void clean() throws Exception {
+ Tracing current = Tracing.current();
+ if (current != null) {
+ current.close();
+ }
+ try {
+ this.applicationContext.close();
+ }
+ catch (Exception ig) {
+
+ }
+ }
+
+ public enum Instrumentation {
+
+ // @formatter:off
+ noSleuthSimple("/simple", Pair.noSleuth()),
+ sleuthSimpleOnHooks("/simple"),
+ sleuthSimpleOnEach("/simple", Pair.noHook(), Pair.onEach()),
+ sleuthSimpleOnLast("/simple", Pair.noHook(), Pair.onLast()),
+ noSleuthComplex("/complexNoSleuth", Pair.noSleuth()),
+ onHooksComplex("/complex"),
+ onEachComplex("/complex", Pair.noHook(), Pair.onEach()),
+ onLastComplex("/complex", Pair.noHook(), Pair.onLast());
+ // @formatter:on
+
+ private String url;
+
+ private List pairs;
+
+ Instrumentation(String url, Pair... pairs) {
+ this.url = url;
+ this.pairs = Arrays.asList(pairs);
+ }
+
+ String[] asParams() {
+ return this.pairs.stream().map(p -> "--" + p.asProp()).collect(Collectors.toList()).toArray(new String[0]);
+ }
+ }
+
+ }
+
+}
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/SpringWebFluxBenchmarks.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/SpringWebFluxBenchmarksTests.java
similarity index 77%
rename from benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/SpringWebFluxBenchmarks.java
rename to benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/SpringWebFluxBenchmarksTests.java
index ca0cfd1cd..c11f45a15 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/SpringWebFluxBenchmarks.java
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/SpringWebFluxBenchmarksTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2019 the original author or authors.
+ * Copyright 2013-2020 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.
@@ -14,7 +14,7 @@
* limitations under the License.
*/
-package org.springframework.cloud.sleuth.benchmarks.jmh.benchmarks;
+package org.springframework.cloud.sleuth.benchmarks.jmh.webflux;
import java.io.IOException;
import java.util.concurrent.TimeUnit;
@@ -26,6 +26,7 @@ import brave.httpclient.TracingHttpClientBuilder;
import brave.propagation.CurrentTraceContext;
import brave.propagation.TraceContext;
import brave.sampler.Sampler;
+import jmh.mbr.junit5.Microbenchmark;
import org.apache.http.client.methods.HttpGet;
import org.apache.http.impl.client.CloseableHttpClient;
import org.apache.http.impl.client.HttpClients;
@@ -51,32 +52,40 @@ import org.springframework.boot.SpringApplication;
import org.springframework.boot.WebApplicationType;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.cloud.sleuth.benchmarks.app.webflux.SleuthBenchmarkingSpringWebFluxApp;
+import org.springframework.cloud.sleuth.benchmarks.jmh.TracerImplementation;
import org.springframework.context.ConfigurableApplicationContext;
-@Measurement(iterations = 5, time = 1)
+@Measurement(iterations = 10, time = 1)
@Warmup(iterations = 10, time = 1)
-@Fork(3)
+@Fork(4)
@BenchmarkMode(Mode.SampleTime)
@OutputTimeUnit(TimeUnit.MICROSECONDS)
@Threads(2)
@State(Scope.Benchmark)
-public class SpringWebFluxBenchmarks {
+@Microbenchmark
+public abstract class SpringWebFluxBenchmarksTests {
+
static final SpanHandler FAKE_SPAN_HANDLER = new SpanHandler() {
// intentionally anonymous to prevent logging fallback on NOOP
};
- protected static TraceContext defaultTraceContext = TraceContext.newBuilder()
- .traceIdHigh(333L).traceId(444L).spanId(3).sampled(true).build();
+ protected static TraceContext defaultTraceContext = TraceContext.newBuilder().traceIdHigh(333L).traceId(444L)
+ .spanId(3).sampled(true).build();
+
protected ConfigurableApplicationContext applicationContext;
+
protected SleuthBenchmarkingSpringWebFluxApp springWebFluxApp;
+
CloseableHttpClient client;
+
CloseableHttpClient tracedClient;
+
CloseableHttpClient unsampledClient;
+
private String baseUrl;
public static void main(String[] args) throws RunnerException {
- Options opt = new OptionsBuilder()
- .include(".*" + SpringWebFluxBenchmarks.class.getSimpleName() + ".*")
+ Options opt = new OptionsBuilder().include(".*" + SpringWebFluxBenchmarksTests.class.getSimpleName() + ".*")
.build();
new Runner(opt).run();
@@ -87,8 +96,7 @@ public class SpringWebFluxBenchmarks {
}
protected CloseableHttpClient newClient(HttpTracing httpTracing) {
- return TracingHttpClientBuilder.create(httpTracing).disableAutomaticRetries()
- .build();
+ return TracingHttpClientBuilder.create(httpTracing).disableAutomaticRetries().build();
}
protected CloseableHttpClient newClient() {
@@ -107,21 +115,18 @@ public class SpringWebFluxBenchmarks {
public void setup() {
ConfigurableApplicationContext context = initContext();
this.applicationContext = context;
- this.springWebFluxApp = this.applicationContext
- .getBean(SleuthBenchmarkingSpringWebFluxApp.class);
+ this.springWebFluxApp = this.applicationContext.getBean(SleuthBenchmarkingSpringWebFluxApp.class);
baseUrl = "http://127.0.0.1:" + springWebFluxApp.port + "/foo";
client = newClient();
- tracedClient = newClient(HttpTracing
- .create(Tracing.newBuilder().addSpanHandler(FAKE_SPAN_HANDLER).build()));
- unsampledClient = newClient(HttpTracing.create(Tracing.newBuilder()
- .sampler(Sampler.NEVER_SAMPLE).addSpanHandler(FAKE_SPAN_HANDLER).build()));
+ tracedClient = newClient(HttpTracing.create(Tracing.newBuilder().addSpanHandler(FAKE_SPAN_HANDLER).build()));
+ unsampledClient = newClient(HttpTracing
+ .create(Tracing.newBuilder().sampler(Sampler.NEVER_SAMPLE).addSpanHandler(FAKE_SPAN_HANDLER).build()));
postSetUp();
}
protected ConfigurableApplicationContext initContext() {
- SpringApplication application = new SpringApplicationBuilder(
- SleuthBenchmarkingSpringWebFluxApp.class).web(WebApplicationType.REACTIVE)
- .application();
+ SpringApplication application = new SpringApplicationBuilder(SleuthBenchmarkingSpringWebFluxApp.class)
+ .web(WebApplicationType.REACTIVE).application();
customSpringApplication(application);
return application.run(runArgs());
}
@@ -134,9 +139,8 @@ public class SpringWebFluxBenchmarks {
}
protected String[] runArgs() {
- return new String[] { "--spring.jmx.enabled=false",
- "--spring.application.name=defaultTraceContext",
- "--spring.sleuth.enabled=true" };
+ return new String[] { "--spring.jmx.enabled=false", "--spring.application.name=defaultTraceContext",
+ TracerImplementation.brave.toString(), "--spring.sleuth.enabled=true" };
}
@TearDown
@@ -171,8 +175,7 @@ public class SpringWebFluxBenchmarks {
@Benchmark
public void tracedClient_get_resumeTrace() throws Exception {
- try (CurrentTraceContext.Scope scope = Tracing.current().currentTraceContext()
- .newScope(defaultTraceContext)) {
+ try (CurrentTraceContext.Scope scope = Tracing.current().currentTraceContext().newScope(defaultTraceContext)) {
get(tracedClient);
}
}
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/WithOutSleuthSpringWebFluxBenchmarks.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/WithOutReactorSleuthSpringWebFluxBenchmarksTests.java
similarity index 57%
rename from benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/WithOutSleuthSpringWebFluxBenchmarks.java
rename to benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/WithOutReactorSleuthSpringWebFluxBenchmarksTests.java
index 6780eba2c..3629040a1 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/WithOutSleuthSpringWebFluxBenchmarks.java
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/WithOutReactorSleuthSpringWebFluxBenchmarksTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2019 the original author or authors.
+ * Copyright 2013-2020 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.
@@ -14,31 +14,36 @@
* limitations under the License.
*/
-package org.springframework.cloud.sleuth.benchmarks.jmh.benchmarks;
+package org.springframework.cloud.sleuth.benchmarks.jmh.webflux;
+import jmh.mbr.junit5.Microbenchmark;
import org.openjdk.jmh.runner.Runner;
import org.openjdk.jmh.runner.RunnerException;
import org.openjdk.jmh.runner.options.Options;
import org.openjdk.jmh.runner.options.OptionsBuilder;
+import org.springframework.cloud.sleuth.benchmarks.jmh.TracerImplementation;
+
/**
* @author alvin
*/
-public class WithOutSleuthSpringWebFluxBenchmarks extends SpringWebFluxBenchmarks {
+@Microbenchmark
+public class WithOutReactorSleuthSpringWebFluxBenchmarksTests extends SpringWebFluxBenchmarksTests {
public static void main(String[] args) throws RunnerException {
- Options opt = new OptionsBuilder().include(
- ".*" + WithOutSleuthSpringWebFluxBenchmarks.class.getSimpleName() + ".*")
- .build();
+ Options opt = new OptionsBuilder()
+ .include(".*" + WithOutReactorSleuthSpringWebFluxBenchmarksTests.class.getSimpleName() + ".*").build();
new Runner(opt).run();
}
@Override
protected String[] runArgs() {
- return new String[] { "--spring.jmx.enabled=false",
- "--spring.application.name=defaultTraceContext",
- "--spring.sleuth.enabled=false" };
+ return new String[] { "--spring.jmx.enabled=false", "--spring.application.name=defaultTraceContext",
+ TracerImplementation.brave.toString(), "--spring.sleuth.enabled=true",
+ "--spring.sleuth.reactor.enabled=false"
+
+ };
}
@Override
diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/WithOutReactorSleuthSpringWebFluxBenchmarks.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/WithOutSleuthSpringWebFluxBenchmarksTests.java
similarity index 59%
rename from benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/WithOutReactorSleuthSpringWebFluxBenchmarks.java
rename to benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/WithOutSleuthSpringWebFluxBenchmarksTests.java
index 111d86c15..8185d0c20 100644
--- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/jmh/benchmarks/WithOutReactorSleuthSpringWebFluxBenchmarks.java
+++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/webflux/WithOutSleuthSpringWebFluxBenchmarksTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2019 the original author or authors.
+ * Copyright 2013-2020 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.
@@ -14,34 +14,33 @@
* limitations under the License.
*/
-package org.springframework.cloud.sleuth.benchmarks.jmh.benchmarks;
+package org.springframework.cloud.sleuth.benchmarks.jmh.webflux;
+import jmh.mbr.junit5.Microbenchmark;
import org.openjdk.jmh.runner.Runner;
import org.openjdk.jmh.runner.RunnerException;
import org.openjdk.jmh.runner.options.Options;
import org.openjdk.jmh.runner.options.OptionsBuilder;
+import org.springframework.cloud.sleuth.benchmarks.jmh.TracerImplementation;
+
/**
* @author alvin
*/
-public class WithOutReactorSleuthSpringWebFluxBenchmarks extends SpringWebFluxBenchmarks {
+@Microbenchmark
+public class WithOutSleuthSpringWebFluxBenchmarksTests extends SpringWebFluxBenchmarksTests {
public static void main(String[] args) throws RunnerException {
- Options opt = new OptionsBuilder().include(
- ".*" + WithOutReactorSleuthSpringWebFluxBenchmarks.class.getSimpleName()
- + ".*")
- .build();
+ Options opt = new OptionsBuilder()
+ .include(".*" + WithOutSleuthSpringWebFluxBenchmarksTests.class.getSimpleName() + ".*").build();
new Runner(opt).run();
}
@Override
protected String[] runArgs() {
- return new String[] { "--spring.jmx.enabled=false",
- "--spring.application.name=defaultTraceContext",
- "--spring.sleuth.enabled=true", "--spring.sleuth.reactor.enabled=false"
-
- };
+ return new String[] { "--spring.jmx.enabled=false", "--spring.application.name=defaultTraceContext",
+ TracerImplementation.brave.toString(), "--spring.sleuth.enabled=false" };
}
@Override
diff --git a/docs/src/main/asciidoc/spring-cloud-sleuth.adoc b/docs/src/main/asciidoc/spring-cloud-sleuth.adoc
index ad529b01a..f207bf515 100644
--- a/docs/src/main/asciidoc/spring-cloud-sleuth.adoc
+++ b/docs/src/main/asciidoc/spring-cloud-sleuth.adoc
@@ -1526,19 +1526,14 @@ To turn off this feature, set the `spring.sleuth.quartz.enabled` property to `fa
=== Project Reactor
+==== From Spring Cloud Sleuth 2.2.8 (inclusive)
+
+With the new Reactor https://github.com/reactor/reactor-core/pull/2566[queue wrapping mechanism] (Reactor 3.3.14) we're instrumenting the way threads are switched by Reactor. You should observe significant improvement in performance. In order to disable this feature you have to set the `spring.sleuth.reactor.decorate-hooks` option to `false`. You'll fall back to the previous instrumentation mode mechanism.
+
+==== To Spring Cloud Sleuth 2.2.8 (exclusive)
+
For projects depending on Project Reactor such as Spring Cloud Gateway, we suggest turning the `spring.sleuth.reactor.decorate-on-each` option to `false`. That way an increased performance gain should be observed in comparison to the standard instrumentation mechanism. What this option does is it will wrap decorate `onLast` operator instead of `onEach` which will result in creation of far fewer objects. The downside of this is that when Project Reactor will change threads, the trace propagation will continue without issues, however anything relying on the `ThreadLocal` such as e.g. MDC entries can be buggy.
== Configuration properties
To see the list of all Sleuth related configuration properties please check link:appendix.html[the Appendix page].
-
-== Running examples
-
-You can see the running examples deployed in the https://run.pivotal.io/[Pivotal Web Services].
-Check them out at the following links:
-
-* https://docssleuth-zipkin-server.cfapps.io/[Zipkin for apps presented in the samples to the top]. First make
-a request to https://docssleuth-service1.cfapps.io/start[Service 1] and then check out the trace in Zipkin.
-* https://docsbrewing-zipkin-server.cfapps.io/[Zipkin for Brewery on PWS], its https://github.com/spring-cloud-samples/brewery[Github Code].
-Ensure that you've picked the lookback period of 7 days. If there are no traces, go to https://docsbrewing-presenting.cfapps.io/[Presenting application]
-and order some beers. Then check Zipkin for traces.
diff --git a/pom.xml b/pom.xml
index 62f0b08af..4d018f908 100644
--- a/pom.xml
+++ b/pom.xml
@@ -29,7 +29,7 @@
org.springframework.cloud
spring-cloud-build
- 2.3.2.RELEASE
+ 2.3.3.BUILD-SNAPSHOT
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java
index 345bfbb8a..690685c4b 100644
--- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java
@@ -129,6 +129,29 @@ public abstract class ReactorSleuth {
});
}
+ public static Function scopePassingOnScheduleHook(
+ ConfigurableApplicationContext springContext) {
+ LazyBean lazyCurrentTraceContext = LazyBean
+ .create(springContext, CurrentTraceContext.class);
+ return delegate -> {
+ if (springContext.isActive()) {
+ final CurrentTraceContext currentTraceContext = lazyCurrentTraceContext
+ .get();
+ if (currentTraceContext == null) {
+ return delegate;
+ }
+ final TraceContext traceContext = currentTraceContext.get();
+ return () -> {
+ try (CurrentTraceContext.Scope scope = currentTraceContext
+ .maybeScope(traceContext)) {
+ delegate.run();
+ }
+ };
+ }
+ return delegate;
+ };
+ }
+
private static Context context(CoreSubscriber super T> sub) {
try {
return sub.currentContext();
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SleuthReactorProperties.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SleuthReactorProperties.java
index fa1d433d1..4a7315164 100644
--- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SleuthReactorProperties.java
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SleuthReactorProperties.java
@@ -33,11 +33,21 @@ public class SleuthReactorProperties {
*/
private boolean enabled = true;
+ /**
+ * When true uses the new decorate hooks feature from Project Reactor. Should allow
+ * the feature set of {@link SleuthReactorProperties#decorateOnEach} with the least
+ * impact on the performance.
+ */
+ private boolean decorateHooks = true;
+
/**
* When true decorates on each operator, will be less performing, but logging will
* always contain the tracing entries in each operator. When false decorates on last
* operator, will be more performing, but logging might not always contain the tracing
* entries.
+ *
+ * If {@link SleuthReactorProperties#decorateHooks} is used, this decoration mode will
+ * NOT be used.
*/
private boolean decorateOnEach = true;
@@ -49,6 +59,14 @@ public class SleuthReactorProperties {
this.enabled = enabled;
}
+ public boolean isDecorateHooks() {
+ return this.decorateHooks;
+ }
+
+ public void setDecorateHooks(boolean decorateHooks) {
+ this.decorateHooks = decorateHooks;
+ }
+
public boolean isDecorateOnEach() {
return this.decorateOnEach;
}
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java
index 87e0eb510..1a6ee4252 100644
--- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java
@@ -16,9 +16,18 @@
package org.springframework.cloud.sleuth.instrument.reactor;
+import java.io.Closeable;
+import java.io.IOException;
+import java.util.AbstractQueue;
+import java.util.Iterator;
+import java.util.Queue;
+import java.util.function.Function;
+
import javax.annotation.PreDestroy;
import brave.Tracing;
+import brave.propagation.CurrentTraceContext;
+import brave.propagation.TraceContext;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import reactor.core.publisher.Hooks;
@@ -37,15 +46,16 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.cloud.context.scope.refresh.RefreshScope;
import org.springframework.cloud.context.scope.refresh.RefreshScopeRefreshedEvent;
-import org.springframework.cloud.sleuth.instrument.async.TraceableScheduledExecutorService;
import org.springframework.cloud.sleuth.instrument.web.TraceWebFluxAutoConfiguration;
import org.springframework.context.ApplicationListener;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.env.ConfigurableEnvironment;
+import org.springframework.util.ReflectionUtils;
import static org.springframework.cloud.sleuth.instrument.reactor.ReactorSleuth.scopePassingSpanOperator;
+import static org.springframework.cloud.sleuth.instrument.reactor.TraceReactorAutoConfiguration.SLEUTH_REACTOR_EXECUTOR_SERVICE_KEY;
import static org.springframework.cloud.sleuth.instrument.reactor.TraceReactorAutoConfiguration.TraceReactorConfiguration.SLEUTH_TRACE_REACTOR_KEY;
/**
@@ -77,6 +87,8 @@ public class TraceReactorAutoConfiguration {
private static final Log log = LogFactory.getLog(TraceReactorConfiguration.class);
+ static final boolean IS_QUEUE_WRAPPER_ON_THE_CLASSPATH = isQueueWrapperOnTheClasspath();
+
@Autowired
ConfigurableApplicationContext springContext;
@@ -87,6 +99,13 @@ public class TraceReactorAutoConfiguration {
}
SleuthReactorProperties reactorProperties = this.springContext
.getBean(SleuthReactorProperties.class);
+ if (TraceReactorAutoConfiguration.TraceReactorConfiguration.IS_QUEUE_WRAPPER_ON_THE_CLASSPATH
+ && reactorProperties.isDecorateHooks()) {
+ if (log.isTraceEnabled()) {
+ log.trace("Resetting queue wrapper instrumentation");
+ }
+ Hooks.removeQueueWrapper(SLEUTH_TRACE_REACTOR_KEY);
+ }
if (reactorProperties.isDecorateOnEach()) {
if (log.isTraceEnabled()) {
log.trace("Resetting onEach operator instrumentation");
@@ -99,8 +118,11 @@ public class TraceReactorAutoConfiguration {
}
Hooks.resetOnLastOperator(SLEUTH_TRACE_REACTOR_KEY);
}
- Schedulers
- .removeExecutorServiceDecorator(SLEUTH_REACTOR_EXECUTOR_SERVICE_KEY);
+ }
+
+ private static boolean isQueueWrapperOnTheClasspath() {
+ return ReflectionUtils.findMethod(Hooks.class, "addQueueWrapper",
+ String.class, Function.class) != null;
}
@Bean
@@ -152,12 +174,23 @@ class HooksRefresher implements ApplicationListener
}
Hooks.resetOnEachOperator(SLEUTH_TRACE_REACTOR_KEY);
Hooks.resetOnLastOperator(SLEUTH_TRACE_REACTOR_KEY);
- if (this.reactorProperties.isDecorateOnEach()) {
+ Hooks.removeQueueWrapper(SLEUTH_TRACE_REACTOR_KEY);
+ if (this.reactorProperties.isDecorateHooks()
+ && TraceReactorAutoConfiguration.TraceReactorConfiguration.IS_QUEUE_WRAPPER_ON_THE_CLASSPATH) {
+ if (log.isTraceEnabled()) {
+ log.trace("Adding queue wrapper instrumentation");
+ }
+ HookRegisteringBeanDefinitionRegistryPostProcessor.addQueueWrapper(context);
+ }
+ else if (this.reactorProperties.isDecorateOnEach()) {
if (log.isTraceEnabled()) {
log.trace("Decorating onEach operator instrumentation");
}
Hooks.onEachOperator(SLEUTH_TRACE_REACTOR_KEY,
scopePassingSpanOperator(this.context));
+ Schedulers.onScheduleHook(
+ TraceReactorAutoConfiguration.SLEUTH_REACTOR_EXECUTOR_SERVICE_KEY,
+ ReactorSleuth.scopePassingOnScheduleHook(this.context));
}
else {
if (log.isTraceEnabled()) {
@@ -171,7 +204,7 @@ class HooksRefresher implements ApplicationListener
}
class HookRegisteringBeanDefinitionRegistryPostProcessor
- implements BeanDefinitionRegistryPostProcessor {
+ implements BeanDefinitionRegistryPostProcessor, Closeable {
private static final Log log = LogFactory
.getLog(HookRegisteringBeanDefinitionRegistryPostProcessor.class);
@@ -194,27 +227,165 @@ class HookRegisteringBeanDefinitionRegistryPostProcessor
static void setupHooks(ConfigurableApplicationContext springContext) {
ConfigurableEnvironment environment = springContext.getEnvironment();
- boolean decorateOnEach = environment.getProperty(
- "spring.sleuth.reactor.decorate-on-each", Boolean.class, true);
- if (decorateOnEach) {
- if (log.isTraceEnabled()) {
- log.trace("Decorating onEach operator instrumentation");
- }
- Hooks.onEachOperator(SLEUTH_TRACE_REACTOR_KEY,
- scopePassingSpanOperator(springContext));
+ Boolean decorateHooks = environment
+ .getProperty("spring.sleuth.reactor.decorate-hooks", Boolean.class);
+ if (wrapperNotOnClasspathButPropertyHasValue(decorateHooks)) {
+ log.warn(
+ "You have explicitly set the decorate hooks option but you're using an old version of Reactor. Please upgrade to the latest Boot version (at least 2.3.9.RELEASE). Will fall back to the previous reactor instrumentation mode");
}
else {
- if (log.isTraceEnabled()) {
- log.trace("Decorating onLast operator instrumentation");
- }
- Hooks.onLastOperator(SLEUTH_TRACE_REACTOR_KEY,
- scopePassingSpanOperator(springContext));
+ decorateHooks = decorateHooks != null ? decorateHooks : Boolean.TRUE;
}
- Schedulers.setExecutorServiceDecorator(
+ if (wrapperOnClasspathHooksPropertyTurnedOn(decorateHooks)) {
+ if (log.isTraceEnabled()) {
+ log.trace("Adding queue wrapper instrumentation");
+ }
+ addQueueWrapper(springContext);
+ }
+ else {
+ boolean decorateOnEach = environment.getProperty(
+ "spring.sleuth.reactor.decorate-on-each", Boolean.class, true);
+ if (decorateOnEach) {
+ if (log.isTraceEnabled()) {
+ log.trace("Decorating onEach operator instrumentation");
+ }
+ Hooks.onEachOperator(SLEUTH_TRACE_REACTOR_KEY,
+ scopePassingSpanOperator(springContext));
+ }
+ else {
+ if (log.isTraceEnabled()) {
+ log.trace("Decorating onLast operator instrumentation");
+ }
+ Hooks.onLastOperator(SLEUTH_TRACE_REACTOR_KEY,
+ scopePassingSpanOperator(springContext));
+ }
+ }
+ decorateScheduler(springContext);
+ }
+
+ private static boolean wrapperOnClasspathHooksPropertyTurnedOn(Boolean decorateHooks) {
+ return Boolean.TRUE.equals(decorateHooks)
+ && TraceReactorAutoConfiguration.TraceReactorConfiguration.IS_QUEUE_WRAPPER_ON_THE_CLASSPATH;
+ }
+
+ private static boolean wrapperNotOnClasspathButPropertyHasValue(Boolean decorateHooks) {
+ return !TraceReactorAutoConfiguration.TraceReactorConfiguration.IS_QUEUE_WRAPPER_ON_THE_CLASSPATH
+ && decorateHooks != null;
+ }
+
+ static void addQueueWrapper(ConfigurableApplicationContext springContext) {
+ Hooks.addQueueWrapper(SLEUTH_TRACE_REACTOR_KEY,
+ queue -> traceQueue(springContext, queue));
+ }
+
+ @Override
+ public void close() throws IOException {
+ if (log.isTraceEnabled()) {
+ log.trace("Cleaning up hooks");
+ }
+ Hooks.resetOnEachOperator(SLEUTH_TRACE_REACTOR_KEY);
+ Hooks.resetOnLastOperator(SLEUTH_TRACE_REACTOR_KEY);
+ Hooks.removeQueueWrapper(SLEUTH_REACTOR_EXECUTOR_SERVICE_KEY);
+ Schedulers.resetOnScheduleHook(SLEUTH_REACTOR_EXECUTOR_SERVICE_KEY);
+ Schedulers.resetOnScheduleHook(
+ TraceReactorAutoConfiguration.SLEUTH_REACTOR_EXECUTOR_SERVICE_KEY);
+ }
+
+ private static void decorateScheduler(ConfigurableApplicationContext springContext) {
+ Schedulers.onScheduleHook(
TraceReactorAutoConfiguration.SLEUTH_REACTOR_EXECUTOR_SERVICE_KEY,
- (scheduler,
- scheduledExecutorService) -> new TraceableScheduledExecutorService(
- springContext, scheduledExecutorService));
+ ReactorSleuth.scopePassingOnScheduleHook(springContext));
+ }
+
+ private static Queue> traceQueue(ConfigurableApplicationContext springContext,
+ Queue> queue) {
+ if (!springContext.isActive()) {
+ return queue;
+ }
+ CurrentTraceContext currentTraceContext = springContext
+ .getBean(CurrentTraceContext.class);
+ @SuppressWarnings("unchecked")
+ Queue envelopeQueue = queue;
+ return new AbstractQueue