作者:京东物流 张小龙
本文通过介绍分布式应用下各个场景的全局日志ID透传思路,以及介绍分布式日志追踪ID简单实现原理和实战效果,从而达到通过提高日志查询排查问题的效率。
背景
开发排查系统问题用得最多的手段就是查看系统日志,相信不少人都值过班当过小秘吧:给下接口和出入参吧,麻烦看看日志里的有没有异常信息啊等等,但是在并发大时使用日志定位问题还是比较麻烦,由于大量的其他用户/其他线程的日志也一起输出穿行其中导致很难筛选出指定请求的全部相关日志,以及下游线程/服务对应的日志,甚至一些特殊场景的出入参只打印了一些诸如gis坐标、四级地址等没有单据信息的日志,使得日志定位起来非常不便
场景分析
自己所在组负责的系统主要是web应用,其中涉及到的请求方式主要有:springmvc的servlet的http场景、jsf场景、MQ场景、resteasy场景、clover场景、easyjob场景,每一种场景都需要不同的方式进行logTraceId的透传,接下来逐个探析上述各个场景的透传方案。
在这之前我们先要简单了解一下日志中透传和打印logTraceId的方式,一般我们使用MDC进行logTraceId的透传与打印,但是基于MDC内部使用的是ThreadLocal所以只有本线程才有效,子线程服务的MDC里的值会丢失,所以这里我们要么是在所有涉及到父子线程的地方以编码侵入式自行实现值的传递,要么就是通过覆写MDCAdapter:通过阿里的TransmittableThreadLocal来解决父子线程传递问题,而本文采用的是比较粗糙地以编码侵入式来解决此问题。
springmvc的servlet的http场景
这个场景相信大家都已经烂熟到骨子里了,主要思路是通过拦截器的方式进行logTraceId的透传,新建一个类实现HandlerInterceptor
preHandle:在业务处理器处理请求之前被调用,这里实现logTraceId的设置与透传
postHandle:在业务处理器处理请求执行完成后,生成视图之前执行,这里空实现就好
afterCompletion:在DispatcherServlet完全处理完请求后被调用,这里用于清除MDC的logTraceId
1@Slf4j 2public class TraceInterceptor implements HandlerInterceptor { 3 @Override 4 public boolean preHandle(HttpServletRequest request, HttpServletResponse response, Object o) throws Exception { 5 try{ 6 String traceId = MDC.get(LogConstants.MDC_LOG_TRACE_ID_KEY); 7 if (StringUtils.isBlank(traceId)) { 8 MDC.put(LogConstants.MDC_LOG_TRACE_ID_KEY, TraceUtils.getTraceId()); 9 } 10 }catch (RuntimeException e){ 11 log.error("mvc自定义log跟踪拦截器执行异常",e); 12 } 13 return true; 14 } 15 16 @Override 17 public void postHandle(javax.servlet.http.HttpServletRequest httpServletRequest, javax.servlet.http.HttpServletResponse httpServletResponse, Object o, ModelAndView modelAndView) throws Exception { 18 } 19 20 @Override 21 public void afterCompletion(javax.servlet.http.HttpServletRequest httpServletRequest, javax.servlet.http.HttpServletResponse httpServletResponse, Object o, Exception e) throws Exception { 22 try{ 23 MDC.clear(); 24 }catch (RuntimeException ex){ 25 log.error("mvc自定义log跟踪拦截器执行异常",ex); 26 } 27 } 28}
jsf场景
相信大家对于jsf并不陌生,而jsf也支持自定义filter,基于jsf过滤器的运行方式(如下图),可以通过配置全局过滤器(继承AbstractFilter)的方式进行logTraceId的透传,需要注意的是jsf是在线程池中执行的所以一定要信任消息体中的logTraceId

