Sets name of the executor service bean for all created runnables / callables

without this change the names are 'async'
with this change the name will be the name of the executor service bean. It will be easier to reason about the origin of such a span

fixes gh-1465
This commit is contained in:
Marcin Grzejszczak
2020-08-18 13:42:28 +02:00
parent 391059f824
commit fc8efb5f5d
12 changed files with 220 additions and 209 deletions

View File

@@ -82,7 +82,7 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor {
}
if (bean instanceof ThreadPoolTaskExecutor) {
if (isProxyNeeded(beanName)) {
return wrapThreadPoolTaskExecutor(bean);
return wrapThreadPoolTaskExecutor(bean, beanName);
}
else {
log.info("Not instrumenting bean " + beanName);
@@ -90,7 +90,7 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor {
}
else if (bean instanceof ScheduledExecutorService) {
if (isProxyNeeded(beanName)) {
return wrapScheduledExecutorService(bean);
return wrapScheduledExecutorService(bean, beanName);
}
else {
log.info("Not instrumenting bean " + beanName);
@@ -98,7 +98,7 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor {
}
else if (bean instanceof ExecutorService) {
if (isProxyNeeded(beanName)) {
return wrapExecutorService(bean);
return wrapExecutorService(bean, beanName);
}
else {
log.info("Not instrumenting bean " + beanName);
@@ -106,26 +106,26 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor {
}
else if (bean instanceof AsyncTaskExecutor) {
if (isProxyNeeded(beanName)) {
return wrapAsyncTaskExecutor(bean);
return wrapAsyncTaskExecutor(bean, beanName);
}
else {
log.info("Not instrumenting bean " + beanName);
}
}
else if (bean instanceof Executor) {
return wrapExecutor(bean);
return wrapExecutor(bean, beanName);
}
return bean;
}
private Object wrapExecutor(Object bean) {
private Object wrapExecutor(Object bean, String beanName) {
Executor executor = (Executor) bean;
boolean methodFinal = anyFinalMethods(executor);
boolean classFinal = Modifier.isFinal(bean.getClass().getModifiers());
boolean cglibProxy = !methodFinal && !classFinal;
try {
return createProxy(bean, cglibProxy,
new ExecutorMethodInterceptor<>(executor, this.beanFactory));
return createProxy(bean, cglibProxy, new ExecutorMethodInterceptor<>(executor,
this.beanFactory, beanName));
}
catch (AopConfigException ex) {
if (cglibProxy) {
@@ -134,43 +134,43 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor {
"Exception occurred while trying to create a proxy, falling back to JDK proxy",
ex);
}
return createProxy(bean, false,
new ExecutorMethodInterceptor<>(executor, this.beanFactory));
return createProxy(bean, false, new ExecutorMethodInterceptor<>(executor,
this.beanFactory, beanName));
}
throw ex;
}
}
private Object wrapThreadPoolTaskExecutor(Object bean) {
private Object wrapThreadPoolTaskExecutor(Object bean, String beanName) {
ThreadPoolTaskExecutor executor = (ThreadPoolTaskExecutor) bean;
boolean classFinal = Modifier.isFinal(bean.getClass().getModifiers());
boolean methodsFinal = anyFinalMethods(executor);
boolean cglibProxy = !classFinal && !methodsFinal;
return createThreadPoolTaskExecutorProxy(bean, cglibProxy, executor);
return createThreadPoolTaskExecutorProxy(bean, cglibProxy, executor, beanName);
}
private Object wrapExecutorService(Object bean) {
private Object wrapExecutorService(Object bean, String beanName) {
ExecutorService executor = (ExecutorService) bean;
boolean classFinal = Modifier.isFinal(bean.getClass().getModifiers());
boolean methodFinal = anyFinalMethods(executor);
boolean cglibProxy = !classFinal && !methodFinal;
return createExecutorServiceProxy(bean, cglibProxy, executor);
return createExecutorServiceProxy(bean, cglibProxy, executor, beanName);
}
private Object wrapScheduledExecutorService(Object bean) {
private Object wrapScheduledExecutorService(Object bean, String beanName) {
ScheduledExecutorService executor = (ScheduledExecutorService) bean;
boolean classFinal = Modifier.isFinal(bean.getClass().getModifiers());
boolean methodFinal = anyFinalMethods(executor);
boolean cglibProxy = !classFinal && !methodFinal;
return createScheduledExecutorServiceProxy(bean, cglibProxy, executor);
return createScheduledExecutorServiceProxy(bean, cglibProxy, executor, beanName);
}
private Object wrapAsyncTaskExecutor(Object bean) {
private Object wrapAsyncTaskExecutor(Object bean, String beanName) {
AsyncTaskExecutor executor = (AsyncTaskExecutor) bean;
boolean classFinal = Modifier.isFinal(bean.getClass().getModifiers());
boolean methodsFinal = anyFinalMethods(executor);
boolean cglibProxy = !classFinal && !methodsFinal;
return createAsyncTaskExecutorProxy(bean, cglibProxy, executor);
return createAsyncTaskExecutorProxy(bean, cglibProxy, executor, beanName);
}
boolean isProxyNeeded(String beanName) {
@@ -179,57 +179,62 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor {
}
Object createThreadPoolTaskExecutorProxy(Object bean, boolean cglibProxy,
ThreadPoolTaskExecutor executor) {
ThreadPoolTaskExecutor executor, String beanName) {
if (!cglibProxy) {
return new LazyTraceThreadPoolTaskExecutor(this.beanFactory, executor);
return new LazyTraceThreadPoolTaskExecutor(this.beanFactory, executor,
beanName);
}
return getProxiedObject(bean, cglibProxy, executor,
() -> new LazyTraceThreadPoolTaskExecutor(this.beanFactory, executor));
return getProxiedObject(bean, beanName, true, executor,
() -> new LazyTraceThreadPoolTaskExecutor(this.beanFactory, executor,
beanName));
}
Supplier<Executor> createThreadPoolTaskSchedulerProxy(
ThreadPoolTaskScheduler executor) {
return () -> new LazyTraceThreadPoolTaskScheduler(this.beanFactory, executor);
ThreadPoolTaskScheduler executor, String beanName) {
return () -> new LazyTraceThreadPoolTaskScheduler(this.beanFactory, executor,
beanName);
}
Supplier<Executor> createScheduledThreadPoolExecutorProxy(
ScheduledThreadPoolExecutor executor) {
ScheduledThreadPoolExecutor executor, String beanName) {
return () -> new LazyTraceScheduledThreadPoolExecutor(executor.getCorePoolSize(),
executor.getThreadFactory(), executor.getRejectedExecutionHandler(),
this.beanFactory, executor);
this.beanFactory, executor, beanName);
}
Object createExecutorServiceProxy(Object bean, boolean cglibProxy,
ExecutorService executor) {
return getProxiedObject(bean, cglibProxy, executor, () -> {
ExecutorService executor, String beanName) {
return getProxiedObject(bean, beanName, cglibProxy, executor, () -> {
if (executor instanceof ScheduledExecutorService) {
return new TraceableScheduledExecutorService(this.beanFactory, executor);
return new TraceableScheduledExecutorService(this.beanFactory, executor,
beanName);
}
return new TraceableExecutorService(this.beanFactory, executor);
return new TraceableExecutorService(this.beanFactory, executor, beanName);
});
}
Object createScheduledExecutorServiceProxy(Object bean, boolean cglibProxy,
ScheduledExecutorService executor) {
return getProxiedObject(bean, cglibProxy, executor,
() -> new TraceableScheduledExecutorService(this.beanFactory, executor));
ScheduledExecutorService executor, String beanName) {
return getProxiedObject(bean, beanName, cglibProxy, executor,
() -> new TraceableScheduledExecutorService(this.beanFactory, executor,
beanName));
}
Object createAsyncTaskExecutorProxy(Object bean, boolean cglibProxy,
AsyncTaskExecutor executor) {
return getProxiedObject(bean, cglibProxy, executor, () -> {
AsyncTaskExecutor executor, String beanName) {
return getProxiedObject(bean, beanName, cglibProxy, executor, () -> {
if (bean instanceof ThreadPoolTaskScheduler) {
return new LazyTraceThreadPoolTaskScheduler(this.beanFactory,
(ThreadPoolTaskScheduler) executor);
(ThreadPoolTaskScheduler) executor, beanName);
}
return new LazyTraceAsyncTaskExecutor(this.beanFactory, executor);
return new LazyTraceAsyncTaskExecutor(this.beanFactory, executor, beanName);
});
}
private Object getProxiedObject(Object bean, boolean cglibProxy, Executor executor,
Supplier<Executor> supplier) {
ProxyFactoryBean factory = proxyFactoryBean(bean, cglibProxy, executor, supplier);
private Object getProxiedObject(Object bean, String beanName, boolean cglibProxy,
Executor executor, Supplier<Executor> supplier) {
ProxyFactoryBean factory = proxyFactoryBean(bean, beanName, cglibProxy, executor,
supplier);
try {
return getObject(factory);
}
@@ -246,7 +251,7 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor {
"Will wrap ThreadPoolTaskScheduler in its tracing representation due to previous errors");
}
return createThreadPoolTaskSchedulerProxy(
(ThreadPoolTaskScheduler) bean).get();
(ThreadPoolTaskScheduler) bean, beanName).get();
}
else if (bean instanceof ScheduledThreadPoolExecutor) {
if (log.isDebugEnabled()) {
@@ -254,7 +259,7 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor {
"Will wrap ScheduledThreadPoolExecutor in its tracing representation due to previous errors");
}
return createScheduledThreadPoolExecutorProxy(
(ScheduledThreadPoolExecutor) bean).get();
(ScheduledThreadPoolExecutor) bean, beanName).get();
}
}
catch (Exception ex2) {
@@ -268,17 +273,18 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor {
}
}
private ProxyFactoryBean proxyFactoryBean(Object bean, boolean cglibProxy,
Executor executor, Supplier<Executor> supplier) {
private ProxyFactoryBean proxyFactoryBean(Object bean, String beanName,
boolean cglibProxy, Executor executor, Supplier<Executor> supplier) {
ProxyFactoryBean factory = new ProxyFactoryBean();
factory.setProxyTargetClass(cglibProxy);
factory.addAdvice(
new ExecutorMethodInterceptor<Executor>(executor, this.beanFactory) {
@Override
Executor executor(BeanFactory beanFactory, Executor executor) {
return supplier.get();
}
});
factory.addAdvice(new ExecutorMethodInterceptor<Executor>(executor,
this.beanFactory, beanName) {
@Override
Executor executor(BeanFactory beanFactory, Executor executor,
String beanName) {
return supplier.get();
}
});
factory.setTarget(bean);
return factory;
}
@@ -342,14 +348,17 @@ class ExecutorMethodInterceptor<T extends Executor> implements MethodInterceptor
private final BeanFactory beanFactory;
ExecutorMethodInterceptor(T delegate, BeanFactory beanFactory) {
private final String beanName;
ExecutorMethodInterceptor(T delegate, BeanFactory beanFactory, String beanName) {
this.delegate = delegate;
this.beanFactory = beanFactory;
this.beanName = beanName;
}
@Override
public Object invoke(MethodInvocation invocation) throws Throwable {
T executor = executor(this.beanFactory, this.delegate);
T executor = executor(this.beanFactory, this.delegate, this.beanName);
Method methodOnTracedBean = getMethod(invocation, executor);
if (methodOnTracedBean != null) {
try {
@@ -370,8 +379,8 @@ class ExecutorMethodInterceptor<T extends Executor> implements MethodInterceptor
method.getParameterTypes());
}
T executor(BeanFactory beanFactory, T executor) {
return (T) new LazyTraceExecutor(beanFactory, executor);
T executor(BeanFactory beanFactory, T executor, String beanName) {
return (T) new LazyTraceExecutor(beanFactory, executor, beanName);
}
}

View File

@@ -45,6 +45,8 @@ public class LazyTraceAsyncTaskExecutor implements AsyncTaskExecutor {
private final AsyncTaskExecutor delegate;
private final String beanName;
private Tracing tracing;
private SpanNamer spanNamer;
@@ -53,13 +55,21 @@ public class LazyTraceAsyncTaskExecutor implements AsyncTaskExecutor {
AsyncTaskExecutor delegate) {
this.beanFactory = beanFactory;
this.delegate = delegate;
this.beanName = null;
}
public LazyTraceAsyncTaskExecutor(BeanFactory beanFactory, AsyncTaskExecutor delegate,
String beanName) {
this.beanFactory = beanFactory;
this.delegate = delegate;
this.beanName = beanName;
}
@Override
public void execute(Runnable task) {
Runnable taskToRun = task;
if (!ContextUtil.isContextUnusable(this.beanFactory)) {
taskToRun = new TraceRunnable(tracing(), spanNamer(), task);
taskToRun = new TraceRunnable(tracing(), spanNamer(), task, this.beanName);
}
this.delegate.execute(taskToRun);
}
@@ -68,7 +78,7 @@ public class LazyTraceAsyncTaskExecutor implements AsyncTaskExecutor {
public void execute(Runnable task, long startTimeout) {
Runnable taskToRun = task;
if (!ContextUtil.isContextUnusable(this.beanFactory)) {
taskToRun = new TraceRunnable(tracing(), spanNamer(), task);
taskToRun = new TraceRunnable(tracing(), spanNamer(), task, this.beanName);
}
this.delegate.execute(taskToRun, startTimeout);
}
@@ -77,7 +87,7 @@ public class LazyTraceAsyncTaskExecutor implements AsyncTaskExecutor {
public Future<?> submit(Runnable task) {
Runnable taskToRun = task;
if (!ContextUtil.isContextUnusable(this.beanFactory)) {
taskToRun = new TraceRunnable(tracing(), spanNamer(), task);
taskToRun = new TraceRunnable(tracing(), spanNamer(), task, this.beanName);
}
return this.delegate.submit(taskToRun);
}
@@ -86,7 +96,7 @@ public class LazyTraceAsyncTaskExecutor implements AsyncTaskExecutor {
public <T> Future<T> submit(Callable<T> task) {
Callable<T> taskToRun = task;
if (!ContextUtil.isContextUnusable(this.beanFactory)) {
taskToRun = new TraceCallable<>(tracing(), spanNamer(), task);
taskToRun = new TraceCallable<>(tracing(), spanNamer(), task, this.beanName);
}
return this.delegate.submit(taskToRun);
}

View File

@@ -42,6 +42,8 @@ public class LazyTraceExecutor implements Executor {
private final Executor delegate;
private final String beanName;
private Tracing tracing;
private SpanNamer spanNamer;
@@ -49,6 +51,14 @@ public class LazyTraceExecutor implements Executor {
public LazyTraceExecutor(BeanFactory beanFactory, Executor delegate) {
this.beanFactory = beanFactory;
this.delegate = delegate;
this.beanName = null;
}
public LazyTraceExecutor(BeanFactory beanFactory, Executor delegate,
String beanName) {
this.beanFactory = beanFactory;
this.delegate = delegate;
this.beanName = beanName;
}
@Override
@@ -66,7 +76,8 @@ public class LazyTraceExecutor implements Executor {
return;
}
}
this.delegate.execute(new TraceRunnable(this.tracing, spanNamer(), command));
this.delegate.execute(
new TraceRunnable(this.tracing, spanNamer(), command, this.beanName));
}
// due to some race conditions trace keys might not be ready yet

View File

@@ -60,6 +60,8 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor {
private final ScheduledThreadPoolExecutor delegate;
private final String beanName;
private final Method decorateTaskRunnable;
private final Method decorateTaskCallable;
@@ -81,10 +83,11 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor {
private SpanNamer spanNamer;
LazyTraceScheduledThreadPoolExecutor(int corePoolSize, BeanFactory beanFactory,
ScheduledThreadPoolExecutor delegate) {
ScheduledThreadPoolExecutor delegate, String beanName) {
super(corePoolSize);
this.beanFactory = beanFactory;
this.delegate = delegate;
this.beanName = beanName;
this.decorateTaskRunnable = ReflectionUtils.findMethod(
ScheduledThreadPoolExecutor.class, "decorateTask", Runnable.class,
RunnableScheduledFuture.class);
@@ -122,82 +125,12 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor {
}
LazyTraceScheduledThreadPoolExecutor(int corePoolSize, ThreadFactory threadFactory,
BeanFactory beanFactory, ScheduledThreadPoolExecutor delegate) {
super(corePoolSize, threadFactory);
this.beanFactory = beanFactory;
this.delegate = delegate;
this.decorateTaskRunnable = ReflectionUtils.findMethod(
ScheduledThreadPoolExecutor.class, "decorateTask", Runnable.class,
RunnableScheduledFuture.class);
makeAccessibleIfNotNull(this.decorateTaskRunnable);
this.decorateTaskCallable = ReflectionUtils.findMethod(
ScheduledThreadPoolExecutor.class, "decorateTaskCallable", Callable.class,
RunnableScheduledFuture.class);
makeAccessibleIfNotNull(this.decorateTaskCallable);
this.finalize = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class,
"finalize");
makeAccessibleIfNotNull(this.finalize);
this.beforeExecute = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class,
"beforeExecute");
makeAccessibleIfNotNull(this.beforeExecute);
this.afterExecute = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class,
"afterExecute", null);
makeAccessibleIfNotNull(this.afterExecute);
this.terminated = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class,
"terminated", null);
makeAccessibleIfNotNull(this.terminated);
this.newTaskForRunnable = ReflectionUtils.findMethod(
ScheduledThreadPoolExecutor.class, "newTaskFor", Runnable.class,
Object.class);
makeAccessibleIfNotNull(this.newTaskForRunnable);
this.newTaskForCallable = ReflectionUtils.findMethod(
ScheduledThreadPoolExecutor.class, "newTaskFor", Callable.class,
Object.class);
makeAccessibleIfNotNull(this.newTaskForCallable);
}
LazyTraceScheduledThreadPoolExecutor(int corePoolSize,
RejectedExecutionHandler handler, BeanFactory beanFactory,
ScheduledThreadPoolExecutor delegate) {
super(corePoolSize, handler);
this.beanFactory = beanFactory;
this.delegate = delegate;
this.decorateTaskRunnable = ReflectionUtils.findMethod(
ScheduledThreadPoolExecutor.class, "decorateTask", Runnable.class,
RunnableScheduledFuture.class);
makeAccessibleIfNotNull(this.decorateTaskRunnable);
this.decorateTaskCallable = ReflectionUtils.findMethod(
ScheduledThreadPoolExecutor.class, "decorateTaskCallable", Callable.class,
RunnableScheduledFuture.class);
makeAccessibleIfNotNull(this.decorateTaskCallable);
this.finalize = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class,
"finalize", null);
makeAccessibleIfNotNull(this.finalize);
this.beforeExecute = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class,
"beforeExecute", null);
makeAccessibleIfNotNull(this.beforeExecute);
this.afterExecute = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class,
"afterExecute", null);
makeAccessibleIfNotNull(this.afterExecute);
this.terminated = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class,
"terminated", null);
makeAccessibleIfNotNull(this.terminated);
this.newTaskForRunnable = ReflectionUtils.findMethod(
ScheduledThreadPoolExecutor.class, "newTaskFor", Runnable.class,
Object.class);
makeAccessibleIfNotNull(this.newTaskForRunnable);
this.newTaskForCallable = ReflectionUtils.findMethod(
ScheduledThreadPoolExecutor.class, "newTaskFor", Callable.class,
Object.class);
makeAccessibleIfNotNull(this.newTaskForCallable);
}
LazyTraceScheduledThreadPoolExecutor(int corePoolSize, ThreadFactory threadFactory,
RejectedExecutionHandler handler, BeanFactory beanFactory,
ScheduledThreadPoolExecutor delegate) {
ScheduledThreadPoolExecutor delegate, String beanName) {
super(corePoolSize, threadFactory, handler);
this.beanFactory = beanFactory;
this.delegate = delegate;
this.beanName = beanName;
this.decorateTaskRunnable = ReflectionUtils.findMethod(
ScheduledThreadPoolExecutor.class, "decorateTask", Runnable.class,
RunnableScheduledFuture.class);
@@ -234,7 +167,7 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor {
RunnableScheduledFuture<V> task) {
return (RunnableScheduledFuture<V>) ReflectionUtils.invokeMethod(
this.decorateTaskRunnable, this.delegate,
new TraceRunnable(tracing(), spanNamer(), runnable), task);
new TraceRunnable(tracing(), spanNamer(), runnable, this.beanName), task);
}
@Override
@@ -243,57 +176,63 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor {
RunnableScheduledFuture<V> task) {
return (RunnableScheduledFuture<V>) ReflectionUtils.invokeMethod(
this.decorateTaskCallable, this.delegate,
new TraceCallable<>(tracing(), spanNamer(), callable), task);
new TraceCallable<>(tracing(), spanNamer(), callable, this.beanName),
task);
}
@Override
public ScheduledFuture<?> schedule(Runnable command, long delay, TimeUnit unit) {
return this.delegate.schedule(new TraceRunnable(tracing(), spanNamer(), command),
delay, unit);
return this.delegate.schedule(
new TraceRunnable(tracing(), spanNamer(), command, this.beanName), delay,
unit);
}
@Override
public <V> ScheduledFuture<V> schedule(Callable<V> callable, long delay,
TimeUnit unit) {
return this.delegate.schedule(
new TraceCallable<>(tracing(), spanNamer(), callable), delay, unit);
new TraceCallable<>(tracing(), spanNamer(), callable, this.beanName),
delay, unit);
}
@Override
public ScheduledFuture<?> scheduleAtFixedRate(Runnable command, long initialDelay,
long period, TimeUnit unit) {
return this.delegate.scheduleAtFixedRate(
new TraceRunnable(tracing(), spanNamer(), command), initialDelay, period,
unit);
new TraceRunnable(tracing(), spanNamer(), command, this.beanName),
initialDelay, period, unit);
}
@Override
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable command, long initialDelay,
long delay, TimeUnit unit) {
return this.delegate.scheduleWithFixedDelay(
new TraceRunnable(tracing(), spanNamer(), command), initialDelay, delay,
unit);
new TraceRunnable(tracing(), spanNamer(), command, this.beanName),
initialDelay, delay, unit);
}
@Override
public void execute(Runnable command) {
this.delegate.execute(new TraceRunnable(tracing(), spanNamer(), command));
this.delegate.execute(
new TraceRunnable(tracing(), spanNamer(), command, this.beanName));
}
@Override
public Future<?> submit(Runnable task) {
return this.delegate.submit(new TraceRunnable(tracing(), spanNamer(), task));
return this.delegate
.submit(new TraceRunnable(tracing(), spanNamer(), task, this.beanName));
}
@Override
public <T> Future<T> submit(Runnable task, T result) {
return this.delegate.submit(new TraceRunnable(tracing(), spanNamer(), task),
result);
return this.delegate.submit(
new TraceRunnable(tracing(), spanNamer(), task, this.beanName), result);
}
@Override
public <T> Future<T> submit(Callable<T> task) {
return this.delegate.submit(new TraceCallable<>(tracing(), spanNamer(), task));
return this.delegate
.submit(new TraceCallable<>(tracing(), spanNamer(), task, this.beanName));
}
@Override
@@ -479,13 +418,13 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor {
@Override
public void beforeExecute(Thread t, Runnable r) {
ReflectionUtils.invokeMethod(this.beforeExecute, this.delegate, t,
new TraceRunnable(tracing(), spanNamer(), r));
new TraceRunnable(tracing(), spanNamer(), r, this.beanName));
}
@Override
public void afterExecute(Runnable r, Throwable t) {
ReflectionUtils.invokeMethod(this.afterExecute, this.delegate,
new TraceRunnable(tracing(), spanNamer(), r), t);
new TraceRunnable(tracing(), spanNamer(), r, this.beanName), t);
}
@Override
@@ -497,7 +436,8 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor {
@SuppressWarnings("unchecked")
public <T> RunnableFuture<T> newTaskFor(Runnable runnable, T value) {
return (RunnableFuture<T>) ReflectionUtils.invokeMethod(this.newTaskForRunnable,
this.delegate, new TraceRunnable(tracing(), spanNamer(), runnable),
this.delegate,
new TraceRunnable(tracing(), spanNamer(), runnable, this.beanName),
value);
}
@@ -505,7 +445,8 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor {
@SuppressWarnings("unchecked")
public <T> RunnableFuture<T> newTaskFor(Callable<T> callable) {
return (RunnableFuture<T>) ReflectionUtils.invokeMethod(this.newTaskForCallable,
this.delegate, new TraceCallable<>(tracing(), spanNamer(), callable));
this.delegate,
new TraceCallable<>(tracing(), spanNamer(), callable, this.beanName));
}
@Override
@@ -519,7 +460,7 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor {
List<Callable<T>> ts = new ArrayList<>();
for (Callable<T> task : tasks) {
if (!(task instanceof TraceCallable)) {
ts.add(new TraceCallable<>(tracing(), spanNamer(), task));
ts.add(new TraceCallable<>(tracing(), spanNamer(), task, this.beanName));
}
}
return ts;

View File

@@ -51,6 +51,8 @@ public class LazyTraceThreadPoolTaskExecutor extends ThreadPoolTaskExecutor {
private final ThreadPoolTaskExecutor delegate;
private final String beanName;
private Tracing tracing;
private SpanNamer spanNamer;
@@ -59,44 +61,55 @@ public class LazyTraceThreadPoolTaskExecutor extends ThreadPoolTaskExecutor {
ThreadPoolTaskExecutor delegate) {
this.beanFactory = beanFactory;
this.delegate = delegate;
this.beanName = null;
}
public LazyTraceThreadPoolTaskExecutor(BeanFactory beanFactory,
ThreadPoolTaskExecutor delegate, String beanName) {
this.beanFactory = beanFactory;
this.delegate = delegate;
this.beanName = beanName;
}
@Override
public void execute(Runnable task) {
this.delegate.execute(ContextUtil.isContextUnusable(this.beanFactory) ? task
: new TraceRunnable(tracing(), spanNamer(), task));
: new TraceRunnable(tracing(), spanNamer(), task, this.beanName));
}
@Override
public void execute(Runnable task, long startTimeout) {
this.delegate.execute(ContextUtil.isContextUnusable(this.beanFactory) ? task
: new TraceRunnable(tracing(), spanNamer(), task), startTimeout);
this.delegate.execute(
ContextUtil.isContextUnusable(this.beanFactory) ? task
: new TraceRunnable(tracing(), spanNamer(), task, this.beanName),
startTimeout);
}
@Override
public Future<?> submit(Runnable task) {
return this.delegate.submit(ContextUtil.isContextUnusable(this.beanFactory) ? task
: new TraceRunnable(tracing(), spanNamer(), task));
: new TraceRunnable(tracing(), spanNamer(), task, this.beanName));
}
@Override
public <T> Future<T> submit(Callable<T> task) {
return this.delegate.submit(ContextUtil.isContextUnusable(this.beanFactory) ? task
: new TraceCallable<>(tracing(), spanNamer(), task));
: new TraceCallable<>(tracing(), spanNamer(), task, this.beanName));
}
@Override
public ListenableFuture<?> submitListenable(Runnable task) {
return this.delegate
.submitListenable(ContextUtil.isContextUnusable(this.beanFactory) ? task
: new TraceRunnable(tracing(), spanNamer(), task));
: new TraceRunnable(tracing(), spanNamer(), task, this.beanName));
}
@Override
public <T> ListenableFuture<T> submitListenable(Callable<T> task) {
return this.delegate
.submitListenable(ContextUtil.isContextUnusable(this.beanFactory) ? task
: new TraceCallable<>(tracing(), spanNamer(), task));
: new TraceCallable<>(tracing(), spanNamer(), task,
this.beanName));
}
@Override

View File

@@ -62,6 +62,8 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler {
private final ThreadPoolTaskScheduler delegate;
private final String beanName;
private final Method initializeExecutor;
private final Method createExecutor;
@@ -77,9 +79,10 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler {
private SpanNamer spanNamer;
LazyTraceThreadPoolTaskScheduler(BeanFactory beanFactory,
ThreadPoolTaskScheduler delegate) {
ThreadPoolTaskScheduler delegate, String beanName) {
this.beanFactory = beanFactory;
this.delegate = delegate;
this.beanName = beanName;
this.initializeExecutor = ReflectionUtils
.findMethod(ThreadPoolTaskScheduler.class, "initializeExecutor", null);
makeAccessibleIfNotNull(this.initializeExecutor);
@@ -127,11 +130,13 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler {
if (executorService instanceof TraceableScheduledExecutorService) {
return executorService;
}
return new TraceableExecutorService(this.beanFactory, executorService);
return new TraceableExecutorService(this.beanFactory, executorService,
this.beanName);
}
private ThreadFactory traceThreadFactory(ThreadFactory threadFactory) {
return r -> threadFactory.newThread(new TraceRunnable(tracing(), spanNamer(), r));
return r -> threadFactory
.newThread(new TraceRunnable(tracing(), spanNamer(), r, this.beanName));
}
@Override
@@ -144,14 +149,16 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler {
if (executorService instanceof TraceableScheduledExecutorService) {
return executorService;
}
return new TraceableScheduledExecutorService(this.beanFactory, executorService);
return new TraceableScheduledExecutorService(this.beanFactory, executorService,
this.beanName);
}
@Override
public ScheduledExecutorService getScheduledExecutor() throws IllegalStateException {
ScheduledExecutorService executor = this.delegate.getScheduledExecutor();
return executor instanceof TraceableScheduledExecutorService ? executor
: new TraceableScheduledExecutorService(this.beanFactory, executor);
: new TraceableScheduledExecutorService(this.beanFactory, executor,
this.beanName);
}
@Override
@@ -164,7 +171,7 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler {
}
return new LazyTraceScheduledThreadPoolExecutor(executor.getCorePoolSize(),
executor.getThreadFactory(), executor.getRejectedExecutionHandler(),
this.beanFactory, executor);
this.beanFactory, executor, this.beanName);
}
@Override
@@ -184,41 +191,45 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler {
@Override
public void execute(Runnable task) {
this.delegate.execute(new TraceRunnable(tracing(), spanNamer(), task));
this.delegate
.execute(new TraceRunnable(tracing(), spanNamer(), task, this.beanName));
}
@Override
public void execute(Runnable task, long startTimeout) {
this.delegate.execute(new TraceRunnable(tracing(), spanNamer(), task),
this.delegate.execute(
new TraceRunnable(tracing(), spanNamer(), task, this.beanName),
startTimeout);
}
@Override
public Future<?> submit(Runnable task) {
return this.delegate.submit(new TraceRunnable(tracing(), spanNamer(), task));
return this.delegate
.submit(new TraceRunnable(tracing(), spanNamer(), task, this.beanName));
}
@Override
public <T> Future<T> submit(Callable<T> task) {
return this.delegate.submit(new TraceCallable<>(tracing(), spanNamer(), task));
return this.delegate
.submit(new TraceCallable<>(tracing(), spanNamer(), task, this.beanName));
}
@Override
public ListenableFuture<?> submitListenable(Runnable task) {
return this.delegate
.submitListenable(new TraceRunnable(tracing(), spanNamer(), task));
return this.delegate.submitListenable(
new TraceRunnable(tracing(), spanNamer(), task, this.beanName));
}
@Override
public <T> ListenableFuture<T> submitListenable(Callable<T> task) {
return this.delegate
.submitListenable(new TraceCallable<>(tracing(), spanNamer(), task));
return this.delegate.submitListenable(
new TraceCallable<>(tracing(), spanNamer(), task, this.beanName));
}
@Override
public void cancelRemainingTask(Runnable task) {
ReflectionUtils.invokeMethod(this.cancelRemainingTask, this.delegate,
new TraceRunnable(tracing(), spanNamer(), task));
new TraceRunnable(tracing(), spanNamer(), task, this.beanName));
}
@Override
@@ -229,13 +240,14 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler {
@Override
@Nullable
public ScheduledFuture<?> schedule(Runnable task, Trigger trigger) {
return this.delegate.schedule(new TraceRunnable(tracing(), spanNamer(), task),
trigger);
return this.delegate.schedule(
new TraceRunnable(tracing(), spanNamer(), task, this.beanName), trigger);
}
@Override
public ScheduledFuture<?> schedule(Runnable task, Date startTime) {
return this.delegate.schedule(new TraceRunnable(tracing(), spanNamer(), task),
return this.delegate.schedule(
new TraceRunnable(tracing(), spanNamer(), task, this.beanName),
startTime);
}
@@ -243,26 +255,28 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler {
public ScheduledFuture<?> scheduleAtFixedRate(Runnable task, Date startTime,
long period) {
return this.delegate.scheduleAtFixedRate(
new TraceRunnable(tracing(), spanNamer(), task), startTime, period);
new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTime,
period);
}
@Override
public ScheduledFuture<?> scheduleAtFixedRate(Runnable task, long period) {
return this.delegate.scheduleAtFixedRate(
new TraceRunnable(tracing(), spanNamer(), task), period);
new TraceRunnable(tracing(), spanNamer(), task, this.beanName), period);
}
@Override
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable task, Date startTime,
long delay) {
return this.delegate.scheduleWithFixedDelay(
new TraceRunnable(tracing(), spanNamer(), task), startTime, delay);
new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTime,
delay);
}
@Override
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable task, long delay) {
return this.delegate.scheduleWithFixedDelay(
new TraceRunnable(tracing(), spanNamer(), task), delay);
new TraceRunnable(tracing(), spanNamer(), task, this.beanName), delay);
}
@Override
@@ -385,7 +399,8 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler {
@Override
public ScheduledFuture<?> schedule(Runnable task, Instant startTime) {
return this.delegate.schedule(new TraceRunnable(tracing(), spanNamer(), task),
return this.delegate.schedule(
new TraceRunnable(tracing(), spanNamer(), task, this.beanName),
startTime);
}
@@ -393,26 +408,28 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler {
public ScheduledFuture<?> scheduleAtFixedRate(Runnable task, Instant startTime,
Duration period) {
return this.delegate.scheduleAtFixedRate(
new TraceRunnable(tracing(), spanNamer(), task), startTime, period);
new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTime,
period);
}
@Override
public ScheduledFuture<?> scheduleAtFixedRate(Runnable task, Duration period) {
return this.delegate.scheduleAtFixedRate(
new TraceRunnable(tracing(), spanNamer(), task), period);
new TraceRunnable(tracing(), spanNamer(), task, this.beanName), period);
}
@Override
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable task, Instant startTime,
Duration delay) {
return this.delegate.scheduleWithFixedDelay(
new TraceRunnable(tracing(), spanNamer(), task), startTime, delay);
new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTime,
delay);
}
@Override
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable task, Duration delay) {
return this.delegate.scheduleWithFixedDelay(
new TraceRunnable(tracing(), spanNamer(), task), delay);
new TraceRunnable(tracing(), spanNamer(), task, this.beanName), delay);
}
private Tracing tracing() {

View File

@@ -42,7 +42,7 @@ public class TraceableExecutorService implements ExecutorService {
final ExecutorService delegate;
private final String spanName;
final String spanName;
Tracing tracing;

View File

@@ -39,43 +39,52 @@ public class TraceableScheduledExecutorService extends TraceableExecutorService
super(beanFactory, delegate);
}
public TraceableScheduledExecutorService(BeanFactory beanFactory,
final ExecutorService delegate, String beanName) {
super(beanFactory, delegate, beanName);
}
private ScheduledExecutorService getScheduledExecutorService() {
return (ScheduledExecutorService) this.delegate;
}
@Override
public ScheduledFuture<?> schedule(Runnable command, long delay, TimeUnit unit) {
return getScheduledExecutorService().schedule(
ContextUtil.isContextUnusable(this.beanFactory) ? command
: new TraceRunnable(tracing(), spanNamer(), command),
delay, unit);
return getScheduledExecutorService()
.schedule(ContextUtil.isContextUnusable(this.beanFactory) ? command
: new TraceRunnable(tracing(), spanNamer(), command,
this.spanName),
delay, unit);
}
@Override
public <V> ScheduledFuture<V> schedule(Callable<V> callable, long delay,
TimeUnit unit) {
return getScheduledExecutorService().schedule(
ContextUtil.isContextUnusable(this.beanFactory) ? callable
: new TraceCallable<>(tracing(), spanNamer(), callable),
delay, unit);
return getScheduledExecutorService()
.schedule(ContextUtil.isContextUnusable(this.beanFactory) ? callable
: new TraceCallable<>(tracing(), spanNamer(), callable,
this.spanName),
delay, unit);
}
@Override
public ScheduledFuture<?> scheduleAtFixedRate(Runnable command, long initialDelay,
long period, TimeUnit unit) {
return getScheduledExecutorService().scheduleAtFixedRate(
ContextUtil.isContextUnusable(this.beanFactory) ? command
: new TraceRunnable(tracing(), spanNamer(), command),
initialDelay, period, unit);
return getScheduledExecutorService()
.scheduleAtFixedRate(ContextUtil.isContextUnusable(this.beanFactory)
? command : new TraceRunnable(tracing(), spanNamer(), command,
this.spanName),
initialDelay, period, unit);
}
@Override
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable command, long initialDelay,
long delay, TimeUnit unit) {
return getScheduledExecutorService().scheduleWithFixedDelay(
ContextUtil.isContextUnusable(this.beanFactory) ? command
: new TraceRunnable(tracing(), spanNamer(), command),
initialDelay, delay, unit);
return getScheduledExecutorService()
.scheduleWithFixedDelay(ContextUtil.isContextUnusable(this.beanFactory)
? command : new TraceRunnable(tracing(), spanNamer(), command,
this.spanName),
initialDelay, delay, unit);
}
}

