SleuthXXL-JOB中不生效的主要原因‌是由于XXL-JOB的调度和执行机制与Spring Cloud Sleuth的集成存在问题。XXL-JOB的调度中心负责发起调度请求,而任务执行器负责接收请求并执行任务。这种设计导致每次任务执行时,都会创建一个新的线程来处理任务,而Sleuth的TraceContext无法跨线程传递,从而导致traceId在任务执行过程中丢失或变化‌

解决方案:

为了不变更xxjob的源码结构,这里我通过spring aop的方式对@XxlJob注解进行切面编程,而注解@Xxljob就是一个切点。代码如下所示:

将Tracing注入到Aspect切面类中

@EnableAspectJAutoProxy
@Configuration
public class AOPConfiguration {


  @Bean
  public XxljobTraceAspect xxljobTraceAspect(Tracing tracing) {
    log.info(">>>>>>>>>>> XxljobTraceAspect init,tracing is {}",tracing.getClass().getName());
    log.info(">>>>>>>>>>> XxljobTraceAspect init,tracing is {}", JSON.toJSONString(tracing));
    XxljobTraceAspect aspect = XxljobTraceAspect.builder().tracing(tracing).build();
    log.info(">>>>>>>>>>> XxljobTraceAspect init completed.");
    return aspect;
  }
}






>>>>>>>>>>> XxljobTraceAspect init,tracing is brave.Tracing$Default
>>>>>>>>>>> XxljobTraceAspect init,tracing is {"noop":false}

TraceXXJobAspect :

import org.aspectj.lang.ProceedingJoinPoint;
import org.aspectj.lang.annotation.Around;
import org.aspectj.lang.annotation.Aspect;
import org.springframework.cloud.sleuth.util.SpanNameUtil;

import brave.Span;
import brave.Tracer;
import brave.Tracing;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor;
import lombok.extern.slf4j.Slf4j;

@AllArgsConstructor
@NoArgsConstructor
@Data
@Builder
@EqualsAndHashCode(callSuper = false)
@Aspect
@Slf4j
public class TraceXXJobAspect {

    private Tracing tracing;

    public TraceXXJobAspect(Tracing tracing) {
        this.tracing = tracing;
    }

    private static final String CLASS_KEY = "class";

    private static final String METHOD_KEY = "method";


    @Around("execution (@com.xxl.job.core.handler.annotation.XxlJob  * *.*(..))")
    public Object traceBackgroundThread(final ProceedingJoinPoint pjp) throws Throwable {

        Tracer tracer = this.tracing.tracer();
        //清除scop中原span
        tracer.withSpanInScope(null);
        String spanName = SpanNameUtil.toLowerHyphen(pjp.getSignature().getName());
        Span span = tracer.newTrace().name(spanName);
        try (Tracer.SpanInScope ws = tracer.withSpanInScope(span.start())) {
            span.tag(CLASS_KEY, pjp.getTarget().getClass().getSimpleName());
            span.tag(METHOD_KEY, pjp.getSignature().getName());
            return pjp.proceed();
        }catch (Throwable ex) {
            String message = ex.getMessage() == null ? ex.getClass().getSimpleName()
                    : ex.getMessage();
            span.tag("error", message);
          
            throw ex;
        }
        finally {
            span.finish();
        }
    }

}

加入上面两个类,就可以实现xxl中的链路跟踪了。

这里解释下XxlJob初始化的过程,我们在使用xxljob的时候需要自己创建XxlJobSpringExecutor类实例对象,并注入到sping中(xxljob依赖于spring),查看下XxlJobSpringExecutor源码:

public class XxlJobSpringExecutor extends XxlJobExecutor implements ApplicationContextAware, SmartInitializingSingleton, DisposableBean {
    private static final Logger logger = LoggerFactory.getLogger(XxlJobSpringExecutor.class);


    // start
    @Override
    public void afterSingletonsInstantiated() {

        // init JobHandler Repository  在2.2之后@JobHandler方式被xxljob给取消了
        /*initJobHandlerRepository(applicationContext);*/

        // 扫描xxljob注解并为每一个被xxljob生成 IJobHandler对象
        initJobHandlerMethodRepository(applicationContext);

        // refresh GlueFactory
        GlueFactory.refreshInstance(1);

        // super start
        try {
            super.start();
        } catch (Exception e) {
            throw new RuntimeException(e);
        }
    }

    // destroy
    @Override
    public void destroy() {
        super.destroy();
    }


    /*private void initJobHandlerRepository(ApplicationContext applicationContext) {
        if (applicationContext == null) {
            return;
        }

        // init job handler action
        Map<String, Object> serviceBeanMap = applicationContext.getBeansWithAnnotation(JobHandler.class);

        if (serviceBeanMap != null && serviceBeanMap.size() > 0) {
            for (Object serviceBean : serviceBeanMap.values()) {
                if (serviceBean instanceof IJobHandler) {
                    String name = serviceBean.getClass().getAnnotation(JobHandler.class).value();
                    IJobHandler handler = (IJobHandler) serviceBean;
                    if (loadJobHandler(name) != null) {
                        throw new RuntimeException("xxl-job jobhandler[" + name + "] naming conflicts.");
                    }
                    registJobHandler(name, handler);
                }
            }
        }
    }*/

