XxlJob添加Sleuth链路追踪
Sleuth在XXL-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);
}
}
更多推荐


所有评论(0)