View File

@@ -212,7 +212,7 @@ public class ExecutorBeanPostProcessorTests {
ExecutorBeanPostProcessor bpp = new ExecutorBeanPostProcessor(this.beanFactory) {
@Override
Object createThreadPoolTaskExecutorProxy(Object bean, boolean cglibProxy,
ThreadPoolTaskExecutor executor) {
ThreadPoolTaskExecutor executor, String beanName) {
throw new AopConfigException("foo");
}
};

View File

@@ -39,7 +39,8 @@ public class LazyTraceScheduledThreadPoolExecutorTests {
};
BeanFactory beanFactory = BDDMockito.mock(BeanFactory.class);
new LazyTraceScheduledThreadPoolExecutor(10, beanFactory, executor).finalize();
new LazyTraceScheduledThreadPoolExecutor(10, beanFactory, executor, null)
.finalize();
BDDAssertions.then(wasCalled).isFalse();
BDDAssertions.then(executor.isShutdown()).isFalse();

View File

@@ -58,8 +58,8 @@ public class LazyTraceThreadPoolTaskSchedulerTests {
@BeforeEach
public void setup() {
this.executor = new LazyTraceThreadPoolTaskScheduler(beanFactory(),
this.delegate);
this.executor = new LazyTraceThreadPoolTaskScheduler(beanFactory(), this.delegate,
null);
}
@AfterEach

View File

@@ -78,7 +78,7 @@ public class TraceableExecutorServiceTests {
@BeforeEach
public void setup() {
this.traceManagerableExecutorService = new TraceableExecutorService(
beanFactory(true), this.executorService);
beanFactory(true), this.executorService, "foo");
this.spans.clear();
this.spanVerifyingRunnable.clear();
}