上节介绍了skywalking的链路追踪、数据持久化、日志查看的功能,已经能满足绝大部分的需求,但是链路追踪这里有个问题:无法支持异步调用的链路追踪。

获取文章列表的这个接口的业务方法中涉及到了异步调用,代码如下:

	@SneakyThrows
    @Override
    public PageData<ArticleVo> list(ArticleListReq req) {
			.......
        //① 异步获取文章作者的信息
        CompletableFuture<List<SysUser>> futureSysUsers = CompletableFuture.supplyAsync(() -> {
            //将主线程的请求信息设置到异步线程中,否则会丢失请求上下文,导致调用失败
            RequestContextHolder.setRequestAttributes(attributes);
            //获取所有的作者ID
            List<String> authorIds = articles.stream().map(Article::getAuthorId).collect(Collectors.toList());
            List<SysUser> sysUsers = new ArrayList<>();
            //调用user服务获取作者信息
            ResultMsg<List<SysUser>> authorRes = userFeign.listByUserId(authorIds);
            if (ResultMsg.isSuccess(authorRes) && authorRes.getData() != null) {
                sysUsers = authorRes.getData();
            }
            return sysUsers;
        }, asyncTaskExecutor);

        //② 异步获取文章的评论总数
        CompletableFuture<List<TotalVo>> futureCommentTotal = CompletableFuture.supplyAsync(() -> {
            //将主线程的请求信息设置到异步线程中,否则会丢失请求上下文,导致调用失败
            RequestContextHolder.setRequestAttributes(attributes);
            List<TotalVo> commentVos = new ArrayList<>();
            //获取所有的文章ID
            List<CommentListReq> params = articles.stream().map(p -> {
                CommentListReq param = new CommentListReq();
                param.setArticleId(p.getArticleId());
                return param;
            }).collect(Collectors.toList());
            //调用评论信息获取评论总数
            ResultMsg<List<TotalVo>> commentRes = commentFeign.listTotal(params);
            if (ResultMsg.isSuccess(commentRes) && commentRes.getData() != null) {
                commentVos = commentRes.getData();
            }
            return commentVos;
        }, asyncTaskExecutor);

        CompletableFuture.allOf(futureSysUsers,futureCommentTotal).join();
		...............
        
    }

①处的代码:异步调用用户服务获取作者信息

②处的代码:异步调用评论服务获取文章的评论总数

以上两处都涉及了异步调用,笔者调用了/blog-article/front/article/list这个接口,在skywalking中的链路显示如下:

这是为什么?因为异步调用导致traceId改变了,traceId是通过ThreadLocal存储,只能在当前线程获取,无法在子线程中获取。

如何解决呢?

解决方案有多种,笔者这里介绍一种常用且推荐使用的方案:异步调用时候通过Skywalking提供的包装类进行调用。有如下三种包装类:

原始类提供的包装类拦截方法使用技巧
CallableCallableWrappercallCallableWrapper.of(xxxCallable)
RunnableRunnableWrapperrunRunnableWrapper.of(xxxRunable)
SupplierSupplierWrappergetSupplierWrapper.of(xxxSupplier)

当然要是这些包装类需要添加如下依赖:

<!--异步链路追踪依赖-->
<dependency>
    <groupId>org.apache.skywalking</groupId>
    <artifactId>apm-toolkit-trace</artifactId>
</dependency>

笔者将这个依赖同样放在了blog-skywalking-starter中,每个微服务只需要依赖:

<dependency>
    <groupId>com.mugu.blog</groupId>
    <artifactId>blog-skywalking-starter</artifactId>
</dependency>

接下来对上面的接口进行改造,由于使用的是CompletableFuture.supplyAsync(),因此需要使用SupplierWrapper进行包装一层,改造后的代码如下:

			//① 异步获取文章作者的信息
        CompletableFuture<List<SysUser>> futureSysUsers = CompletableFuture.supplyAsync(SupplierWrapper.of(() -> {
            ......
        }), asyncTaskExecutor);

        //② 异步获取文章的评论总数
        CompletableFuture<List<TotalVo>> futureCommentTotal = CompletableFuture.supplyAsync(SupplierWrapper.of(() -> {
            .......
        }), asyncTaskExecutor);

看出来变化了吗?其实就是包装了一层,CompletableFuture.supplyAsync()调用变成了如下:

CompletableFuture.supplyAsync(SupplierWrapper.of(()->{}),asyncTaskExecutor);

同样的CompletableFuture.runAsync()调用变成了如下这样:

CompletableFuture.runAsync(RunnableWrapper.of(()->{}),asyncTaskExecutor);

改造之后,分为如下两个场景测试:

1、blog-comments和blog-user-boot不启动

这两个服务不启动会触发降级,来看下skywalking中的链路状态,是否会提示调用错误,如下图:

可以看到异步调用也被链路追踪了,并且降级也通过飘红提示调用失败了。点击红色那条链路,也会出现异常日志,如下图:

2、blog-comments和blog-user-boot启动

两个服务全部启动,再次看下链路的状态,如下图:

异步调用已经被追踪了,同时链路正常,并没有报错。

为什么用SupplierWrapper包装类包装一下就能被追踪呢?

进入源码一探究竟,源码如下:

@TraceCrossThread
public class SupplierWrapper<V> implements Supplier<V> {

看到这么一个注解**@TraceCrossThread**,看这个名字就知道是和线程有关;这个注解是skywalking内部拦截的标记。

另外两个包装类都被**@TraceCrossThread**这个注解标记了。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