skywalking收集异步任务信息
·
1. 使用@TraceCrossThread注解
import org.apache.skywalking.apm.toolkit.trace.TraceCrossThread;
import org.apache.skywalking.apm.toolkit.trace.TraceContext;
import org.apache.skywalking.apm.toolkit.trace.ActiveSpan;
import java.util.concurrent.Callable;
import java.util.function.Supplier;
@TraceCrossThread
public class TraceableSupplier<T> implements Supplier<T> {
private final Supplier<T> delegate;
public TraceableSupplier(Supplier<T> delegate) {
this.delegate = delegate;
}
@Override
public T get() {
// 这里会自动传递Trace上下文
ActiveSpan.setOperationName("AsyncOperation");
return delegate.get();
}
}
@TraceCrossThread
public class TraceableCallable<T> implements Callable<T> {
private final Callable<T> delegate;
public TraceableCallable(Callable<T> delegate) {
this.delegate = delegate;
}
@Override
public T call() throws Exception {
ActiveSpan.setOperationName("AsyncCallable");
return delegate.call();
}
}
@TraceCrossThread
public class TraceableRunnable implements Runnable {
private final Runnable delegate;
public TraceableRunnable(Runnable delegate) {
this.delegate = delegate;
}
@Override
public void run() {
ActiveSpan.setOperationName("AsyncRunnable");
delegate.run();
}
}
工具类
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;
import java.util.function.Supplier;
public class TraceableCompletableFuture {
/**
* 可追踪的supplyAsync
*/
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier) {
return CompletableFuture.supplyAsync(new TraceableSupplier<>(supplier));
}
/**
* 可追踪的supplyAsync(指定线程池)
*/
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier, Executor executor) {
return CompletableFuture.supplyAsync(new TraceableSupplier<>(supplier), executor);
}
/**
* 可追踪的runAsync
*/
public static CompletableFuture<Void> runAsync(Runnable runnable) {
return CompletableFuture.runAsync(new TraceableRunnable(runnable));
}
/**
* 可追踪的runAsync(指定线程池)
*/
public static CompletableFuture<Void> runAsync(Runnable runnable, Executor executor) {
return CompletableFuture.runAsync(new TraceableRunnable(runnable), executor);
}
}
2. 业务使用示例
import org.springframework.stereotype.Service;
import java.util.concurrent.CompletableFuture;
@Service
public class UserService {
/**
* 简单的异步查询
*/
public CompletableFuture<String> getUserInfoAsync(Long userId) {
return TraceableCompletableFuture.supplyAsync(() -> {
// 模拟业务逻辑
ActiveSpan.tag("user_id", String.valueOf(userId));
System.out.println("查询用户信息,TraceId: " + TraceContext.traceId());
// 模拟耗时操作
try {
Thread.sleep(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
return "UserInfo-" + userId;
});
}
/**
* 复杂的异步处理链
*/
public CompletableFuture<String> processUserDataAsync(Long userId) {
return TraceableCompletableFuture.supplyAsync(() -> {
ActiveSpan.setOperationName("ProcessUserData");
ActiveSpan.tag("user_id", String.valueOf(userId));
// 第一步处理
return "raw-data-" + userId;
})
.thenApplyAsync(rawData -> {
// 第二步处理
ActiveSpan.setOperationName("TransformData");
return rawData.toUpperCase();
})
.thenApplyAsync(transformedData -> {
// 第三步处理
ActiveSpan.setOperationName("FinalProcess");
return "Processed: " + transformedData;
});
}
}
3. pom.xml依赖
<dependencies>
<!-- SkyWalking Toolkit -->
<dependency>
<groupId>org.apache.skywalking</groupId>
<artifactId>apm-toolkit-trace</artifactId>
<version>8.5.0</version>
</dependency>
<dependency>
<groupId>org.apache.skywalking</groupId>
<artifactId>apm-toolkit-logback-1.x</artifactId>
<version>8.5.0</version>
</dependency>
<!-- 如果需要Spring支持 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
<version>2.7.0</version>
</dependency>
</dependencies>
关键点说明
-
@TraceCrossThread 是核心注解,确保Trace上下文跨线程传递
-
ActiveSpan 用于设置操作名称和标签,增强可观测性
-
TraceContext 用于获取当前Trace信息
-
通过包装模式,业务代码无需关心Trace传递细节
更多推荐

所有评论(0)