    private void initJobHandlerMethodRepository(ApplicationContext applicationContext) {
        if (applicationContext == null) {
            return;
        }
        // init job handler from method
        //这里是获取springioc容器的 bean对象
        String[] beanDefinitionNames = applicationContext.getBeanNamesForType(Object.class, false, true);
        for (String beanDefinitionName : beanDefinitionNames) {
            //获取bean
            Object bean = applicationContext.getBean(beanDefinitionName);

            Map<Method, XxlJob> annotatedMethods = null;   // referred to :org.springframework.context.event.EventListenerMethodProcessor.processBean
            try {
                //从Bean中判断是否有方法使用了@Xxljob注解
                annotatedMethods = MethodIntrospector.selectMethods(bean.getClass(),
                        new MethodIntrospector.MetadataLookup<XxlJob>() {
                            @Override
                            public XxlJob inspect(Method method) {
                                return AnnotatedElementUtils.findMergedAnnotation(method, XxlJob.class);
                            }
                        });
            } catch (Throwable ex) {
                logger.error("xxl-job method-jobhandler resolve error for bean[" + beanDefinitionName + "].", ex);
            }
            //没有获取到开始下个循环
            if (annotatedMethods==null || annotatedMethods.isEmpty()) {
                continue;
            }

            for (Map.Entry<Method, XxlJob> methodXxlJobEntry : annotatedMethods.entrySet()) {
                //获取xxljob对应的method对象
                Method method = methodXxlJobEntry.getKey();
                //获取xxljob注解对象
                XxlJob xxlJob = methodXxlJobEntry.getValue();
                if (xxlJob == null) {
                    continue;
                }
                //对获取name
                String name = xxlJob.value();
                if (name.trim().length() == 0) {
                    throw new RuntimeException("xxl-job method-jobhandler name invalid, for[" + bean.getClass() + "#" + method.getName() + "] .");
                }
                //验证name的唯一性 
                if (loadJobHandler(name) != null) {
                    throw new RuntimeException("xxl-job jobhandler[" + name + "] naming conflicts.");
                }

                // execute method
                //日志记录
                if (!(method.getParameterTypes().length == 1 && method.getParameterTypes()[0].isAssignableFrom(String.class))) {
                    throw new RuntimeException("xxl-job method-jobhandler param-classtype invalid, for[" + bean.getClass() + "#" + method.getName() + "] , " +
                            "The correct method format like \" public ReturnT<String> execute(String param) \" .");
                }
                if (!method.getReturnType().isAssignableFrom(ReturnT.class)) {
                    throw new RuntimeException("xxl-job method-jobhandler return-classtype invalid, for[" + bean.getClass() + "#" + method.getName() + "] , " +
                            "The correct method format like \" public ReturnT<String> execute(String param) \" .");
                }
                method.setAccessible(true);

                // init and destory
                Method initMethod = null;
                Method destroyMethod = null;
                //根据init属性判断是否存在对应init方法
                if (xxlJob.init().trim().length() > 0) {
                    try {
                        initMethod = bean.getClass().getDeclaredMethod(xxlJob.init());
                        initMethod.setAccessible(true);
                    } catch (NoSuchMethodException e) {
                        throw new RuntimeException("xxl-job method-jobhandler initMethod invalid, for[" + bean.getClass() + "#" + method.getName() + "] .");
                    }
                }
                //根据init属性判断是否存在对应destroy方法

                if (xxlJob.destroy().trim().length() > 0) {
                    try {
                        destroyMethod = bean.getClass().getDeclaredMethod(xxlJob.destroy());
                        destroyMethod.setAccessible(true);
                    } catch (NoSuchMethodException e) {
                        throw new RuntimeException("xxl-job method-jobhandler destroyMethod invalid, for[" + bean.getClass() + "#" + method.getName() + "] .");
                    }
                }

                // registry jobhandler
                registJobHandler(name, new MethodJobHandler(bean, method, initMethod, destroyMethod));
            }
        }

    }

    // ---------------------- applicationContext ----------------------
    private static ApplicationContext applicationContext;

    @Override
    public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
        this.applicationContext = applicationContext;
    }

    public static ApplicationContext getApplicationContext() {
        return applicationContext;
    }

}

 注册jobhandler,实际上写入到本地缓存中jobHandlerRepository,如下所示:

 // ---------------------- job handler repository ----------------------
    private static ConcurrentMap<String, IJobHandler> jobHandlerRepository = new ConcurrentHashMap<String, IJobHandler>();
    public static IJobHandler registJobHandler(String name, IJobHandler jobHandler){
        logger.info(">>>>>>>>>>> xxl-job register jobhandler success, name:{}, jobHandler:{}", name, jobHandler);
        return jobHandlerRepository.put(name, jobHandler);
    }

xxljob的调用时基于Netty的,客户端在初始化的时候会本地构建一个Netty服务端(详细可参考NettyHttpServer),并与admin保持一个长连接,在Admin来被定时调度的触发后通过Netty发送调度的信息给客户端(详细参考NettyServerHandler与XxlRpcProviderFactory),客户端收到信息后,拿着信息去jobHandlerRepository中找到对应的IjobHandler对象,执行对应的方法

网上说的一些MDC加traceId的方法,个人测试发现,是不好用的

@Aspect
@Component
public class JobTraceLogAspect {
    @Before("execution(public * com.xxl.job..*.execute(String))")
    public void beforeMethod(JoinPoint joinPoint) {
        String traceId = UUID.randomUUID().toString().replace("-", "");
        MDC.put("traceId", traceId);
    }
}

更多推荐