jsf消费者过滤器:主要从上下文环境中获取logTraceId并进行透传,实现代码如下
1@Slf4j 2public class TraceIdGlobalJsfFilter extends AbstractFilter { 3 @Override 4 public ResponseMessage invoke(RequestMessage requestMessage) { 5 //设置traceId 6 setAndGetTraceId(requestMessage); 7 try{ 8 return this.getNext().invoke(requestMessage); 9 }finally { 10 } 11 } 12 13 /** 14 * 设置并返回traceId 15 * @param requestMessage 16 * @return 17 */ 18 private void setAndGetTraceId(RequestMessage requestMessage) { 19 try{ 20 String logTraceId = MDC.get(LogConstants.MDC_LOG_TRACE_ID_KEY); 21 Object logTraceIdObj = requestMessage.getInvocationBody().getAttachment(LogConstants.JSF_LOG_TRACE_ID_KEY); 22 if(StringUtils.isBlank(logTraceId) && logTraceIdObj == null){ 23 //如果filter和MDC都没有获取到则说明有遗漏,打印日志 24 if(log.isDebugEnabled()){ 25 log.debug("jsf消费者自定义log跟踪拦截器预警,filter和MDC都没有traceId,jsf信息:{}", JSON.toJSONString(requestMessage)); 26 } 27 } else if(StringUtils.isBlank(logTraceId) && logTraceIdObj != null) { 28 //如果MDC没有,filter有,打印日志 29 if(log.isDebugEnabled()){ 30 log.debug("jsf消费者自定义log跟踪拦截器预警,MDC没有filter有traceId,jsf信息:{}", JSON.toJSONString(requestMessage)); 31 } 32 } else if(StringUtils.isNotBlank(logTraceId) && logTraceIdObj == null){ 33 //如果MDC有,filter没有,说明是源头已经有了,但是jsf是第一次调,透传 34 requestMessage.getInvocationBody().addAttachment(LogConstants.JSF_LOG_TRACE_ID_KEY, logTraceId); 35 }else if(StringUtils.isNotBlank(logTraceId) && logTraceIdObj != null){ 36 //MDC和fitler都有,但是并不相等,则存在问题打印日志 37 if(log.isDebugEnabled()){ 38 log.debug("jsf消费者自定义log跟踪拦截器预警,MDC和filter都有traceId,jsf信息:{}", JSON.toJSONString(requestMessage)); 39 } 40 } 41 }catch (RuntimeException e){ 42 log.error("jsf消费者自定义log跟踪拦截器执行异常",e); 43 } 44 } 45}
jsf提供者过滤器:通过拿到消费者在消息体中透传的logTraceId来实现,实现代码如下
1@Slf4j 2public class TraceIdGlobalJsfProducerFilter extends AbstractFilter { 3 @Override 4 public ResponseMessage invoke(RequestMessage requestMessage) { 5 //设置traceId 6 boolean isNeedClearMdc = transferTraceId(requestMessage); 7 try{ 8 return this.getNext().invoke(requestMessage); 9 }finally { 10 if(isNeedClearMdc){ 11 clear(); 12 } 13 } 14 } 15 /** 16 * 设置并返回traceId 17 * @param requestMessage 18 * @return 19 */ 20 private boolean transferTraceId(RequestMessage requestMessage) { 21 boolean isNeedClearMdc = false; 22 try{ 23 String logTraceId = MDC.get(LogConstants.MDC_LOG_TRACE_ID_KEY); 24 Object logTraceIdObj = requestMessage.getInvocationBody().getAttachment(LogConstants.JSF_LOG_TRACE_ID_KEY); 25 if(StringUtils.isBlank(logTraceId) && logTraceIdObj == null){ 26 //如果filter和MDC都没有获取到,说明存在遗漏场景或是提供给外部系统调用的接口,打印日志进行观察 27 String traceId = TraceUtils.getTraceId(); 28 MDC.put(LogConstants.MDC_LOG_TRACE_ID_KEY,traceId); 29 requestMessage.getInvocationBody().addAttachment(LogConstants.JSF_LOG_TRACE_ID_KEY, traceId); 30 if(log.isDebugEnabled()){ 31 log.debug("jsf生产者自定义log跟踪拦截器预警,filter和MDC都没有traceId,jsf信息:{}", JSON.toJSONString(requestMessage)); 32 } 33 isNeedClearMdc = true; 34 } else if(StringUtils.isBlank(logTraceId) && logTraceIdObj != null) { 35 //如果MDC没有,filter有,说明是被调用方,需要透传下去 36 MDC.put(LogConstants.MDC_LOG_TRACE_ID_KEY,logTraceIdObj.toString()); 37 isNeedClearMdc = true; 38 } else if(StringUtils.isNotBlank(logTraceId) && logTraceIdObj == null){ 39 //如果MDC有,filter没有,存在问题,打印日志 40 if(log.isDebugEnabled()){ 41 log.debug("jsf生产者自定义log跟踪拦截器预警,MDC有filter没有traceId,jsf信息:{}", JSON.toJSONString(requestMessage)); 42 } 43 isNeedClearMdc = true; 44 }else if(StringUtils.isNotBlank(logTraceId) && logTraceIdObj != null && !logTraceId.equals(logTraceIdObj.toString())){ 45 //MDC和fitler都有,但是并不相等,则信任filter透传结果 46 TraceUtils.resetTraceId(logTraceIdObj.toString()); 47 if(log.isDebugEnabled()){ 48 log.debug("jsf生产者自定义log跟踪拦截器预警,MDC和fitler都有traceId,但是并不相等,jsf信息:{}", JSON.toJSONString(requestMessage)); 49 } 50 } 51 return isNeedClearMdc; 52 }catch (RuntimeException e){ 53 log.error("jsf生产者自定义log跟踪拦截器执行异常",e); 54 return false; 55 } 56 } 57 58 /** 59 * 清除MDC 60 */ 61 private void clear() { 62 try{ 63 MDC.clear(); 64 }catch (RuntimeException e){ 65 log.error("jsf生产者自定义log跟踪拦截器执行异常",e); 66 } 67 } 68}
MQ场景
说到MQ相信大家对于此就更不陌生了,此种场景主要通过在提供者发送消息时拿到上下文中的logTraceId,将其以扩展信息的方式设置进消息体中进行透传,而消费者则从消息体中进行获取
生产者:新建一个抽象类继承MessageProducer,覆写父类中的两个send方法(批量发送、单条发送),send方法中主要调用抽象加工消息体的方法(logTraceId属性赋值)和日志打印,在子类中进行发送前对消息体的加工处理,具体代码如下
1@Slf4j 2public abstract class BaseTraceIdProducer extends MessageProducer { 3 4 private static final String SEPARATOR_COMMA = ","; 5 6 public BaseTraceIdProducer() { 7 } 8 9 public BaseTraceIdProducer(TransportManager transportManager) { 10 super(transportManager); 11 } 12 13 /** 14 * 获取消息体-单个 15 * @param messageContext 16 * @return 17 */ 18 protected abstract Message getMessage(MessageContext messageContext); 19 20 /** 获取消息体-批量 21 * 22 * @param messageContext 23 * @return 24 */ 25 protected abstract List<Message> getMessages(MessageContext messageContext); 26 27 /** 28 * 填充消息体上下文信息 29 * @param message 30 * @param messageContext 31 */ 32 protected void fillContext(Message message,MessageContext messageContext) { 33 if(message == null){ 34 return; 35 } 36 if(StringUtils.isBlank(messageContext.getLogTraceId())){ 37 String logTraceId = message.getAttribute(LogConstants.JMQ2_LOG_TRACE_ID_KEY); 38 messageContext.setLogTraceId(logTraceId); 39 } 40 if(StringUtils.isBlank(messageContext.getTopic())){ 41 String topic = message.getTopic(); 42 messageContext.setTopic(topic); 43 } 44 String businessId = message.getBusinessId(); 45 messageContext.getBusinessIdBuf().append(SEPARATOR_COMMA).append(businessId); 46 } 47 48 /** 49 * traceId嵌入消息体中 50 * @param message 51 */ 52 protected void generateTraceIdIntoMessage(Message message){ 53 if(message == null){ 54 return; 55 } 56 try{ 57 String logTraceId = MDC.get(LogConstants.MDC_LOG_TRACE_ID_KEY); 58 if(StringUtils.isBlank(logTraceId)){ 59 logTraceId = TraceUtils.getTraceId(); 60 MDC.put(LogConstants.MDC_LOG_TRACE_ID_KEY,logTraceId); 61 } 62 message.setAttribute(LogConstants.JMQ2_LOG_TRACE_ID_KEY,logTraceId); 63 }catch (RuntimeException e){ 64 log.error("jmq2自定义log跟踪拦截器执行异常",e); 65 } 66 } 67 68 /** 69 * 批量发送消息-无回调 70 * @param messages 71 * @param timeout 72 * @throws JMQException 73 */ 74 public void send(List<Message> messages, int timeout) throws JMQException { 75 MessageContext messageContext = new MessageContext(); 76 messageContext.setMessages(messages); 77 List<Message> messageList = this.getMessages(messageContext); 78 //打印日志,方便排查问题 79 printLog(messageContext); 80 super.send(messageList, timeout); 81 } 82 83 /** 84 * 单个发送消息 85 * @param message 86 * @param transaction 87 * @param <T> 88 * @return 89 * @throws JMQException 90 */ 91 public <T> T send(Message message, LocalTransaction<T> transaction) throws JMQException { 92 MessageContext messageContext = new MessageContext(); 93 messageContext.setMessage(message); 94 Message msg = this.getMessage(messageContext); 95 //打印日志,方便排查问题 96 printLog(messageContext); 97 return super.send(msg, transaction); 98 } 99 100 /** 101 * 批量发送消息-有回调 102 * @param messages 103 * @param timeout 104 * @param callback 105 * @throws JMQException 106 */ 107 public void send(List<Message> messages, int timeout, AsyncSendCallback callback) throws JMQException { 108 MessageContext messageContext = new MessageContext(); 109 messageContext.setMessages(messages); 110 List<Message> messageList = this.getMessages(messageContext); 111 //打印日志,方便排查问题 112 printLog(messageContext); 113 super.send(messageList, timeout, callback); 114 } 115 116 /** 117 * 打印日志,方便排查问题 118 * @param messageContext 119 */ 120 private void printLog(MessageContext messageContext) { 121 if(messageContext==null){ 122 return; 123 } 124 if(log.isInfoEnabled()){ 125 log.info("MQ发送:traceId:{},topic:{},businessIds:[{}]",messageContext.getLogTraceId(),messageContext.getTopic(),messageContext.getBusinessIdBuf()==null?"":messageContext.getBusinessIdBuf().toString()); 126 } 127 } 128 129} 130 131@Slf4j 132public class TraceIdEnvMessageProducer extends BaseTraceIdProducer { 133 134 private static final String UAT_TRUE = String.valueOf(true); 135 private boolean uat = false; 136 137 public TraceIdEnvMessageProducer() { 138 } 139 140 public TraceIdEnvMessageProducer(TransportManager transportManager) { 141 super(transportManager); 142 } 143 144 /** 145 * 环境变量打标-单个消息体 146 * @param message 147 */ 148 private void convertUatMessage(Message message) { 149 if (message != null) { 150 message.setAttribute(SplitMessage.JMQ_SPLIT_KEY_IS_UAT, UAT_TRUE); 151 } 152 } 153 154 155 /** 156 * 消息转换-批量消息体 157 * @param messageContext 158 * @return 159 */ 160 private List<Message> convertMessages(MessageContext messageContext) { 161 List<Message> messages = messageContext.getMessages(); 162 if (!CollectionUtils.isEmpty(messages)) { 163 Iterator messageIterator = messages.iterator(); 164 while(messageIterator.hasNext()) { 165 Message message = (Message)messageIterator.next(); 166 if(this.isUat()){ 167 this.convertUatMessage(message); 168 } 169 super.generateTraceIdIntoMessage(message); 170 super.fillContext(message,messageContext); 171 } 172 } 173 return messageContext.getMessages(); 174 } 175 176 /** 177 * 消息转换-单个消息体 178 * @param messageContext 179 * @return 180 */ 181 private Message convertMessage(MessageContext messageContext){ 182 Message message = messageContext.getMessage(); 183 if(this.isUat()){ 184 this.convertUatMessage(message); 185 } 186 super.generateTraceIdIntoMessage(message); 187 super.fillContext(message,messageContext); 188 return message; 189 } 190 191 protected Message getMessage(MessageContext messageContext) { 192 if(log.isDebugEnabled()){ 193 log.debug("current environment is UAT : {}", this.isUat()); 194 } 195 return this.convertMessage(messageContext); 196 } 197 198 protected List<Message> getMessages(MessageContext messageContext) { 199 if(log.isDebugEnabled()){ 200 log.debug("current environment is UAT : {}", this.isUat()); 201 } 202 return this.convertMessages(messageContext); 203 } 204 205 public void setUat(boolean uat) { 206 this.uat = uat; 207 } 208 209 boolean isUat() { 210 return this.uat; 211 } 212 213}
消费者:新建一个抽象类继承MessageListener,覆写父类中的onMessage方法,主要进行设置日志traceId和消费完成后的traceId清理等,而在子类中进行一些自定义处理,具体代码如下
1@Slf4j 2public abstract class BaseTraceIdMessageListener implements MessageListener { 3 4 public BaseTraceIdMessageListener() { 5 } 6 7 public abstract void onMessageList(List<Message> messages) throws Exception; 8 9 @Override 10 public final void onMessage(List<Message> messages) throws Exception { 11 try{ 12 if(CollectionUtils.isEmpty(messages)){ 13 return; 14 } 15 //设置日志traceId 16 setLogTraceId(messages); 17 this.onMessageList(messages); 18 //消费完后清除traceId 19 clear(); 20 }catch (Exception e){ 21 throw e; 22 }finally { 23 MDC.clear(); 24 } 25 } 26 27 /** 28 * 设置日志traceId 29 * @param messages 30 */ 31 private void setLogTraceId(List<Message> messages) { 32 try{ 33 Message message = messages.get(0); 34 String logTraceId = message.getAttribute(LogConstants.JMQ2_LOG_TRACE_ID_KEY); 35 if(StringUtils.isBlank(logTraceId)){ 36 logTraceId = TraceUtils.getTraceId(); 37 } 38 MDC.put(LogConstants.MDC_LOG_TRACE_ID_KEY,logTraceId); 39 }catch (RuntimeException e){ 40 log.error("jmq2自定义log跟踪拦截器执行异常",e); 41 } 42 } 43 44 /** 45 * 清除traceId 46 */ 47 private void clear() { 48 try{ 49 MDC.clear(); 50 }catch (RuntimeException e){ 51 log.error("jmq2自定义log跟踪拦截器执行异常",e); 52 } 53 } 54 55} 56 57@Slf4j 58public abstract class TraceIdEnvMessageListener extends BaseTraceIdMessageListener{ 59 60 private String uat; 61 62 public TraceIdEnvMessageListener() { 63 } 64 65 public abstract void onMessages(List<Message> var1) throws Exception; 66 67 @Override 68 public void onMessageList(List<Message> messages) throws Exception { 69 Iterator iterator; 70 Message message; 71 if (this.getUat() != null && Boolean.valueOf(this.getUat())) { 72 iterator = messages.iterator(); 73 74 while(true) { 75 while(iterator.hasNext()) { 76 message = (Message)iterator.next(); 77 if (message != null && Boolean.valueOf(message.getAttribute(SplitMessage.JMQ_SPLIT_KEY_IS_UAT))) { 78 this.onMessages(Arrays.asList(message)); 79 } else { 80 log.debug("Ignore message: [BusinessId: {}, Text: {}]", message.getBusinessId(), message.getText()); 81 } 82 } 83 84 return; 85 } 86 } else if (this.getUat() != null && !Boolean.valueOf(this.getUat())) { 87 iterator = messages.iterator(); 88 89 while(true) { 90 while(iterator.hasNext()) { 91 message = (Message)iterator.next(); 92 if (message != null && !Boolean.valueOf(message.getAttribute(SplitMessage.JMQ_SPLIT_KEY_IS_UAT))) { 93 this.onMessages(Arrays.asList(message)); 94 } else { 95 log.debug("Ignore message: [BusinessId: {}, Text: {}]", message.getBusinessId(), message.getText()); 96 } 97 } 98 99 return; 100 } 101 } else { 102 this.onMessages(messages); 103 } 104 } 105 106 public void setUat(String uat) { 107 if (!"true".equals(uat) && !"false".equals(uat)) { 108 throw new IllegalArgumentException("uat 属性值只能为 true 或 false."); 109 } else { 110 this.uat = uat; 111 } 112 } 113 114 public String getUat() { 115 return this.uat; 116 } 117}
resteasy场景
此场景类似于spinrg-mvc场景,也是http请求,需要通过拦截器在消息头中进行logTraceId的透传,主要有客户端拦截器,服务端:预处理拦截器、后置拦截器,代码如下
1@ClientInterceptor 2@Provider 3@Slf4j 4public class ResteasyClientInterceptor implements ClientExecutionInterceptor { 5 @Override 6 public ClientResponse execute(ClientExecutionContext clientExecutionContext) throws Exception { 7 try{ 8 String logTraceId = MDC.get(LogConstants.MDC_LOG_TRACE_ID_KEY); 9 ClientRequest request = clientExecutionContext.getRequest(); 10 String headerTraceId = request.getHeaders().getFirst(LogConstants.HEADER_LOG_TRACE_ID_KEY); 11 if(StringUtils.isBlank(logTraceId) && StringUtils.isBlank(headerTraceId)){ 12 //如果filter和MDC都没有获取到则说明是调用源头 13 String traceId = TraceUtils.getTraceId(); 14 TraceUtils.resetTraceId(traceId); 15 request.header(LogConstants.HEADER_LOG_TRACE_ID_KEY,traceId); 16 } else if(StringUtils.isBlank(headerTraceId)){ 17 //如果MDC有但是filter没有则需要传递 18 request.header(LogConstants.HEADER_LOG_TRACE_ID_KEY,logTraceId); 19 } 20 }catch (RuntimeException e){ 21 log.error("resteasy客户端log跟踪拦截器执行异常",e); 22 } 23 return clientExecutionContext.proceed(); 24 } 25} 26 27@Slf4j 28@Provider 29@ServerInterceptor 30public class RestEasyPreInterceptor implements PreProcessInterceptor { 31 @Override 32 public ServerResponse preProcess(HttpRequest request, ResourceMethod resourceMethod) throws Failure, WebApplicationException { 33 try{ 34 MultivaluedMap<String, String> requestHeaders = request.getHttpHeaders().getRequestHeaders(); 35 String headerTraceId = requestHeaders.getFirst(LogConstants.HEADER_LOG_TRACE_ID_KEY); 36 if(StringUtils.isNotBlank(headerTraceId)){ 37 //如果filter则透传 38 TraceUtils.resetTraceId(headerTraceId); 39 } 40 }catch (RuntimeException e){ 41 log.error("resteasy服务端log跟踪前置拦截器执行异常",e); 42 } 43 return null; 44 } 45} 46 47@Slf4j 48@Provider 49@ServerInterceptor 50public class ResteasyPostInterceptor implements PostProcessInterceptor { 51 @Override 52 public void postProcess(ServerResponse serverResponse) { 53 try{ 54 MDC.clear(); 55 }catch (RuntimeException e){ 56 log.error("resteasy服务端log跟踪后置拦截器执行异常",e); 57 } 58 } 59}
clover场景
clover的大体机制主要是在项目启动的时候扫描到带有注解@HessianWebService的类进行服务注册并维持心跳检测,而clover端则通过servlet请求方式进行任务的回调,同时继承AbstractScheduleTaskProcess方式的任务是以线程池的方式进行业务的处理
基于上述原理我们需要解决两个问题:1.新建一个类继承ServiceExporterServlet,并在web.xml配置中进行servlet配置,代码如下;
1@Slf4j 2public class ServiceExporterTraceIdServlet extends ServiceExporterServlet { 3 4 @Override 5 public void service(ServletRequest req, ServletResponse res) throws ServletException, IOException { 6 try { 7 String traceId = MDC.get("traceId"); 8 if (StringUtils.isBlank(traceId)) { 9 MDC.put("traceId", TraceUtils.getTraceId()); 10 } 11 } catch (Exception e) { 12 log.error("clover请求servlet执行异常", e); 13 } 14 try { 15 super.service(req, res); 16 } catch (Throwable e) { 17 log.error("clover请求servlet执行异常", e); 18 throw e; 19 }finally { 20 try{ 21 MDC.clear(); 22 }catch (RuntimeException ex){ 23 log.error("clover请求servlet执行异常",ex); 24 } 25 } 26 } 27}
2.新建一个抽象类继承AbstractScheduleTaskProcess,在类中以编码形式进行父子线程的透传(可优化:通过覆写MDCAdapter:通过阿里的TransmittableThreadLocal来解决父子线程传递问题),所有任务均改为继承此类,关键代码如下
1try{ 2 traceId = MDC.get(LogConstants.MDC_LOG_TRACE_ID_KEY); 3 if (StringUtils.isBlank(traceId)) { 4 log.warn("clover自定义log跟踪拦截器预警,mdc没有traceId"); 5 } 6 }catch (RuntimeException e){ 7 log.error("clover自定义log跟踪拦截器执行异常",e); 8 } 9 final String logTraceId = traceId; 10 while(iterator.hasNext()) { 11 final List<TcTask> list = (List<TcTask>)iterator.next(); 12 this.executor.submit(new Callable<Object>() { 13 public Object call() throws Exception { 14 try{ 15 if (StringUtils.isNotBlank(logTraceId)) { 16 MDC.put(LogConstants.MDC_LOG_TRACE_ID_KEY, logTraceId); 17 } 18 }catch (RuntimeException e){ 19 log.error("clover自定义log跟踪拦截器执行异常",e); 20 } 21 Object var1; 22 try { 23 if (BaseTcTaskProcessWorker.logger.isInfoEnabled()) { 24 BaseTcTaskProcessWorker.logger.info("正在执行任务[" + this.getClass().getName() + "],条数:" + list.size() + "..."); 25 } 26 27 28 BaseTcTaskProcessWorker.this.executeTasks(list); 29 30 if (BaseTcTaskProcessWorker.logger.isInfoEnabled()) { 31 BaseTcTaskProcessWorker.logger.info("执行任务[" + this.getClass().getName() + "],条数:" + list.size() + "成功!"); 32 } 33 34 var1 = null; 35 } catch (Exception var5) { 36 BaseTcTaskProcessWorker.logger.error(var5.getMessage(), var5); 37 throw var5; 38 } finally { 39 try{ 40 MDC.clear(); 41 }catch (RuntimeException ex){ 42 log.error("clover自定义log跟踪拦截器执行异常",ex); 43 } 44 latch.countDown(); 45 } 46 47 return var1; 48 } 49 }); 50 }
easyjob场景
easyjob的大体机制是在项目启动的时候通过扫描实现接口Scheduler的类进行上报注册,同时启动一个acceptor(获取任务的线程池),而acceptor拉取到任务后会将父任务放进一个叫executor的线程池,子任务范进一个叫slowExecutor的线程池,我们可以新建一个抽奖类实现接口ScheduleFlowTask,复用clover场景硬编码方式进行父子线程logTraceId的透传处理(可优化:通过覆写MDCAdapter:通过阿里的TransmittableThreadLocal来解决父子线程传递问题),示例代码如下
1@Slf4j 2public abstract class AbstractEasyjobOnlyScheduleProcess<T> implements ScheduleFlowTask { 3 4 /** 5 * EASYJOB平台UMP监控key前缀 6 */ 7 private static final String EASYJOB_UMP_KEY_RREFIX = "trans.easyjob.dotask."; 8 9 /** 10 * EASYJOB单个任务处理分布式锁前缀 11 */ 12 private static final String EASYJOB_SINGLE_TASK_LOCK_PREFIX = "basic_easyjob_single_task_lock_prefix_"; 13 14 /** 15 * 环境标识-开关配置进行环境隔离 16 */ 17 @Value("${spring.profiles.active}") 18 private String activeEnv; 19 20 @Value("${task.scene.mark}") 21 private String sceneMark = TaskSceneMarkEnum.PRODUCTION.getDesc(); 22 23 /** 24 * easyJob维度线程池变量 25 */ 26 private ThreadPoolExecutor easyJobExecutor; 27 /** 28 * easyJob维度服务器个数-分片个数 29 */ 30 private volatile int easyJobLastThreadCount = 0; 31 32 /** 33 * easyjob多线程名称 34 */ 35 private static final String EASYJOB_THREAD_NAME = "dts.easyJobs"; 36 37 /** 38 * 子类的泛型参数类型 39 */ 40 private Class<T> argumentType; 41 42 /** 43 * 无参构造 44 */ 45 public AbstractEasyjobOnlyScheduleProcess() { 46 //设置子类泛型参数类型 47 argumentType = this.getArgumentType(); 48 } 49 50 @Autowired 51 private RedisHelper redisHelper; 52 53 /** 54 * 非task表扫描待处理的任务数据 55 * @param taskServerParam 56 * @param curServer 57 * @return 58 */ 59 protected abstract List<T> loadTasks(TaskServerParam taskServerParam, int curServer); 60 61 /** 62 * 业务处理抽象方法-单个 63 * @param task 64 */ 65 protected abstract void doSingleTask(T task); 66 67 /** 68 * 业务处理抽象方法-批量 69 * @param tasks 70 */ 71 protected abstract void doBatchTasks(List<T> tasks); 72 73 /** 74 * 拼装ump监控key 75 * @param prefix 76 * @param taskNameKey 77 * @return 78 */ 79 private String getUmpKey(String prefix,String taskNameKey) { 80 StringBuffer umpKeyBuf = new StringBuffer(); 81 umpKeyBuf.append(prefix).append(taskNameKey); 82 return umpKeyBuf.toString(); 83 } 84 85 /** 86 * easyjob平台异步任务回调方法 87 * @param scheduleContext 88 * @return 89 * @throws Exception 90 */ 91 @Override 92 public TaskResult doTask(ScheduleContext scheduleContext) throws Exception { 93 String requestNo = TraceUtils.getTraceId(); 94 try { 95 String traceId = MDC.get(LogConstants.MDC_LOG_TRACE_ID_KEY); 96 if (StringUtils.isBlank(traceId)) { 97 MDC.put(LogConstants.MDC_LOG_TRACE_ID_KEY, requestNo); 98 } 99 } catch (Exception e) { 100 log.error("easyjob执行异常", e); 101 } 102 EasyJobTaskServerParam taskServerParam = null; 103 104 CallerInfo callerinfo = null; 105 try { 106 //条件转换 107 taskServerParam = EasyJobCoreUtil.transTaskServerParam(scheduleContext); 108 String taskNameKey = getTaskNameKey(); 109 String umpKey = getUmpKey(EASYJOB_UMP_KEY_RREFIX,taskNameKey); 110 callerinfo = Profiler.registerInfo(umpKey, Constants.TRANS_BASIC, false, true); 111 //多服务器,并且非子任务,本次不执行,提交子任务 112 if (taskServerParam.getServerCount() > 1 && !taskServerParam.isSubTask()) { 113 submitSubTask(scheduleContext, taskServerParam,requestNo); 114 return TaskResult.success(); 115 } 116 117 if (log.isInfoEnabled()) { 118 log.info("请求编号[{}],开始获取任务,任务ID[{}],任务名称[{}],执行参数[{}]", requestNo, taskServerParam.getTaskId(), taskServerParam.getTaskName(), JSON.toJSONString(taskServerParam)); 119 } 120 TaskServerParam cloverTaskServerParam = EasyJobCoreUtil.transferCloverTaskServerParam(taskServerParam); 121 122 List<T> tasks = this.selectTasks(cloverTaskServerParam, taskServerParam.getCurServer()); 123 124 if (log.isInfoEnabled()) { 125 log.info("请求编号[{}],获取任务ID[{}],任务名称[{}]共{}条", requestNo, taskServerParam.getTaskId(), taskServerParam.getTaskName(), tasks == null ? 0 : tasks.size()); 126 } 127 128 if (CollectionUtils.isNotEmpty(tasks)) { 129 if (log.isInfoEnabled()) { 130 log.info("请求编号[{}],开始执行任务,任务ID[{}],任务名称[{}]", requestNo, taskServerParam.getTaskId(), taskServerParam.getTaskName()); 131 } 132 133 this.easyJobExecuteTasksInner(taskServerParam, tasks,requestNo); 134 if (log.isInfoEnabled()) { 135 log.info("请求编号[{}],执行任务,任务ID[{}],任务名称[{}],执行数量[{}]完成....", requestNo, taskServerParam.getTaskId(), taskServerParam.getTaskName(), tasks.size()); 136 } 137 138 } 139 return TaskResult.success(); 140 } catch (Exception e) { 141 Profiler.functionError(callerinfo); 142 if (log.isInfoEnabled()) { 143 log.error("请求编号[{}],任务执行失败,任务ID[{}],任务名称[{}]", requestNo, taskServerParam == null ? "" : taskServerParam.getTaskId(), taskServerParam == null ? "" :taskServerParam.getTaskName(), e); 144 } 145 return TaskResult.fail(e.getMessage()); 146 }finally { 147 try{ 148 MDC.clear(); 149 }catch (RuntimeException ex){ 150 log.error("easyjob执行异常",ex); 151 } 152 Profiler.registerInfoEnd(callerinfo); 153 } 154 } 155 156 /** 157 * 多分片提交子任务 158 * @param scheduleContext 调度任务上下文参数 159 * @param taskServerParam 调度任务参数 160 * @param requestNo 调度任务参数 161 * @return void 162 */ 163 private void submitSubTask(ScheduleContext scheduleContext, EasyJobTaskServerParam taskServerParam,String requestNo) throws IOException { 164 165 log.info("请求编号[{}],执行任务,任务ID[{}],任务名称[{}],子任务个数[{}],开始提交子任务", requestNo, taskServerParam.getTaskId(), taskServerParam.getTaskName(), taskServerParam.getServerCount()); 166 167 String jobClass = scheduleContext.getTaskGetResponse().getJobClass(); 168 169 if (StringUtils.isBlank(jobClass)) { 170 throw new RuntimeException("jobClass get error"); 171 } 172 173 for (int i = 0; i < taskServerParam.getServerCount(); i++) { 174 Map<String, String> dataMap = scheduleContext.getParameters(); 175 //提交子任务标识 176 dataMap.put("isSubTask", "true"); 177 //给子任务进行编号 178 dataMap.put("curServer", String.valueOf(i)); 179 //父任务名称传递子任务 180 dataMap.put("taskName", taskServerParam.getTaskName()); 181 scheduleContext.commitSubTask(jobClass, dataMap, taskServerParam.getExpected(), taskServerParam.getTransactionalAccept()); 182 } 183 // 父任务等待子任务执行完毕再更改状态,如果执行时间超过等待时间,抛异常 184 //scheduleContext.waitForSubtaskCompleted((long) taskServerParam.getServerCount() * taskServerParam.getExpected()); 185 log.info("请求编号[{}],执行任务,任务ID[{}],任务名称[{}],子任务个数[{}],提交完成....", requestNo, taskServerParam.getTaskId(), taskServerParam.getTaskName(), taskServerParam.getServerCount()); 186 } 187 188 /** 189 * 创建线程池,按配置参数执行task 190 * @param param 执行参数 191 * @param tasks 任务集合 192 * @param requestNoStr 193 * @return void 194 */ 195 private void easyJobExecuteTasksInner(final EasyJobTaskServerParam param, List<T> tasks,String requestNoStr) { 196 int threadCount = param.getThreadCount(); 197 synchronized (this) { 198 if (this.easyJobExecutor == null) { 199 this.easyJobExecutor = (ThreadPoolExecutor) EasyJobCoreUtil.createCustomeasyJobExecutorService(threadCount, EASYJOB_THREAD_NAME); 200 this.easyJobLastThreadCount = threadCount; 201 } else if (threadCount > this.easyJobLastThreadCount) { 202 this.easyJobExecutor.setMaximumPoolSize(threadCount); 203 this.easyJobExecutor.setCorePoolSize(threadCount); 204 this.easyJobLastThreadCount = threadCount; 205 } else if (threadCount < this.easyJobLastThreadCount) { 206 this.easyJobExecutor.setCorePoolSize(threadCount); 207 this.easyJobExecutor.setMaximumPoolSize(threadCount); 208 this.easyJobLastThreadCount = threadCount; 209 } 210 } 211 212 List<List<T>> lists = Lists.partition(tasks, param.getExecuteCount()); 213 final CountDownLatch latch = new CountDownLatch(lists.size()); 214 final String requestNo = requestNoStr; 215 for (final List<T> list : lists) { 216 this.easyJobExecutor.submit( 217 new Callable<Object>() { 218 public Object call() throws Exception { 219 try{ 220 if (StringUtils.isNotBlank(requestNo)) { 221 MDC.put(LogConstants.MDC_LOG_TRACE_ID_KEY, requestNo); 222 } 223 }catch (RuntimeException e){ 224 log.error("easyjob自定义log跟踪拦截器执行异常",e); 225 } 226 try { 227 if (log.isInfoEnabled()) { 228 log.info("请求编号[{}],正在执行任务,任务ID[{}],任务名称[{}],[{}],条数:[{}]...", requestNo, param.getTaskId(), param.getTaskName(), Thread.currentThread().getName(), list.size()); 229 } 230 executeTasks(list); 231 if (log.isInfoEnabled()) { 232 log.info("请求编号[{}],执行任务,任务ID[{}],任务名称[{}],[{}],条数:[{}]成功!", requestNo, param.getTaskId(), param.getTaskName(), Thread.currentThread().getName(), list.size()); 233 } 234 } catch (Exception e) { 235 log.error(e.getMessage(), e); 236 throw e; 237 } finally { 238 try{ 239 MDC.clear(); 240 }catch (RuntimeException ex){ 241 log.error("easyjob自定义log跟踪拦截器执行异常",ex); 242 } 243 latch.countDown(); 244 } 245 return null; 246 } 247 248 } 249 ); 250 } 251 252 try { 253 latch.await(); 254 } catch (InterruptedException e) { 255 throw new RuntimeException("interrupted when processing data access request in concurrency", e); 256 } 257 } 258 259 /** 260 * 获取任务名称 261 * @return 262 */ 263 private String getTaskNameKey(){ 264 StringBuffer keyBuf = new StringBuffer(); 265 keyBuf.append(activeEnv) 266 .append(Constants.SEPARATOR_UNDERLINE) 267 .append(this.getClass().getSimpleName()); 268 return keyBuf.toString(); 269 } 270 271 protected void executeTasks(List<T> taskList) { 272 if(CollectionUtils.isEmpty(taskList)) { 273 return; 274 } 275 this.doTasks(taskList); 276 } 277 278 /** 279 * 业务处理抽象方法 280 * @param list 281 */ 282 protected void doTasks(List<T> list){ 283 if(isDoBatchTasks()){ 284 CallerInfo info = Profiler.registerInfo(getClass().getName()+"_batch", Constants.TRANS_BASIC,false, true); 285 try { 286 /** 开始执行各个子类真正业务逻辑 */ 287 this.doBatchTasks(list); 288 } catch(CommonBusinessException ex){ 289 log.warn(ex.getMessage()); 290 } catch (Exception e) { 291 Profiler.functionError(info); 292 log.error("任务处理失败,方法:{},任务:{}",ClassHelper.getMethod(),JSON.toJSONString(list), e); 293 } finally { 294 Profiler.registerInfoEnd(info); 295 } 296 }else{ 297 for (T task : list) { 298 CallerInfo info = Profiler.registerInfo(getClass().getName(), Constants.TRANS_BASIC,false, true); 299 if(task == null) { continue; } 300 String lockKey = ""; 301 try { 302 /** 开始执行各个子类真正业务逻辑 */ 303 if (useConcurrentLock()) { 304 lockKey = getLockKey(task); 305 if (redisHelper.lock(RedisKeyDef.SyncLockKeyPrefix.TASK_PROCESS_LOCK_PREFIX, lockKey)) { 306 this.doSingleTask(task); 307 }else{ 308 lockKey = ""; 309 log.warn("lockKey:{},加载失败,正在被其他用户锁定,请重试!",lockKey); 310 } 311 } else { 312 this.doSingleTask(task); 313 } 314 } catch(CommonBusinessException ex){ 315 log.warn(ex.getMessage()); 316 } catch (Exception e) { 317 Profiler.functionError(info); 318 log.error("任务处理失败,方法:{},任务:{}",ClassHelper.getMethod(),JSON.toJSONString(task), e); 319 } finally { 320 Profiler.registerInfoEnd(info); 321 if (StringUtils.isNotBlank(lockKey)) { 322 redisHelper.unlock(RedisKeyDef.SyncLockKeyPrefix.TASK_PROCESS_LOCK_PREFIX, lockKey); 323 } 324 } 325 } 326 } 327 } 328 329 /** 330 * 获取实体类的实际类型 331 * 332 * @return 333 */ 334 private Class<T> getArgumentType() { 335 return (Class<T>) ((ParameterizedType) getClass().getGenericSuperclass()).getActualTypeArguments()[0]; 336 } 337 338 /** 339 * 是否使用防并发锁 340 * 默认不使用,如需使用子类重写该方法 341 * @return 342 */ 343 protected boolean useConcurrentLock() { 344 return false; 345 } 346 347 /** 348 * 根所注解获取LockKey,可被子类重写,提高效率 349 * 350 * @param businessObj 业务对象 351 * @return concurrent lock key 352 */ 353 protected String getLockKey( T businessObj) { 354 StringBuilder lockKey = new StringBuilder(EASYJOB_SINGLE_TASK_LOCK_PREFIX); 355 //若存在注解指定的防重字段,则使用这些字段拼装防重Key,否则使用MQ业务主键防重 356 List<ValueEntryInfo> valueEntries = getAnnotaionConcurrentKeys(businessObj); 357 if (!CollectionUtils.isEmpty(valueEntries)) { 358 for (ValueEntryInfo valueEntry : valueEntries) { 359 lockKey.append(Constants.SEPARATOR_UNDERLINE); 360 lockKey.append(valueEntry.getValue()); 361 } 362 } else { 363 throw new CommonBusinessException(String.format("此任务处理需要加分布式锁,但是未设置锁key,所以不做业务处理,请检查,任务信息:%s",JSON.toJSONString(businessObj))); 364 } 365 return lockKey.toString(); 366 } 367 368 /** 369 * 查找对象的ConccurentKey注解,获取防重字段,并排序返回 370 * 371 * @param businessObj 业务对象 372 * @return 有序的业务字段值列表 373 */ 374 private List<ValueEntryInfo> getAnnotaionConcurrentKeys(T businessObj) { 375 List<ValueEntryInfo> valueEntries = new ArrayList<ValueEntryInfo>(); 376 Field[] fields = businessObj.getClass().getDeclaredFields(); 377 for (int i = 0; i < fields.length; i++) { 378 ConcurrentKey concurrentKey = fields[i].getAnnotation(ConcurrentKey.class); 379 if (concurrentKey != null) { 380 fields[i].setAccessible(true); 381 Object fieldVal = null; 382 try { 383 ValueEntryInfo valueEntry = new ValueEntryInfo(); 384 fieldVal = fields[i].get(businessObj); 385 if (fieldVal != null) { 386 valueEntry.setValue(String.format("%1$s", fieldVal)); 387 valueEntry.setOrder(concurrentKey.order()); 388 valueEntries.add(valueEntry); 389 } 390 } catch (IllegalAccessException e) { 391 log.error("IllegalAccess-{}.{}", businessObj.getClass().getName(), fields[i].getName()); 392 } 393 } 394 } 395 if (valueEntries.size() > 1) { 396 //排序ConcurrentKey 397 Collections.sort(valueEntries, new Comparator<ValueEntryInfo>() { 398 @Override 399 public int compare(ValueEntryInfo o1, ValueEntryInfo o2) { 400 if (o1.getOrder() > o2.getOrder()) { 401 return 1; 402 } else if (o1.getOrder() == o2.getOrder()) { 403 return 0; 404 } else { 405 return -1; 406 } 407 } 408 }); 409 } 410 return valueEntries; 411 } 412 413 protected List<T> selectTasks(TaskServerParam taskServerParam, int curServer) { 414 return this.loadTasks(taskServerParam, curServer); 415 } 416 417 /** 418 * 获取select时的任务创建开始时间 419 * @param serverArg 420 * @return 421 */ 422 protected Date getCreateTimeFrom(String serverArg){ 423 return null; 424 } 425 426 /** 427 * 是否以批量方式处理任务 428 * @return 429 */ 430 protected boolean isDoBatchTasks(){ 431 return false; 432 } 433 434}
实战结果
上述所述均为透传ID场景的原理和示例代码,实战效果如下图:调用jsf超时,跨系统查看日志进行排查,得知为慢sql引起


上述大部分场景已经抽出一个通用jar包,详细使用教程见我的另一篇文章:分布式日志追踪ID使用教程
