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>

关键点说明

  1. @TraceCrossThread 是核心注解,确保Trace上下文跨线程传递

  2. ActiveSpan 用于设置操作名称和标签,增强可观测性

  3. TraceContext 用于获取当前Trace信息

  4. 通过包装模式,业务代码无需关心Trace传递细节

更多推荐