diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java index f80d3a1d2..d788dce60 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java @@ -23,6 +23,7 @@ import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo; import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil; import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter; +import com.ai.cloud.skywalking.analysis.categorize2chain.util.VersionIdentifier; import com.ai.cloud.skywalking.analysis.config.ConfigInitializer; import com.ai.cloud.skywalking.protocol.Span; @@ -39,6 +40,10 @@ public class Categorize2ChainMapper extends TableMapper { @Override protected void map(ImmutableBytesWritable key, Result value, Context context) throws IOException, InterruptedException { + if(!VersionIdentifier.enableAnaylsis(Bytes.toString(key.get()))){ + return; + } + List spanList = new ArrayList(); ChainInfo chainInfo = null; try { diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/util/VersionIdentifier.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/util/VersionIdentifier.java new file mode 100644 index 000000000..a24b236fd --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/util/VersionIdentifier.java @@ -0,0 +1,31 @@ +package com.ai.cloud.skywalking.analysis.categorize2chain.util; + +/** + * 版本识别器 + * + * @author wusheng + * + */ +public class VersionIdentifier { + /** + * 根据tid识别数据是否可分析
+ * 目前允许分析所有1.x的版本号 + * + * @param tid + * @return + */ + public static boolean enableAnaylsis(String tid){ + if(tid != null){ + String[] tidSections = tid.split("\\."); + if(tidSections.length == 7){ + String version = tidSections[0]; + String subVersion = tidSections[1]; + + if("1".equals(version) && subVersion.length() > 0){ + return true; + } + } + } + return false; + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Chain2SummaryMapper.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Chain2SummaryMapper.java index 0bafb3576..6324d9a89 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Chain2SummaryMapper.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Chain2SummaryMapper.java @@ -17,7 +17,7 @@ import java.io.IOException; public class Chain2SummaryMapper extends TableMapper { private Logger logger = LoggerFactory - .getLogger(Chain2SummaryMapper.class.getName()); + .getLogger(Chain2SummaryMapper.class); @Override @@ -36,7 +36,6 @@ public class Chain2SummaryMapper extends TableMapper { private Logger logger = LoggerFactory - .getLogger(Chain2SummaryReducer.class.getName()); + .getLogger(Chain2SummaryReducer.class); @Override protected void setup(Context context) throws IOException, InterruptedException { @@ -36,7 +36,6 @@ public class Chain2SummaryReducer extends Reducer ThreadTraceIdSequence = new ThreadLocal(); + private static final ThreadLocal ThreadTraceIdSequence = new ThreadLocal(); - private static final String PROCESS_UUID; + private static final String PROCESS_UUID; - static { - String uuid = UUID.randomUUID().toString().replaceAll("-", ""); - PROCESS_UUID = uuid.substring(uuid.length() - 7); - } + static { + String uuid = UUID.randomUUID().toString().replaceAll("-", ""); + PROCESS_UUID = uuid.substring(uuid.length() - 7); + } - private TraceIdGenerator() { - } + private TraceIdGenerator() { + } - public static String generate() { - Integer seq = ThreadTraceIdSequence.get(); - if (seq == null || seq == 10000 || seq > 10000) { - seq = 0; - } - seq++; - ThreadTraceIdSequence.set(seq); + /** + * TraceId由以下规则组成
+ * 2位version号 + 1位时间戳(毫秒数) + 1位进程随机号(UUID后7位) + 1位进程数号 + 1位线程号 + 1位线程内序号 + * + * 注意:这里的位,是指“.”作为分隔符所占的位数,非字符串长度的位数。 + * TraceId为不定长字符串,但保证在分布式集群条件下的唯一性 + * + * @return + */ + public static String generate() { + Integer seq = ThreadTraceIdSequence.get(); + if (seq == null || seq == 10000 || seq > 10000) { + seq = 0; + } + seq++; + ThreadTraceIdSequence.set(seq); - return Constants.SDK_VERSION + "." + System.currentTimeMillis() - + "."+ PROCESS_UUID - + "."+ BuriedPointMachineUtil.getProcessNo() - + "."+ Thread.currentThread().getId() - + "."+ seq; - } + return Constants.SDK_VERSION + + "." + System.currentTimeMillis() + + "." + PROCESS_UUID + + "." + BuriedPointMachineUtil.getProcessNo() + + "." + Thread.currentThread().getId() + + "." + seq; + } }