From f40fd0302d3cb303f64ee69bf22d3bd81d28949a Mon Sep 17 00:00:00 2001 From: ascrutae Date: Sun, 6 Mar 2016 18:46:15 +0800 Subject: [PATCH] =?UTF-8?q?=E6=8F=90=E4=BA=A4=E6=9C=80=E6=96=B0=E7=9A=84ma?= =?UTF-8?q?per-reduce=E4=BB=A3=E7=A0=81=EF=BC=8C=E7=BB=9F=E8=AE=A1?= =?UTF-8?q?=E9=83=A8=E5=88=86=E6=9C=AA=E5=AE=8C=E6=88=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- skywalking-analysis/skywalking-analysis.iml | 117 ------- .../Categorize2ChainMapper.java | 112 ------ .../Categorize2ChainReducer.java | 57 ---- .../entity/CategorizedChainInfo.java | 77 ----- .../ChainNodeSpecificTimeWindowSummary.java | 71 ---- ...ainNodeSpecificTimeWindowSummaryValue.java | 51 --- .../entity/ChainRelationship.java | 152 --------- .../ChainSpecificTimeWindowSummary.java | 60 ---- .../ChainSummaryWithoutRelationship.java | 65 ---- .../entity/UncategorizeChainInfo.java | 69 ---- .../filter/impl/AppendBusinessKeyFilter.java | 16 - .../analysis/chainbuild/CallChainTree.java | 88 +++++ .../analysis/chainbuild/ChainBuildMapper.java | 74 +++- .../chainbuild/ChainBuildReducer.java | 42 +++ .../DBCallChainInfoDao.java | 17 +- .../SpanEntry.java | 4 +- .../entity/BranchTraceSpanNode.java | 39 --- .../entity/CallChainDetail.java} | 14 +- .../entity/ChainNodeForSummary.java | 18 + .../chainbuild/entity/TraceSpanNode.java | 323 ------------------ .../chainbuild/entity/TraceSpanTree.java | 310 ----------------- .../entity/VisualTraceSpanNode.java | 24 -- .../BuildTraceSpanTreeException.java | 13 - .../TraceSpanTreeNotFountException.java | 13 - .../TraceSpanTreeSerializeException.java | 13 - .../filter/SpanNodeProcessChain.java | 2 +- .../filter/SpanNodeProcessFilter.java | 8 +- .../filter/impl/AppendBusinessKeyFilter.java | 16 + .../filter/impl/CopyAttrFilter.java | 10 +- .../filter/impl/ProcessCostTimeFilter.java | 10 +- .../filter/impl/ReplaceAddressFilter.java | 10 +- .../filter/impl/TokenGenerateFilter.java | 10 +- .../chainbuild/po/CallChainTreeNode.java | 46 +++ .../po/ChainInfo.java | 23 +- .../po/ChainNode.java | 2 +- .../util/HBaseUtil.java | 189 ++++------ .../analysis/chainbuild/util/StringUtil.java | 22 -- .../util/SubLevelSpanCostCounter.java | 2 +- .../chainbuild/util/TokenGenerator.java | 10 +- .../chainbuild/util/VersionIdentifier.java | 50 ++- .../analysis/config/HBaseTableMetaData.java | 20 +- 41 files changed, 446 insertions(+), 1823 deletions(-) delete mode 100644 skywalking-analysis/skywalking-analysis.iml delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainReducer.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/CategorizedChainInfo.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainNodeSpecificTimeWindowSummary.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainNodeSpecificTimeWindowSummaryValue.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainRelationship.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainSpecificTimeWindowSummary.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainSummaryWithoutRelationship.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/UncategorizeChainInfo.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/AppendBusinessKeyFilter.java create mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/CallChainTree.java create mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/ChainBuildReducer.java rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/DBCallChainInfoDao.java (82%) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/SpanEntry.java (97%) delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/BranchTraceSpanNode.java rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain/entity/ChainDetail.java => chainbuild/entity/CallChainDetail.java} (76%) create mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/ChainNodeForSummary.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/TraceSpanNode.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/TraceSpanTree.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/VisualTraceSpanNode.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/BuildTraceSpanTreeException.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/TraceSpanTreeNotFountException.java delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/TraceSpanTreeSerializeException.java rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/filter/SpanNodeProcessChain.java (98%) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/filter/SpanNodeProcessFilter.java (64%) create mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/AppendBusinessKeyFilter.java rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/filter/impl/CopyAttrFilter.java (63%) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/filter/impl/ProcessCostTimeFilter.java (72%) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/filter/impl/ReplaceAddressFilter.java (66%) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/filter/impl/TokenGenerateFilter.java (56%) create mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/CallChainTreeNode.java rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/po/ChainInfo.java (88%) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/po/ChainNode.java (98%) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/util/HBaseUtil.java (66%) delete mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/StringUtil.java rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/{categorize2chain => chainbuild}/util/SubLevelSpanCostCounter.java (87%) diff --git a/skywalking-analysis/skywalking-analysis.iml b/skywalking-analysis/skywalking-analysis.iml deleted file mode 100644 index aa8205c55..000000000 --- a/skywalking-analysis/skywalking-analysis.iml +++ /dev/null @@ -1,117 +0,0 @@ - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - \ No newline at end of file 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 deleted file mode 100644 index 14248280e..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java +++ /dev/null @@ -1,112 +0,0 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain; - -import java.io.IOException; -import java.util.ArrayList; -import java.util.Collections; -import java.util.Comparator; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.Map; - -import org.apache.hadoop.hbase.Cell; -import org.apache.hadoop.hbase.client.Result; -import org.apache.hadoop.hbase.io.ImmutableBytesWritable; -import org.apache.hadoop.hbase.mapreduce.TableMapper; -import org.apache.hadoop.hbase.util.Bytes; -import org.apache.hadoop.io.Text; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessChain; -import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter; -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.chainbuild.util.VersionIdentifier; -import com.ai.cloud.skywalking.analysis.config.ConfigInitializer; -import com.ai.cloud.skywalking.protocol.Span; - -public class Categorize2ChainMapper extends TableMapper { - private Logger logger = LoggerFactory - .getLogger(Categorize2ChainMapper.class.getName()); - - @Override - protected void setup(Context context) throws IOException, - InterruptedException { - ConfigInitializer.initialize(); - } - - @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 { - for (Cell cell : value.rawCells()) { - Span span = new Span(Bytes.toString(cell.getValueArray(), - cell.getValueOffset(), cell.getValueLength())); - spanList.add(span); - } - - chainInfo = spanToChainInfo(Bytes.toString(key.get()), spanList); - logger.info("Success convert span to chain info...." - + chainInfo.getCID() + " TraceId : " + Bytes.toString(key.get())); - context.write( - new Text(chainInfo.getUserId() + ":" - + chainInfo.getEntranceNodeToken()), chainInfo); - } catch (Exception e) { - logger.error("Failed to mapper call chain[" + key.toString() + "]", - e); - } - } - - public static ChainInfo spanToChainInfo(String key, List spanList) { - SubLevelSpanCostCounter costMap = new SubLevelSpanCostCounter(); - ChainInfo chainInfo = new ChainInfo(); - Collections.sort(spanList, new Comparator() { - @Override - public int compare(Span span1, Span span2) { - String span1TraceLevel = span1.getParentLevel() + "." - + span1.getLevelId(); - String span2TraceLevel = span2.getParentLevel() + "." - + span2.getLevelId(); - return span1TraceLevel.compareTo(span2TraceLevel); - } - }); - - Map spanEntryMap = mergeSpanDataSet(spanList); - for (Map.Entry entry : spanEntryMap.entrySet()) { - ChainNode chainNode = new ChainNode(); - SpanNodeProcessFilter filter = SpanNodeProcessChain - .getProcessChainByCallType(entry.getValue().getSpanType()); - filter.doFilter(entry.getValue(), chainNode, costMap); - chainInfo.addNodes(chainNode); - } - - chainInfo.generateChainToken(); - HBaseUtil.saveCidTidMapping(key, chainInfo); - return chainInfo; - } - - private static Map mergeSpanDataSet(List spanList) { - Map spanEntryMap = new LinkedHashMap(); - for (int i = spanList.size() - 1; i >= 0; i--) { - Span span = spanList.get(i); - SpanEntry spanEntry = spanEntryMap.get(span.getParentLevel() + "." - + span.getLevelId()); - if (spanEntry == null) { - spanEntry = new SpanEntry(); - spanEntryMap.put( - span.getParentLevel() + "." + span.getLevelId(), - spanEntry); - } - spanEntry.setSpan(span); - } - return spanEntryMap; - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainReducer.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainReducer.java deleted file mode 100644 index 1623ca84a..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainReducer.java +++ /dev/null @@ -1,57 +0,0 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain; - -import java.io.IOException; -import java.util.Iterator; - -import org.apache.hadoop.io.IntWritable; -import org.apache.hadoop.io.Text; -import org.apache.hadoop.mapreduce.Reducer; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainRelationship; -import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainSummaryWithoutRelationship; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo; -import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil; -import com.ai.cloud.skywalking.analysis.config.ConfigInitializer; - -public class Categorize2ChainReducer extends Reducer { - private static Logger logger = LoggerFactory.getLogger(Categorize2ChainReducer.class.getName()); - - @Override - protected void setup(Context context) throws IOException, - InterruptedException { - ConfigInitializer.initialize(); - } - - @Override - protected void reduce(Text key, Iterable values, Context context) throws IOException, InterruptedException { - int totalCount = reduceAction(key.toString(), values.iterator()); - context.write(new Text(key.toString()), new IntWritable(totalCount)); - } - - public static int reduceAction(String key, Iterator chainInfoIterator) throws IOException, InterruptedException { - int totalCount = 0; - try { - ChainRelationship chainRelate = HBaseUtil.loadCallChainRelationship(key.toString()); - ChainSummaryWithoutRelationship summary = new ChainSummaryWithoutRelationship(); - while (chainInfoIterator.hasNext()) { - ChainInfo chainInfo = chainInfoIterator.next(); - try { - chainRelate.categoryChain(chainInfo); - summary.summary(chainInfo); - } catch (Exception e) { - continue; - } - totalCount++; - } - - chainRelate.save(); - summary.save(); - } catch (Exception e) { - logger.error("Failed to reduce key[" + key + "]", e); - } - - return totalCount; - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/CategorizedChainInfo.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/CategorizedChainInfo.java deleted file mode 100644 index 48f7ac12f..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/CategorizedChainInfo.java +++ /dev/null @@ -1,77 +0,0 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.entity; - -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.google.gson.Gson; -import com.google.gson.JsonObject; -import com.google.gson.JsonParser; -import com.google.gson.reflect.TypeToken; - -import java.text.SimpleDateFormat; -import java.util.ArrayList; -import java.util.Date; -import java.util.List; -import java.util.regex.Pattern; - -public class CategorizedChainInfo { - private String cid; - private String chainFullToken; - - private List children_Token; - - public CategorizedChainInfo(ChainInfo chainInfo) { - cid = chainInfo.getCID(); - - StringBuilder stringBuilder = new StringBuilder(); - boolean flag = false; - for (ChainNode chainNode : chainInfo.getNodes()) { - if (flag) { - stringBuilder.append(";"); - } - stringBuilder.append(chainNode.getTraceLevelId() + "-" + chainNode.getNodeToken()); - flag = true; - } - - chainFullToken = stringBuilder.toString(); - children_Token = new ArrayList(); - } - - public CategorizedChainInfo(String value) { - JsonObject jsonObject = (JsonObject) new JsonParser().parse(value); - cid = jsonObject.get("chainToken").getAsString(); - chainFullToken = jsonObject.get("chainFullToken").getAsString(); - children_Token = new Gson().fromJson(jsonObject.get("children_Token"), - new TypeToken>() { - }.getType()); - } - - public String getChainFullToken() { - return chainFullToken; - } - - public boolean isContained(UncategorizeChainInfo uncategorizeChainInfo) { - Pattern pattern = Pattern.compile(uncategorizeChainInfo.getNodeRegEx()); - return pattern.matcher(this.chainFullToken).find(); - } - - public static void main(String[] args) { - System.out.println(new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(new Date(1453507835957L))); - } - - public boolean isAlreadyContained(UncategorizeChainInfo uncategorizeChainInfo) { - return children_Token.contains(uncategorizeChainInfo.getCID()); - } - - public void add(UncategorizeChainInfo uncategorizeChainInfo) { - children_Token.add(uncategorizeChainInfo.getCID()); - } - - @Override - public String toString() { - return new Gson().toJson(this); - } - - public List getChildren_Token() { - return children_Token; - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainNodeSpecificTimeWindowSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainNodeSpecificTimeWindowSummary.java deleted file mode 100644 index 05e81815a..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainNodeSpecificTimeWindowSummary.java +++ /dev/null @@ -1,71 +0,0 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.entity; - -import java.util.HashMap; -import java.util.Map; - -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.google.gson.Gson; -import com.google.gson.JsonObject; -import com.google.gson.JsonParser; -import com.google.gson.reflect.TypeToken; - -public class ChainNodeSpecificTimeWindowSummary { - public static final long INTERVAL = 1L; - - private String traceLevelId; - private String nodeToken; - // key : 分钟 - private Map summerValueMap; - - public static ChainNodeSpecificTimeWindowSummary newInstance(String traceLevelId, String nodeToken) { - ChainNodeSpecificTimeWindowSummary cns = new ChainNodeSpecificTimeWindowSummary(); - cns.traceLevelId = traceLevelId; - cns.nodeToken = nodeToken; - return cns; - } - - private ChainNodeSpecificTimeWindowSummary() { - summerValueMap = new HashMap(); - } - - public ChainNodeSpecificTimeWindowSummary(String value) { - JsonObject jsonObject = (JsonObject) new JsonParser().parse(value); - traceLevelId = jsonObject.get("traceLevelId").getAsString(); - summerValueMap = new Gson().fromJson(jsonObject.get("summerValueMap").toString(), - new TypeToken>() { - }.getType()); - nodeToken = jsonObject.get("nodeToken").getAsString(); - } - - public String getTraceLevelId() { - return traceLevelId; - } - - public void summary(ChainNode node) { - String key = generateKey(node.getStartDate()); - ChainNodeSpecificTimeWindowSummaryValue summaryResult = summerValueMap.get(key); - if (summaryResult == null) { - summaryResult = new ChainNodeSpecificTimeWindowSummaryValue(); - summerValueMap.put(key, summaryResult); - } - summaryResult.summary(node); - } - - private String generateKey(long startTime) { - long minutes = (startTime % (1000 * 60 * 60)) / (1000 * 60); - return String.valueOf(minutes / INTERVAL); - } - - @Override - public String toString() { - return new Gson().toJson(this); - } - - public String getNodeToken() { - return nodeToken; - } - - public Map getSummerValueMap() { - return summerValueMap; - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainNodeSpecificTimeWindowSummaryValue.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainNodeSpecificTimeWindowSummaryValue.java deleted file mode 100644 index 0d69053df..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainNodeSpecificTimeWindowSummaryValue.java +++ /dev/null @@ -1,51 +0,0 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.entity; - -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; - -public class ChainNodeSpecificTimeWindowSummaryValue { - private long totalCall; - private long totalCostTime; - private long correctNumber; - private long humanInterruptionNumber; - - public ChainNodeSpecificTimeWindowSummaryValue() { - totalCall = 0; - totalCostTime = 0; - correctNumber = 0; - humanInterruptionNumber = 0; - } - - public long getTotalCall() { - return totalCall; - } - - public long getTotalCostTime() { - return totalCostTime; - } - - public long getCorrectNumber() { - return correctNumber; - } - - public long getHumanInterruptionNumber() { - return humanInterruptionNumber; - } - - public void summary(ChainNode node) { - totalCall++; - if (node.getStatus() == ChainNode.NodeStatus.NORMAL) { - correctNumber++; - } - if (node.getStatus() == ChainNode.NodeStatus.HUMAN_INTERRUPTION) { - humanInterruptionNumber++; - } - totalCostTime += node.getCost(); - } - - public void accumulate(ChainNodeSpecificTimeWindowSummaryValue value) { - this.totalCall += value.getTotalCall(); - this.correctNumber += value.getCorrectNumber(); - this.totalCostTime += value.getTotalCostTime(); - this.humanInterruptionNumber += value.getHumanInterruptionNumber(); - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainRelationship.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainRelationship.java deleted file mode 100644 index e7e88516e..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainRelationship.java +++ /dev/null @@ -1,152 +0,0 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.entity; - -import java.io.IOException; -import java.sql.SQLException; -import java.util.ArrayList; -import java.util.HashMap; -import java.util.HashSet; -import java.util.Iterator; -import java.util.List; -import java.util.Map; -import java.util.Set; - -import org.apache.hadoop.hbase.client.Put; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo; -import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil; -import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData; -import com.google.gson.GsonBuilder; - -public class ChainRelationship { - private static Logger logger = LoggerFactory.getLogger(ChainRelationship.class.getName()); - - private String key; - private Map categorizedChainInfoMap = new HashMap(); - private Set uncategorizeChainInfoSet = new HashSet(); - private Map chainDetailMap = new HashMap(); - - public ChainRelationship(String key) { - this.key = key; - } - - private void categoryAllUncategorizedChainInfo(CategorizedChainInfo parentChains) { - if (uncategorizeChainInfoSet != null && uncategorizeChainInfoSet.size() > 0) { - Iterator uncategorizeChainInfoIterator = uncategorizeChainInfoSet.iterator(); - while (uncategorizeChainInfoIterator.hasNext()) { - UncategorizeChainInfo uncategorizeChainInfo = uncategorizeChainInfoIterator.next(); - if (parentChains.isContained(uncategorizeChainInfo)) { - parentChains.add(uncategorizeChainInfo); - uncategorizeChainInfoIterator.remove(); - } - } - } - } - - private void try2CategoryUncategorizedChainInfo(UncategorizeChainInfo child) { - boolean isContained = false; - for (Map.Entry entry : categorizedChainInfoMap.entrySet()) { - if (entry.getValue().isAlreadyContained(child)) { - isContained = true; - } else if (entry.getValue().isContained(child)) { - logger.info("There has contained :" + entry.getKey() + " " + child.getCID()); - entry.getValue().add(child); - chainDetailMap.put(child.getCID(), new ChainDetail(child.getChainInfo(), false)); - isContained = true; - } - } - - if (!isContained) { - if (!uncategorizeChainInfoSet.contains(child)) { - chainDetailMap.put(child.getCID(), new ChainDetail(child.getChainInfo(), false)); - uncategorizeChainInfoSet.add(child); - } - } - - } - - private CategorizedChainInfo addCategorizedChain(ChainInfo chainInfo) { - if (!categorizedChainInfoMap.containsKey(chainInfo.getCID())) { - categorizedChainInfoMap.put(chainInfo.getCID(), - new CategorizedChainInfo(chainInfo)); - - chainDetailMap.put(chainInfo.getCID(), new ChainDetail(chainInfo, true)); - } - return categorizedChainInfoMap.get(chainInfo.getCID()); - } - - public void categoryChain(ChainInfo chainInfo) { - if (chainInfo.getChainStatus() == ChainInfo.ChainStatus.NORMAL) { - CategorizedChainInfo categorizedChainInfo = addCategorizedChain(chainInfo); - categoryAllUncategorizedChainInfo(categorizedChainInfo); - } else { - UncategorizeChainInfo uncategorizeChainInfo = new UncategorizeChainInfo(chainInfo); - try2CategoryUncategorizedChainInfo(uncategorizeChainInfo); - } - } - - public void save() throws SQLException, IOException, InterruptedException { - saveChainRelationship(); - saveChainDetail(); - } - - private void saveChainDetail() throws SQLException, IOException, InterruptedException { - List puts = new ArrayList(); - for (Map.Entry entry : chainDetailMap.entrySet()) { - Put put1 = new Put(entry.getKey().getBytes()); - entry.getValue().save(put1); - puts.add(put1); - } - - - try { - HBaseUtil.saveChainDetails(puts); - } catch (IOException e) { - logger.error("Faild to save chain detail to hbase.", e); - throw e; - } catch (InterruptedException e) { - logger.error("Faild to save chain detail to hbase.", e); - throw e; - } - } - - private void saveChainRelationship() throws IOException { - Put put = new Put(getKey().getBytes()); - - put.addColumn(HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.COLUMN_FAMILY_NAME.getBytes(), HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.UNCATEGORIZE_COLUMN_NAME.getBytes() - , new GsonBuilder().excludeFieldsWithoutExposeAnnotation().create().toJson(getUncategorizeChainInfoList()).getBytes()); - - for (Map.Entry entry : getCategorizedChainInfoMap().entrySet()) { - put.addColumn(HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.COLUMN_FAMILY_NAME.getBytes(), entry.getKey().getBytes() - , entry.getValue().toString().getBytes()); - } - - try { - HBaseUtil.saveChainRelationship(put); - } catch (IOException e) { - logger.error("Faild to save chain relationship to hbase.", e); - throw e; - } - } - - public void addCategorizeChain(String qualifierName, CategorizedChainInfo categorizedChainInfo) { - categorizedChainInfoMap.put(qualifierName, categorizedChainInfo); - } - - public String getKey() { - return key; - } - - public Map getCategorizedChainInfoMap() { - return categorizedChainInfoMap; - } - - public Set getUncategorizeChainInfoList() { - return uncategorizeChainInfoSet; - } - - public void addUncategorizeChain(List uncategorizeChainInfos) { - uncategorizeChainInfoSet.addAll(uncategorizeChainInfos); - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainSpecificTimeWindowSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainSpecificTimeWindowSummary.java deleted file mode 100644 index 5c552ec00..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainSpecificTimeWindowSummary.java +++ /dev/null @@ -1,60 +0,0 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.entity; - -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil; -import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData; - -import org.apache.hadoop.hbase.client.Put; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.io.IOException; -import java.util.HashMap; -import java.util.Map; - -public class ChainSpecificTimeWindowSummary { - - private static Logger logger = LoggerFactory.getLogger(ChainSpecificTimeWindowSummary.class.getName()); - /** - * key : cid + uid + 时间窗口 - */ - private Map chainNodeSummaryResultMap; - - public ChainSpecificTimeWindowSummary() { - chainNodeSummaryResultMap = new HashMap(); - } - - public static ChainSpecificTimeWindowSummary load(String cid_uid_time) { - ChainSpecificTimeWindowSummary result = null; - try { - result = HBaseUtil.selectChainSummaryResult(cid_uid_time); - } catch (IOException e) { - logger.error("Failed to load the key[" + cid_uid_time + "] summary result.", e); - } - - if (result == null) { - result = new ChainSpecificTimeWindowSummary(); - } - return result; - } - - public void addNodeSummaryResult(ChainNodeSpecificTimeWindowSummary chainNodeSummaryResult) { - chainNodeSummaryResultMap.put(chainNodeSummaryResult.getTraceLevelId(), chainNodeSummaryResult); - } - - public void summaryNodeValue(ChainNode node) { - String tlid = node.getTraceLevelId(); - ChainNodeSpecificTimeWindowSummary chainNodeSummaryResult = chainNodeSummaryResultMap.get(tlid); - if (chainNodeSummaryResult == null) { - chainNodeSummaryResult = ChainNodeSpecificTimeWindowSummary.newInstance(tlid, node.getNodeToken()); - chainNodeSummaryResultMap.put(tlid, chainNodeSummaryResult); - } - chainNodeSummaryResult.summary(node); - } - - public void save(Put put) { - for (Map.Entry entry : chainNodeSummaryResultMap.entrySet()) { - put.addColumn(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME.getBytes(), entry.getKey().getBytes(), entry.getValue().toString().getBytes()); - } - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainSummaryWithoutRelationship.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainSummaryWithoutRelationship.java deleted file mode 100644 index f047946b1..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainSummaryWithoutRelationship.java +++ /dev/null @@ -1,65 +0,0 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.entity; - -import com.ai.cloud.skywalking.analysis.categorize2chain.DBCallChainInfoDao; -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 org.apache.hadoop.hbase.client.Put; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.io.IOException; -import java.sql.SQLException; -import java.sql.Timestamp; -import java.text.SimpleDateFormat; -import java.util.*; - -public class ChainSummaryWithoutRelationship { - - private static Logger logger = LoggerFactory.getLogger(ChainSummaryWithoutRelationship.class.getName()); - private Map loadedChainSpecificTimeWindowSummary; - private Map updateChainInfo; - - public ChainSummaryWithoutRelationship() { - loadedChainSpecificTimeWindowSummary = new HashMap(); - updateChainInfo = new HashMap(); - } - - public void summary(ChainInfo chainInfo) { - for (ChainNode node : chainInfo.getNodes()) { - String csk = generateChainSummaryKey(chainInfo, node.getStartDate()); - if (!loadedChainSpecificTimeWindowSummary.containsKey(csk)) { - loadedChainSpecificTimeWindowSummary.put(csk, ChainSpecificTimeWindowSummary.load(csk)); - } - loadedChainSpecificTimeWindowSummary.get(csk).summaryNodeValue(node); - } - updateChainInfo.put(chainInfo.getCID(), new Timestamp(System.currentTimeMillis())); - } - - private String generateChainSummaryKey(ChainInfo chainInfo, long startDate) { - return chainInfo.getCID() + "-" + chainInfo.getUserId() + "-" + new SimpleDateFormat("yyyy/MM/dd HH:mm:ss"). - format(new Date(startDate / (1000 * 60 * 60) * (1000 * 60 * 60))); - } - - public void save() throws IOException, InterruptedException, SQLException { - batchSaveChainSpecificTimeWindowSummary(); - updateChainLastActiveTime(); - } - - private void updateChainLastActiveTime() throws SQLException { - DBCallChainInfoDao.updateChainLastActiveTime(updateChainInfo); - } - - private void batchSaveChainSpecificTimeWindowSummary() throws IOException, InterruptedException { - List puts = new ArrayList(); - logger.info("There are [" + loadedChainSpecificTimeWindowSummary.size() + "] summary data will be storage to HBase"); - for (Map.Entry entry : loadedChainSpecificTimeWindowSummary.entrySet()) { - Put put = new Put(entry.getKey().getBytes()); - entry.getValue().save(put); - puts.add(put); - } - - HBaseUtil.batchSaveChainSpecificTimeWindowSummary(puts); - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/UncategorizeChainInfo.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/UncategorizeChainInfo.java deleted file mode 100644 index 1890b8c0b..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/UncategorizeChainInfo.java +++ /dev/null @@ -1,69 +0,0 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.entity; - -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.google.gson.GsonBuilder; -import com.google.gson.annotations.Expose; - -public class UncategorizeChainInfo { - @Expose - private String cid; - @Expose - private String nodeRegEx; - - private ChainInfo chainInfo; - - public UncategorizeChainInfo() { - } - - public UncategorizeChainInfo(ChainInfo chainInfo) { - this.cid = chainInfo.getCID(); - StringBuilder stringBuilder = new StringBuilder(); - boolean flag = false; - for (ChainNode node : chainInfo.getNodes()) { - if (flag) { - stringBuilder.append(";*"); - } - stringBuilder.append((node.getTraceLevelId() + "-" + node.getNodeToken())); - flag = true; - } - - nodeRegEx = stringBuilder.toString(); - - this.chainInfo = chainInfo; - } - - public String getCID() { - return cid; - } - - public String getNodeRegEx() { - return nodeRegEx; - } - - public ChainInfo getChainInfo() { - return chainInfo; - } - - @Override - public boolean equals(Object o) { - if (this == o) return true; - if (!(o instanceof UncategorizeChainInfo)) return false; - - UncategorizeChainInfo that = (UncategorizeChainInfo) o; - - return cid != null ? cid.equals(that.cid) : that.cid == null; - } - - @Override - public int hashCode() { - return cid != null ? cid.hashCode() : 0; - } - - @Override - public String toString() { - GsonBuilder gsonBuilder = new GsonBuilder(); - gsonBuilder.excludeFieldsWithoutExposeAnnotation(); - return gsonBuilder.create().toJson(this); - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/AppendBusinessKeyFilter.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/AppendBusinessKeyFilter.java deleted file mode 100644 index 4becf8e06..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/AppendBusinessKeyFilter.java +++ /dev/null @@ -1,16 +0,0 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl; - -import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry; -import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter; - -public class AppendBusinessKeyFilter extends SpanNodeProcessFilter { - - @Override - public void doFilter(SpanEntry spanEntry, ChainNode node, SubLevelSpanCostCounter costMap) { - node.setViewPoint(node.getViewPoint() + spanEntry.getBusinessKey()); - - this.doNext(spanEntry, node, costMap); - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/CallChainTree.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/CallChainTree.java new file mode 100644 index 000000000..faec57924 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/CallChainTree.java @@ -0,0 +1,88 @@ +package com.ai.cloud.skywalking.analysis.chainbuild; + +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.po.CallChainTreeNode; +import com.ai.cloud.skywalking.analysis.chainbuild.util.HBaseUtil; +import org.apache.hadoop.hbase.client.Put; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +public class CallChainTree { + + private String callEntrance; + + //存放已经合并过的调用链ID + private List hasBeenMergedChainIds; + + // 本次Reduce合并过的调用链 + private Map combineChains; + + //合并之后的节点 + // key : trace level Id + private Map nodes; + + public CallChainTree(String callEntrance) { + hasBeenMergedChainIds = new ArrayList(); + combineChains = new HashMap(); + nodes = new HashMap(); + this.callEntrance = callEntrance; + } + + public static CallChainTree load(String callEntrance) throws IOException { + CallChainTree chain = HBaseUtil.loadMergedCallChain(callEntrance); + chain.hasBeenMergedChainIds.addAll(HBaseUtil.loadHasBeenMergeChainIds(callEntrance)); + if (chain == null) { + chain = new CallChainTree(callEntrance); + } + return chain; + } + + public void processMerge(ChainInfo chainInfo) { + if (hasBeenMergedChainIds.contains(chainInfo.getCID())) { + return; + } + + for (ChainNode node : chainInfo.getNodes()) { + CallChainTreeNode callChainTreeNode = nodes.get(node.getTraceLevelId()); + if (callChainTreeNode != null) { + callChainTreeNode.mergeIfNess(node); + } else { + nodes.put(node.getTraceLevelId(), new CallChainTreeNode(node)); + } + } + + hasBeenMergedChainIds.add(chainInfo.getChainToken()); + combineChains.put(chainInfo.getChainToken(), chainInfo); + } + + public void summary(ChainInfo chainInfo) { + for (ChainNode node : chainInfo.getNodes()) { + CallChainTreeNode callChainTreeNode = nodes.get(node.getTraceLevelId()); + callChainTreeNode.summary(node); + } + } + + public void saveToHbase() { + List chainInfoPuts = new ArrayList(); + for (Map.Entry entry : combineChains.entrySet()) { + Put put = new Put(entry.getKey().getBytes()); + entry.getValue().saveToHBase(put); + chainInfoPuts.add(put); + } + + HBaseUtil.saveMergedCallChain(this); + } + + public String getCallEntrance() { + return callEntrance; + } + + public void addMergedChainNode(CallChainTreeNode chainNode) { + nodes.put(chainNode.getTraceLevelId(), chainNode); + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/ChainBuildMapper.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/ChainBuildMapper.java index 9daf3e75a..9a2c09fea 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/ChainBuildMapper.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/ChainBuildMapper.java @@ -1,6 +1,10 @@ package com.ai.cloud.skywalking.analysis.chainbuild; -import com.ai.cloud.skywalking.analysis.chainbuild.entity.TraceSpanTree; +import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessChain; +import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter; import com.ai.cloud.skywalking.analysis.chainbuild.util.VersionIdentifier; import com.ai.cloud.skywalking.analysis.config.ConfigInitializer; import com.ai.cloud.skywalking.protocol.Span; @@ -14,12 +18,12 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.IOException; -import java.util.ArrayList; -import java.util.List; +import java.util.*; + +public class ChainBuildMapper extends TableMapper { -public class ChainBuildMapper extends TableMapper { private Logger logger = LoggerFactory - .getLogger(ChainBuildMapper.class); + .getLogger(ChainBuildMapper.class.getName()); @Override protected void setup(Context context) throws IOException, @@ -27,6 +31,7 @@ public class ChainBuildMapper extends TableMapper { ConfigInitializer.initialize(); } + @Override protected void map(ImmutableBytesWritable key, Result value, Context context) throws IOException, InterruptedException { @@ -34,21 +39,68 @@ public class ChainBuildMapper extends TableMapper { return; } + List spanList = new ArrayList(); + ChainInfo chainInfo = null; try { - List spanList = new ArrayList(); for (Cell cell : value.rawCells()) { Span span = new Span(Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength())); spanList.add(span); - } - TraceSpanTree tree = new TraceSpanTree(); - tree.build(spanList); - context.write(new Text(tree.getCid()), tree); - } catch (Throwable e) { + chainInfo = spanToChainInfo(Bytes.toString(key.get()), spanList); + logger.info("Success convert span to chain info...." + + chainInfo.getCID() + " TraceId : " + Bytes.toString(key.get())); + context.write( + new Text(chainInfo.getUserId() + ":" + + chainInfo.getEntranceNodeToken()), chainInfo); + } catch (Exception e) { logger.error("Failed to mapper call chain[" + key.toString() + "]", e); } } + + public static ChainInfo spanToChainInfo(String key, List spanList) { + SubLevelSpanCostCounter costMap = new SubLevelSpanCostCounter(); + ChainInfo chainInfo = new ChainInfo(); + Collections.sort(spanList, new Comparator() { + @Override + public int compare(Span span1, Span span2) { + String span1TraceLevel = span1.getParentLevel() + "." + + span1.getLevelId(); + String span2TraceLevel = span2.getParentLevel() + "." + + span2.getLevelId(); + return span1TraceLevel.compareTo(span2TraceLevel); + } + }); + + Map spanEntryMap = mergeSpanDataSet(spanList); + for (Map.Entry entry : spanEntryMap.entrySet()) { + ChainNode chainNode = new ChainNode(); + SpanNodeProcessFilter filter = SpanNodeProcessChain + .getProcessChainByCallType(entry.getValue().getSpanType()); + filter.doFilter(entry.getValue(), chainNode, costMap); + chainInfo.addNodes(chainNode); + } + //chainInfo.generateChainToken(); + //HBaseUtil.saveCidTidMapping(key, chainInfo); + return chainInfo; + } + + private static Map mergeSpanDataSet(List spanList) { + Map spanEntryMap = new LinkedHashMap(); + for (int i = spanList.size() - 1; i >= 0; i--) { + Span span = spanList.get(i); + SpanEntry spanEntry = spanEntryMap.get(span.getParentLevel() + "." + + span.getLevelId()); + if (spanEntry == null) { + spanEntry = new SpanEntry(); + spanEntryMap.put( + span.getParentLevel() + "." + span.getLevelId(), + spanEntry); + } + spanEntry.setSpan(span); + } + return spanEntryMap; + } } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/ChainBuildReducer.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/ChainBuildReducer.java new file mode 100644 index 000000000..f94f1b67e --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/ChainBuildReducer.java @@ -0,0 +1,42 @@ +package com.ai.cloud.skywalking.analysis.chainbuild; + +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo; +import com.ai.cloud.skywalking.analysis.config.ConfigInitializer; +import org.apache.hadoop.hbase.util.Bytes; +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Reducer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.IOException; +import java.util.Iterator; + +public class ChainBuildReducer extends Reducer { + + private Logger logger = LoggerFactory + .getLogger(ChainBuildReducer.class.getName()); + + @Override + protected void setup(Context context) throws IOException, + InterruptedException { + ConfigInitializer.initialize(); + } + + @Override + protected void reduce(Text key, Iterable values, Context context) throws IOException, + InterruptedException { + CallChainTree chainTree = CallChainTree.load(Bytes.toString(key.getBytes())); + Iterator chainInfoIterator = values.iterator(); + while (chainInfoIterator.hasNext()) { + ChainInfo chainInfo = chainInfoIterator.next(); + if (chainInfo.getChainStatus() == ChainInfo.ChainStatus.NORMAL) { + chainTree.processMerge(chainInfo); + } + //合并数据 + chainTree.summary(chainInfo); + } + + chainTree.saveToHbase(); + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/DBCallChainInfoDao.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/DBCallChainInfoDao.java similarity index 82% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/DBCallChainInfoDao.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/DBCallChainInfoDao.java index 9a587c7a3..6214d8747 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/DBCallChainInfoDao.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/DBCallChainInfoDao.java @@ -1,9 +1,8 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain; +package com.ai.cloud.skywalking.analysis.chainbuild; -import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainDetail; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.entity.CallChainDetail; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; import com.ai.cloud.skywalking.analysis.config.Config; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -31,16 +30,16 @@ public class DBCallChainInfoDao { } } - public synchronized static void saveChainDetail(ChainDetail chainDetail) + public synchronized static void saveChainDetail(CallChainDetail callChainDetail) throws SQLException { PreparedStatement preparedStatement = null; try { preparedStatement = connection .prepareStatement("INSERT INTO sw_chain_detail(cid,uid,traceLevelId,viewpoint,create_time)" + " VALUES(?,?,?,?,?)"); - for (ChainNode chainNode : chainDetail.getChainNodes()) { - preparedStatement.setString(1, chainDetail.getChainToken()); - preparedStatement.setString(2, chainDetail.getUserId()); + for (ChainNode chainNode : callChainDetail.getChainNodes()) { + preparedStatement.setString(1, callChainDetail.getChainToken()); + preparedStatement.setString(2, callChainDetail.getUserId()); preparedStatement.setString(3, chainNode.getTraceLevelId()); preparedStatement.setString(4, chainNode.getViewPoint() + ":" + chainNode.getBusinessKey()); @@ -52,7 +51,7 @@ public class DBCallChainInfoDao { for (int i : result) { if (i != 1) { logger.error("Failed to save chain detail [" - + chainDetail.getChainToken() + "]"); + + callChainDetail.getChainToken() + "]"); } } } finally { diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/SpanEntry.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/SpanEntry.java similarity index 97% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/SpanEntry.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/SpanEntry.java index 3df7b4758..7f90ea56b 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/SpanEntry.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/SpanEntry.java @@ -1,6 +1,6 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain; +package com.ai.cloud.skywalking.analysis.chainbuild; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; import com.ai.cloud.skywalking.protocol.CallType; import com.ai.cloud.skywalking.protocol.Span; diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/BranchTraceSpanNode.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/BranchTraceSpanNode.java deleted file mode 100644 index a2d30af19..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/BranchTraceSpanNode.java +++ /dev/null @@ -1,39 +0,0 @@ -package com.ai.cloud.skywalking.analysis.chainbuild.entity; - -import java.util.List; - -public class BranchTraceSpanNode extends TraceSpanNode { - - - protected BranchTraceSpanNode(TraceSpanNode origin, TraceSpanNode dest, - List spanContainer) { - - setNextBranchNode(dest); - dest.parent = this; - - dest.setNextBranchNode(origin); - origin.parent = this; - - this.setParent(dest.parent); - this.branchNode = true; - spanContainer.add(this); - } - - public boolean hasNextBranch() { - return nextBranchNode != null; - } - - public TraceSpanNode nextBranch() { - return nextBranchNode; - } - - public void addBranch(TraceSpanNode branch) { - TraceSpanNode lastBranchNode = null; - while (hasNextBranch()) { - lastBranchNode = nextBranchNode; - } - - lastBranchNode.nextBranchNode = branch; - branch.parent = this; - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainDetail.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainDetail.java similarity index 76% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainDetail.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainDetail.java index 34beb6dfe..ed8ab9e66 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/entity/ChainDetail.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainDetail.java @@ -1,12 +1,10 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.entity; +package com.ai.cloud.skywalking.analysis.chainbuild.entity; -import com.ai.cloud.skywalking.analysis.categorize2chain.DBCallChainInfoDao; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.ai.cloud.skywalking.analysis.config.Config; +import com.ai.cloud.skywalking.analysis.chainbuild.DBCallChainInfoDao; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData; import com.google.gson.Gson; - import org.apache.hadoop.hbase.client.Put; import java.sql.SQLException; @@ -14,13 +12,13 @@ import java.util.Collection; import java.util.HashMap; import java.util.Map; -public class ChainDetail { +public class CallChainDetail { private boolean isNormal = true; private String chainToken; private Map chainNodeMap = new HashMap(); private String userId; - public ChainDetail(ChainInfo chainInfo, boolean isNormal) { + public CallChainDetail(ChainInfo chainInfo, boolean isNormal) { chainToken = chainInfo.getCID(); for (ChainNode chainNode : chainInfo.getNodes()) { chainNodeMap.put(chainNode.getTraceLevelId(), chainNode); diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/ChainNodeForSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/ChainNodeForSummary.java new file mode 100644 index 000000000..157582a6f --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/ChainNodeForSummary.java @@ -0,0 +1,18 @@ +package com.ai.cloud.skywalking.analysis.chainbuild.entity; + +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; + +public class ChainNodeForSummary { + + private String traceLevelId; + private String viewPointId; + + public ChainNodeForSummary(ChainNode node) { + this.traceLevelId = node.getTraceLevelId(); + this.viewPointId = node.getViewPoint(); + } + + public void summary(ChainNode node) { + + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/TraceSpanNode.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/TraceSpanNode.java deleted file mode 100644 index 38588632d..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/TraceSpanNode.java +++ /dev/null @@ -1,323 +0,0 @@ -package com.ai.cloud.skywalking.analysis.chainbuild.entity; - -import java.util.List; - -import com.ai.cloud.skywalking.analysis.chainbuild.exception.TraceSpanTreeNotFountException; -import com.ai.cloud.skywalking.analysis.chainbuild.exception.TraceSpanTreeSerializeException; -import com.ai.cloud.skywalking.analysis.chainbuild.util.StringUtil; -import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator; -import com.ai.cloud.skywalking.protocol.CallType; -import com.ai.cloud.skywalking.protocol.Span; -import com.google.gson.annotations.Expose; - -public class TraceSpanNode { - - protected TraceSpanNode prev = null; - - - protected TraceSpanNode next = null; - - - protected TraceSpanNode parent = null; - - - protected TraceSpanNode sub = null; - @Expose - protected String prevNodeRefToken = null; - @Expose - protected String nextNodeRefToken = null; - @Expose - protected String parentNodeRefToken = null; - @Expose - protected String subNodeRefToken = null; - - @Expose - protected String nodeRefToken = null; - - @Expose - protected boolean visualNode = true; - - @Expose - protected String parentLevel; - - @Expose - protected int levelId; - - @Expose - protected String viewPointId = ""; - - @Expose - protected long cost = 0; - - @Expose - protected long callTimes = 0; - - /** - * 节点调用的状态
- * 0:成功
- * 1:异常
- * 异常判断原则:代码产生exception,并且此exception不在忽略列表中 - */ - @Expose - protected byte statusCode = 0; - - /** - * 节点调用的错误堆栈
- * 堆栈以JAVA的exception为主要判断依据 - */ - @Expose - protected String exceptionStack; - /** - * 节点类型描述
- * 已字符串的形式描述
- * 如:java,dubbo等 - */ - @Expose - protected String spanType = ""; - - /** - * 节点调用过程中的业务字段
- * 如:业务系统设置的订单号,SQL语句等 - */ - @Expose - protected String businessKey = ""; - - /** - * 节点调用所在的系统逻辑名称
- * 由授权文件指定 - */ - @Expose - protected String applicationId = ""; - - /** - * 是否为分支节点 - */ - @Expose - protected boolean branchNode; - - @Expose - protected TraceSpanNode nextBranchNode; - - /** - * Warning: call this constructor ONLY by gson for deserialize - */ - public TraceSpanNode(){ - - } - - public TraceSpanNode(TraceSpanNode parent, TraceSpanNode sub, TraceSpanNode prev, TraceSpanNode next, Span span, List spanContainer) { - this(parent, sub, prev, next, spanContainer); - this.visualNode = false; - this.parentLevel = span.getParentLevel(); - this.levelId = span.getLevelId(); - this.viewPointId = span.getViewPointId(); - this.cost = span.getCost(); - this.callTimes = 1; - this.statusCode = span.getStatusCode(); - if (span.isReceiver()) { - this.exceptionStack = "server stack:"; - } else { - this.exceptionStack = "client stack:"; - } - this.exceptionStack += span.getExceptionStack(); - this.spanType = span.getSpanType(); - this.businessKey = span.getBusinessKey(); - this.applicationId = span.getApplicationId(); - - //nodeToken : MD5(parentLevelId + levelId + viewpoint) - nodeRefToken = TokenGenerator.generateNodeToken(parentLevel + "-" + levelId + "-" + viewPointId); - - } - - protected TraceSpanNode(TraceSpanNode parent, TraceSpanNode sub, TraceSpanNode prev, TraceSpanNode next, List spanContainer) { - this.visualNode = true; - this.setParent(parent); - if (parent != null) { - parent.setSub(this); - } - this.setSub(sub); - if (sub != null) { - sub.setParent(this); - } - this.setPrev(prev); - if (prev != null) { - prev.setNext(this); - } - this.setNext(next); - if (next != null) { - next.setPrev(this); - } - spanContainer.add(this); - } - - protected TraceSpanNode(TraceSpanNode parent, TraceSpanNode sub, TraceSpanNode prev, TraceSpanNode next, String parentLevelId, int levelId, List spanContainer) { - this(parent, sub, prev, next, spanContainer); - this.parentLevel = parentLevelId; - this.levelId = levelId; - this.callTimes = 0; - } - - boolean hasNext() { - if (this.next != null) { - return true; - } else { - return false; - } - } - - boolean hasSub() { - if (this.sub != null) { - return true; - } else { - return false; - } - } - - void mergeSpan(Span span) { - if (CallType.convert(span.getCallType()) == CallType.ASYNC) { - this.cost += span.getCost(); - } - if (span.getStatusCode() != 0 && !StringUtil.isBlank(span.getExceptionStack())) { - if (span.isReceiver()) { - this.exceptionStack += "server stack:"; - } else { - this.exceptionStack += "client stack:"; - } - this.exceptionStack += span.getExceptionStack(); - } - } - - public TraceSpanNode prev(TraceSpanTree tree) throws TraceSpanTreeNotFountException { - if(prev == null){ - if(prevNodeRefToken == null){ - throw new TraceSpanTreeNotFountException(getDesc() + " unexpected prev== null and prevNodeRefToken==null"); - }else{ - prev = tree.findNode(prevNodeRefToken); - } - } - return prev; - } - - public TraceSpanNode next(TraceSpanTree tree) throws TraceSpanTreeNotFountException { - if(next == null){ - if(nextNodeRefToken == null){ - throw new TraceSpanTreeNotFountException(getDesc() + " unexpected next== null and nextNodeRefToken==null"); - }else{ - next = tree.findNode(nextNodeRefToken); - } - } - return next; - } - - public TraceSpanNode parent(TraceSpanTree tree) throws TraceSpanTreeNotFountException { - if(parent == null){ - if(parentNodeRefToken == null){ - throw new TraceSpanTreeNotFountException(getDesc() + " unexpected parent== null and parentNodeRefToken==null"); - }else{ - parent = tree.findNode(parentNodeRefToken); - } - } - return parent; - } - - public TraceSpanNode sub(TraceSpanTree tree) throws TraceSpanTreeNotFountException { - if(sub == null){ - if(subNodeRefToken == null){ - throw new TraceSpanTreeNotFountException(getDesc() + " unexpected sub== null and subNodeRefToken==null"); - }else{ - sub = tree.findNode(subNodeRefToken); - } - } - return sub; - } - - public void setPrev(TraceSpanNode prev) { - this.prev = prev; - } - - public void setNext(TraceSpanNode next) { - this.next = next; - } - - public void setParent(TraceSpanNode parent) { - this.parent = parent; - } - - public void setNextBranchNode(TraceSpanNode nextBranchNode){ - this.nextBranchNode = nextBranchNode; - } - - public void setSub(TraceSpanNode sub) { - this.sub = sub; - } - - public boolean isVisualNode() { - return visualNode; - } - - public String getParentLevel() { - return parentLevel; - } - - public int getLevelId() { - return levelId; - } - - public String getViewPointId() { - return viewPointId; - } - - public long getCost() { - return cost; - } - - public byte getStatusCode() { - return statusCode; - } - - public String getExceptionStack() { - return exceptionStack; - } - - public String getSpanType() { - return spanType; - } - - public String getBusinessKey() { - return businessKey; - } - - public String getApplicationId() { - return applicationId; - } - - public String getNodeRefToken() throws TraceSpanTreeSerializeException { - if (StringUtil.isBlank(nodeRefToken)) { - throw new TraceSpanTreeSerializeException(getDesc() + " ref token is null."); - } - return nodeRefToken; - } - - private String getDesc(){ - return "Node[parentLevel=" + parentLevel + ", levelId=" + levelId + ", viewPointId=" + viewPointId + "]"; - } - - void serializeRef() throws TraceSpanTreeSerializeException { - if (prev != null) { - prevNodeRefToken = prev.getNodeRefToken(); - } - if (parent != null) { - parentNodeRefToken = parent.getNodeRefToken(); - } - if (next != null) { - nextNodeRefToken = next.getNodeRefToken(); - } - if (sub != null) { - subNodeRefToken = sub.getNodeRefToken(); - } - } - - public boolean isBranchNode() { - return branchNode; - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/TraceSpanTree.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/TraceSpanTree.java deleted file mode 100644 index b600508ae..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/TraceSpanTree.java +++ /dev/null @@ -1,310 +0,0 @@ -package com.ai.cloud.skywalking.analysis.chainbuild.entity; - -import com.ai.cloud.skywalking.analysis.chainbuild.exception.BuildTraceSpanTreeException; -import com.ai.cloud.skywalking.analysis.chainbuild.exception.TraceSpanTreeNotFountException; -import com.ai.cloud.skywalking.analysis.chainbuild.exception.TraceSpanTreeSerializeException; -import com.ai.cloud.skywalking.analysis.chainbuild.util.StringUtil; -import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator; -import com.ai.cloud.skywalking.protocol.Span; -import com.google.gson.Gson; -import com.google.gson.GsonBuilder; -import com.google.gson.JsonObject; -import com.google.gson.JsonParser; -import com.google.gson.annotations.Expose; -import com.google.gson.reflect.TypeToken; -import org.apache.hadoop.io.Writable; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.io.DataInput; -import java.io.DataOutput; -import java.io.IOException; -import java.util.*; - -public class TraceSpanTree implements Writable { - private Logger logger = LoggerFactory.getLogger(TraceSpanTree.class); - - @Expose - private String userId = null; - - @Expose - private String cid; - - @Expose - private TraceSpanNode treeRoot; - - @Expose - private List spanContainer = new ArrayList(); - - private Map traceSpanNodeMap = new HashMap(); - - public TraceSpanTree() { - } - - public String build(List spanList) - throws BuildTraceSpanTreeException, TraceSpanTreeNotFountException { - if (spanList.size() == 0) { - throw new BuildTraceSpanTreeException("spanList is empty."); - } - - Collections.sort(spanList, new Comparator() { - @Override - public int compare(Span span1, Span span2) { - String span1TraceLevel = span1.getParentLevel() + "." - + span1.getLevelId(); - String span2TraceLevel = span2.getParentLevel() + "." - + span2.getLevelId(); - return span1TraceLevel.compareTo(span2TraceLevel); - } - }); - Span span = spanList.get(0); - if (!StringUtil.isBlank(span.getUserId())) { - userId = span.getUserId(); - } else { - throw new BuildTraceSpanTreeException( - "spanList[0] 's userId is null"); - } - cid = generateCID(spanList.get(0)); - treeRoot = new TraceSpanNode(null, null, null, null, spanList.get(0), - spanContainer); - if (spanList.size() > 1) { - for (int i = 1; i < spanList.size(); i++) { - this.build(spanList.get(i)); - } - } - - return cid; - } - - private void build(Span span) throws BuildTraceSpanTreeException, - TraceSpanTreeNotFountException { - if (userId == null && !StringUtil.isBlank(span.getUserId())) { - userId = span.getUserId(); - } - - TraceSpanNode clientOrServerNode = findNodeAndCreateVisualNodeIfNess( - span.getParentLevel(), span.getLevelId()); - if (clientOrServerNode != null) { - clientOrServerNode.mergeSpan(span); - } - - if (span.getLevelId() > 0) { - TraceSpanNode foundNode = findNodeAndCreateVisualNodeIfNess( - span.getParentLevel(), span.getLevelId() - 1); - /** - * Create node between foundNode and foundNode.next(maybe - * foundNode.next == null) - */ - new TraceSpanNode(null, null, foundNode, foundNode.next(this), - span, spanContainer); - } else { - /** - * levelId=0 find for parent level if parentLevelId = 0.0.1 then - * find node[parentLevelId=0.0,levelId=1] - */ - String parentLevel = span.getParentLevel(); - int idx = parentLevel.lastIndexOf("\\."); - if (idx < 0) { - throw new BuildTraceSpanTreeException("parentLevel=" - + parentLevel + " is unexpected."); - } - TraceSpanNode foundNode = findNodeAndCreateVisualNodeIfNess( - parentLevel.substring(0, idx), - Integer.parseInt(parentLevel.substring(idx + 1))); - /** - * Create sub node of using span data. FoundNode is parent node. - */ - new TraceSpanNode(foundNode, null, null, null, span, spanContainer); - - } - } - - private TraceSpanNode findNodeAndCreateVisualNodeIfNess( - String parentLevelId, int levelId) - throws TraceSpanTreeNotFountException { - String levelDesc = StringUtil.isBlank(parentLevelId) ? (levelId + "") - : (parentLevelId + "." + levelId); - String[] levelArray = levelDesc.split("\\."); - - TraceSpanNode currentNode = treeRoot; - String contextParentLevelId = ""; - for (String currentLevel : levelArray) { - int currentLevelInt = Integer.parseInt(currentLevel); - for (int i = 0; i < currentLevelInt; i++) { - if (currentNode.hasNext()) { - currentNode = currentNode.next(this); - } else { - // create visual next node - currentNode = new VisualTraceSpanNode(null, null, - currentNode, null, contextParentLevelId, i, - spanContainer); - } - } - contextParentLevelId = contextParentLevelId == "" ? ("" + currentLevelInt) - : (contextParentLevelId + "." + currentLevelInt); - if (currentNode.hasSub()) { - currentNode = currentNode.sub(this); - } else { - // create visual sub node - currentNode = new VisualTraceSpanNode(currentNode, null, null, - null, contextParentLevelId, 0, spanContainer); - } - } - - return currentNode; - } - - private String generateCID(Span level0Span) - throws BuildTraceSpanTreeException { - if (StringUtil.isBlank(level0Span.getParentLevel()) - && level0Span.getLevelId() == 0) { - StringBuilder chainTokenDesc = new StringBuilder(); - chainTokenDesc.append(userId).append("_"); - chainTokenDesc.append(level0Span.getViewPointId()); - return getTSBySpanTraceId(level0Span) + "_" + TokenGenerator.generateCID(chainTokenDesc.toString()); - } else { - throw new BuildTraceSpanTreeException("tid:" - + level0Span.getTraceId() + " level0 span data is illegal"); - } - } - - private static String getTSBySpanTraceId(Span span) - throws BuildTraceSpanTreeException { - try { - Calendar calendar = Calendar.getInstance(); - calendar.setTime(new Date(Long.parseLong(span.getTraceId().split( - "\\.")[2]))); - return calendar.get(Calendar.YEAR) + "-" + (calendar.get(Calendar.MONTH) + 1); - } catch (Throwable t) { - throw new BuildTraceSpanTreeException("tid:" + span.getTraceId() - + " is illegal."); - } - } - - private void beforeSerialize() throws TraceSpanTreeSerializeException { - for (TraceSpanNode treeNode : spanContainer) { - treeNode.serializeRef(); - } - } - - public String serialize() throws TraceSpanTreeSerializeException { - beforeSerialize(); - return new GsonBuilder().excludeFieldsWithoutExposeAnnotation() - .create().toJson(this); - } - - TraceSpanNode findNode(String nodeRefToken) - throws TraceSpanTreeNotFountException { - if (traceSpanNodeMap.containsKey(nodeRefToken)) { - return traceSpanNodeMap.get(nodeRefToken); - } else { - throw new TraceSpanTreeNotFountException("nodeRefToken=" - + nodeRefToken + " not found."); - } - } - - @Override - public void write(DataOutput out) throws IOException { - try { - out.write(serialize().getBytes()); - } catch (TraceSpanTreeSerializeException e) { - logger.error("Failed to serialize Chain Id[" + cid + "]", e); - } - } - - @Override - public void readFields(DataInput in) throws IOException { - String value = in.readLine(); - try { - JsonObject jsonObject = (JsonObject) new JsonParser().parse(value); - userId = jsonObject.get("userId").getAsString(); - cid = jsonObject.get("cid").getAsString(); - treeRoot = new Gson().fromJson(jsonObject.get("treeRoot"), - TraceSpanNode.class); - spanContainer = new Gson().fromJson( - jsonObject.get("spanContainer"), - new TypeToken>() { - }.getType()); - for (TraceSpanNode node : spanContainer) { - traceSpanNodeMap.put(node.getNodeRefToken(), node); - } - } catch (Exception e) { - logger.error("Failed to parse the value[" + value - + "] to TraceSpanTree Object", e); - } - } - - public TraceSpanNode getTreeRoot() { - return treeRoot; - } - - public String getCid() { - return cid; - } - - public void merge(TraceSpanTree spanTree) { - if (spanTree.getTreeRoot().hasNext()) { - SpanTreeMerger.merge(spanTree.getTreeRoot().next, treeRoot.next, spanContainer); - } - - if (spanTree.getTreeRoot().hasSub()) { - SpanTreeMerger.merge(spanTree.getTreeRoot().sub, treeRoot.sub, spanContainer); - } - } - - private static class SpanTreeMerger { - public static boolean merge(TraceSpanNode origin, TraceSpanNode dest, List spanContainer) { - boolean flag = false; - if (origin == null || dest == null) { - if (origin != null && dest == null) { - dest.parent.sub = origin; - origin.parent = dest.parent; - return true; - } - return true; - } - - if (dest.isBranchNode()) { - BranchTraceSpanNode branchTraceSpanNode = (BranchTraceSpanNode) dest; - boolean branchFlag = false; - while (branchTraceSpanNode.hasNextBranch()) { - branchFlag = merge(origin, branchTraceSpanNode.nextBranch(), spanContainer); - if (branchFlag) { - break; - } - } - - if (branchFlag) { - return true; - } else { - branchTraceSpanNode.addBranch(origin); - return false; - } - } - - if (origin.isVisualNode() || dest.isVisualNode()) { - boolean nextFlag = merge(origin.next, dest.next, spanContainer); - boolean subFlag = merge(origin.sub, dest.sub, spanContainer); - - if (subFlag && nextFlag) { - // 合并子树数据 - } else { - new BranchTraceSpanNode(origin, dest, spanContainer); - } - flag = nextFlag && nextFlag; - } else { - if (origin.nodeRefToken.equals(dest.nodeRefToken)) { - // 合并子树数据 - flag = true; - } else { - new BranchTraceSpanNode(origin, dest, spanContainer); - flag = false; - } - flag = flag && merge(origin.next, dest.next, spanContainer); - flag = flag && merge(origin.sub, dest.sub, spanContainer); - } - - return flag; - } - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/VisualTraceSpanNode.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/VisualTraceSpanNode.java deleted file mode 100644 index 4b589b74b..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/VisualTraceSpanNode.java +++ /dev/null @@ -1,24 +0,0 @@ -package com.ai.cloud.skywalking.analysis.chainbuild.entity; - -import java.util.List; - -import com.ai.cloud.skywalking.analysis.chainbuild.util.StringUtil; - -public class VisualTraceSpanNode extends TraceSpanNode { - - protected VisualTraceSpanNode(TraceSpanNode parent, TraceSpanNode sub, - TraceSpanNode prev, TraceSpanNode next, String parentLevelId, - int levelId, List spanContainer) { - super(parent, sub, prev, next, parentLevelId, levelId, spanContainer); - - /**set visual node token.
- * for example:
- * VisualNode[0.0]
- * VisualNode[0.0.1]
- * etc.
- */ - nodeRefToken = "VisualNode[" + (StringUtil.isBlank(parentLevelId) ? "": nodeRefToken + ".") + levelId + "]"; - } - - -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/BuildTraceSpanTreeException.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/BuildTraceSpanTreeException.java deleted file mode 100644 index ba0d2fcf8..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/BuildTraceSpanTreeException.java +++ /dev/null @@ -1,13 +0,0 @@ -package com.ai.cloud.skywalking.analysis.chainbuild.exception; - -public class BuildTraceSpanTreeException extends Exception { - private static final long serialVersionUID = 5816399370389190974L; - - public BuildTraceSpanTreeException(String msg){ - super(msg); - } - - public BuildTraceSpanTreeException(String msg, Exception cause){ - super(msg, cause); - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/TraceSpanTreeNotFountException.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/TraceSpanTreeNotFountException.java deleted file mode 100644 index 9bff13a02..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/TraceSpanTreeNotFountException.java +++ /dev/null @@ -1,13 +0,0 @@ -package com.ai.cloud.skywalking.analysis.chainbuild.exception; - -public class TraceSpanTreeNotFountException extends Exception { - private static final long serialVersionUID = 5559441397011866237L; - - public TraceSpanTreeNotFountException(String msg){ - super(msg); - } - - public TraceSpanTreeNotFountException(String msg, Exception cause){ - super(msg, cause); - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/TraceSpanTreeSerializeException.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/TraceSpanTreeSerializeException.java deleted file mode 100644 index 5f0bdd9e0..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/exception/TraceSpanTreeSerializeException.java +++ /dev/null @@ -1,13 +0,0 @@ -package com.ai.cloud.skywalking.analysis.chainbuild.exception; - -public class TraceSpanTreeSerializeException extends Exception { - private static final long serialVersionUID = 7857716041262993579L; - - public TraceSpanTreeSerializeException(String msg){ - super(msg); - } - - public TraceSpanTreeSerializeException(String msg, Exception cause){ - super(msg, cause); - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/SpanNodeProcessChain.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/SpanNodeProcessChain.java similarity index 98% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/SpanNodeProcessChain.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/SpanNodeProcessChain.java index d3bb6a3d3..c5fa17c74 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/SpanNodeProcessChain.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/SpanNodeProcessChain.java @@ -1,4 +1,4 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.filter; +package com.ai.cloud.skywalking.analysis.chainbuild.filter; import java.io.IOException; import java.util.HashMap; diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/SpanNodeProcessFilter.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/SpanNodeProcessFilter.java similarity index 64% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/SpanNodeProcessFilter.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/SpanNodeProcessFilter.java index eb8b54409..b26c1a022 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/SpanNodeProcessFilter.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/SpanNodeProcessFilter.java @@ -1,8 +1,8 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.filter; +package com.ai.cloud.skywalking.analysis.chainbuild.filter; -import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter; +import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter; public abstract class SpanNodeProcessFilter { diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/AppendBusinessKeyFilter.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/AppendBusinessKeyFilter.java new file mode 100644 index 000000000..d6632af93 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/AppendBusinessKeyFilter.java @@ -0,0 +1,16 @@ +package com.ai.cloud.skywalking.analysis.chainbuild.filter.impl; + +import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry; +import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter; + +public class AppendBusinessKeyFilter extends SpanNodeProcessFilter { + + @Override + public void doFilter(SpanEntry spanEntry, ChainNode node, SubLevelSpanCostCounter costMap) { + node.setViewPoint(node.getViewPoint() + spanEntry.getBusinessKey()); + + this.doNext(spanEntry, node, costMap); + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/CopyAttrFilter.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/CopyAttrFilter.java similarity index 63% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/CopyAttrFilter.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/CopyAttrFilter.java index 92a615216..ac3669312 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/CopyAttrFilter.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/CopyAttrFilter.java @@ -1,9 +1,9 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl; +package com.ai.cloud.skywalking.analysis.chainbuild.filter.impl; -import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry; -import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter; +import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry; +import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter; public class CopyAttrFilter extends SpanNodeProcessFilter { diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/ProcessCostTimeFilter.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/ProcessCostTimeFilter.java similarity index 72% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/ProcessCostTimeFilter.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/ProcessCostTimeFilter.java index 95dd78ce0..7704f087b 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/ProcessCostTimeFilter.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/ProcessCostTimeFilter.java @@ -1,9 +1,9 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl; +package com.ai.cloud.skywalking.analysis.chainbuild.filter.impl; -import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry; -import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter; +import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry; +import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter; public class ProcessCostTimeFilter extends SpanNodeProcessFilter { @Override diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/ReplaceAddressFilter.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/ReplaceAddressFilter.java similarity index 66% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/ReplaceAddressFilter.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/ReplaceAddressFilter.java index 9c12fafc7..d05d6b94c 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/ReplaceAddressFilter.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/ReplaceAddressFilter.java @@ -1,9 +1,9 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl; +package com.ai.cloud.skywalking.analysis.chainbuild.filter.impl; -import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry; -import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter; +import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry; +import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter; public class ReplaceAddressFilter extends SpanNodeProcessFilter { diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/TokenGenerateFilter.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/TokenGenerateFilter.java similarity index 56% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/TokenGenerateFilter.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/TokenGenerateFilter.java index 990de051d..803a0ee30 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/filter/impl/TokenGenerateFilter.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/filter/impl/TokenGenerateFilter.java @@ -1,9 +1,9 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl; +package com.ai.cloud.skywalking.analysis.chainbuild.filter.impl; -import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry; -import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode; -import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter; +import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry; +import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter; import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator; public class TokenGenerateFilter extends SpanNodeProcessFilter { diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/CallChainTreeNode.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/CallChainTreeNode.java new file mode 100644 index 000000000..385d1f2e8 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/CallChainTreeNode.java @@ -0,0 +1,46 @@ +package com.ai.cloud.skywalking.analysis.chainbuild.po; + +import com.ai.cloud.skywalking.analysis.chainbuild.entity.ChainNodeForSummary; +import com.google.gson.Gson; +import com.google.gson.JsonObject; +import com.google.gson.JsonParser; +import com.google.gson.reflect.TypeToken; + +import java.util.Map; + +public class CallChainTreeNode { + + private String traceLevelId; + + // key: nodeToken + private Map chainNodeContainer; + + + public CallChainTreeNode(ChainNode node) { + this.traceLevelId = node.getTraceLevelId(); + chainNodeContainer.put(node.getNodeToken(), new ChainNodeForSummary(node)); + } + + public CallChainTreeNode(String originData) { + JsonObject jsonObject = (JsonObject) new JsonParser().parse(originData); + traceLevelId = jsonObject.get("traceLevelId").getAsString(); + chainNodeContainer = new Gson().fromJson(jsonObject.get("chainNodeContainer").getAsString(), + new TypeToken>() { + }.getType()); + } + + public void mergeIfNess(ChainNode node) { + if (!chainNodeContainer.containsKey(node.getNodeToken())) { + chainNodeContainer.put(node.getNodeToken(), new ChainNodeForSummary(node)); + } + } + + public void summary(ChainNode node) { + ChainNodeForSummary chainNode = chainNodeContainer.get(node.getNodeToken()); + chainNode.summary(node); + } + + public String getTraceLevelId() { + return traceLevelId; + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/po/ChainInfo.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/ChainInfo.java similarity index 88% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/po/ChainInfo.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/ChainInfo.java index 0a0a14294..621977e08 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/po/ChainInfo.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/ChainInfo.java @@ -1,11 +1,11 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.po; +package com.ai.cloud.skywalking.analysis.chainbuild.po; import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator; import com.google.gson.Gson; import com.google.gson.JsonObject; import com.google.gson.JsonParser; import com.google.gson.reflect.TypeToken; - +import org.apache.hadoop.hbase.client.Put; import org.apache.hadoop.io.Writable; import java.io.DataInput; @@ -21,6 +21,7 @@ public class ChainInfo implements Writable { private String userId = null; private ChainNode firstChainNode; private long startDate; + private String chainToken; public ChainInfo(String userId) { super(); @@ -61,15 +62,6 @@ public class ChainInfo implements Writable { } } - public void generateChainToken() { - StringBuilder chainTokenDesc = new StringBuilder(); - for (ChainNode node : nodes) { - chainTokenDesc.append(node.getParentLevelId() + "." - + node.getLevelId() + "-" + node.getNodeToken() + ";"); - } - this.cid = TokenGenerator.generateCID(chainTokenDesc.toString()); - } - public ChainStatus getChainStatus() { return chainStatus; } @@ -92,6 +84,7 @@ public class ChainInfo implements Writable { && chainNode.getLevelId() == 0) { firstChainNode = chainNode; startDate = chainNode.getStartDate(); + cid = firstChainNode.getViewPoint(); } } @@ -107,6 +100,14 @@ public class ChainInfo implements Writable { this.userId = userId; } + public String getChainToken() { + return chainToken; + } + + public void saveToHBase(Put put) { + + } + public enum ChainStatus { NORMAL('N'), ABNORMAL('A'); private char value; diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/po/ChainNode.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/ChainNode.java similarity index 98% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/po/ChainNode.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/ChainNode.java index 5349bd986..f07ecbf4a 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/po/ChainNode.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/ChainNode.java @@ -1,4 +1,4 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.po; +package com.ai.cloud.skywalking.analysis.chainbuild.po; import com.google.gson.GsonBuilder; import com.google.gson.annotations.Expose; diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/util/HBaseUtil.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/HBaseUtil.java similarity index 66% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/util/HBaseUtil.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/HBaseUtil.java index 2156fb1ba..67d853faa 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/util/HBaseUtil.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/HBaseUtil.java @@ -1,19 +1,13 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.util; +package com.ai.cloud.skywalking.analysis.chainbuild.util; -import com.ai.cloud.skywalking.analysis.categorize2chain.*; -import com.ai.cloud.skywalking.analysis.categorize2chain.entity.CategorizedChainInfo; -import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainNodeSpecificTimeWindowSummary; -import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainRelationship; -import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainSpecificTimeWindowSummary; -import com.ai.cloud.skywalking.analysis.categorize2chain.entity.UncategorizeChainInfo; -import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo; -import com.ai.cloud.skywalking.analysis.chain2summary.ChainRelationship4Search; import com.ai.cloud.skywalking.analysis.chain2summary.entity.*; +import com.ai.cloud.skywalking.analysis.chainbuild.CallChainTree; +import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo; +import com.ai.cloud.skywalking.analysis.chainbuild.po.CallChainTreeNode; import com.ai.cloud.skywalking.analysis.config.Config; import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData; import com.google.gson.Gson; import com.google.gson.reflect.TypeToken; - import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.hbase.*; import org.apache.hadoop.hbase.client.*; @@ -36,30 +30,35 @@ public class HBaseUtil { try { initHBaseClient(); - createTableIfNeed(HBaseTableMetaData.TABLE_CID_TID_MAPPING.TABLE_NAME, - HBaseTableMetaData.TABLE_CID_TID_MAPPING.COLUMN_FAMILY_NAME); +// createTableIfNeed(HBaseTableMetaData.TABLE_CID_TID_MAPPING.TABLE_NAME, +// HBaseTableMetaData.TABLE_CID_TID_MAPPING.COLUMN_FAMILY_NAME); +// +// createTableIfNeed(HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.TABLE_NAME, +// HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.COLUMN_FAMILY_NAME); +// +// createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP.TABLE_NAME, +// HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME); +// +// createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_DETAIL.TABLE_NAME, +// HBaseTableMetaData.TABLE_CHAIN_DETAIL.COLUMN_FAMILY_NAME); +// +// createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_ONE_HOUR_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME, +// HBaseTableMetaData.TABLE_CHAIN_ONE_HOUR_SUMMARY_INCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME); +// +// createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_ONE_DAY_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME, +// HBaseTableMetaData.TABLE_CHAIN_ONE_DAY_SUMMARY_INCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME); +// +// createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME, +// HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_INCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME); +// +// createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_ONE_MONTH_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME, +// HBaseTableMetaData.TABLE_CHAIN_ONE_MONTH_SUMMARY_INCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME); - createTableIfNeed(HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.TABLE_NAME, - HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.COLUMN_FAMILY_NAME); - - createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP.TABLE_NAME, - HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME); - - createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_DETAIL.TABLE_NAME, - HBaseTableMetaData.TABLE_CHAIN_DETAIL.COLUMN_FAMILY_NAME); - - createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_ONE_HOUR_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME, - HBaseTableMetaData.TABLE_CHAIN_ONE_HOUR_SUMMARY_INCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME); - - createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_ONE_DAY_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME, - HBaseTableMetaData.TABLE_CHAIN_ONE_DAY_SUMMARY_INCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME); - - createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME, - HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_INCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME); - - createTableIfNeed(HBaseTableMetaData.TABLE_CHAIN_ONE_MONTH_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME, - HBaseTableMetaData.TABLE_CHAIN_ONE_MONTH_SUMMARY_INCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME); + createTableIfNeed(HBaseTableMetaData.TABLE_MERGED_CHAIN_DETAIL.TABLE_NAME, + HBaseTableMetaData.TABLE_MERGED_CHAIN_DETAIL.COLUMN_FAMILY_NAME); + createTableIfNeed(HBaseTableMetaData.TABLE_CALL_CHAIN_TREE_ID_AND_CID_MAPPING.TABLE_NAME, + HBaseTableMetaData.TABLE_CALL_CHAIN_TREE_ID_AND_CID_MAPPING.COLUMN_FAMILY_NAME); } catch (IOException e) { logger.error("Create tables failed", e); } @@ -115,52 +114,6 @@ public class HBaseUtil { return true; } - - public static ChainRelationship loadCallChainRelationship(String key) throws IOException { - ChainRelationship chainRelate = new ChainRelationship(key); - Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.TABLE_NAME)); - Get g = new Get(Bytes.toBytes(key)); - Result r = table.get(g); - for (Cell cell : r.rawCells()) { - if (cell.getValueArray().length > 0) { - - String qualifierName = Bytes.toString(cell.getQualifierArray(), cell.getQualifierOffset(), - cell.getQualifierLength()); - if (HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.UNCATEGORIZE_COLUMN_NAME.equals(qualifierName)) { - List uncategorizeChainInfoList = new Gson().fromJson(Bytes.toString(cell.getValueArray(), - cell.getValueOffset(), cell.getValueLength()), - new TypeToken>() { - }.getType()); - chainRelate.addUncategorizeChain(uncategorizeChainInfoList); - } else { - chainRelate.addCategorizeChain(qualifierName, new CategorizedChainInfo( - Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength()) - )); - } - } - } - return chainRelate; - } - - public static ChainSpecificTimeWindowSummary selectChainSummaryResult(String key) throws IOException { - ChainSpecificTimeWindowSummary result = null; - Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP.TABLE_NAME)); - Get g = new Get(Bytes.toBytes(key)); - Result r = table.get(g); - - if (r.rawCells().length == 0) { - return null; - } - result = new ChainSpecificTimeWindowSummary(); - for (Cell cell : r.rawCells()) { - if (cell.getValueArray().length > 0) - result.addNodeSummaryResult(new ChainNodeSpecificTimeWindowSummary(Bytes.toString(cell.getValueArray(), - cell.getValueOffset(), cell.getValueLength()))); - } - - return result; - } - public static void saveChainRelationship(Put put) throws IOException { Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.TABLE_NAME)); @@ -194,44 +147,6 @@ public class HBaseUtil { } } - public static ChainRelationship4Search queryChainRelationship(String key) throws IOException { - Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP.TABLE_NAME)); - Get g = new Get(key.getBytes()); - Result r = table.get(g); - - if (r.rawCells().length == 0) { - return null; - } - ChainRelationship4Search result = new ChainRelationship4Search(); - - for (Cell cell : r.rawCells()) { - if (cell.getValueArray().length > 0) { - String qualifierName = Bytes.toString(cell.getQualifierArray(), cell.getQualifierOffset(), - cell.getQualifierLength()); - - if (HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.UNCATEGORIZE_COLUMN_NAME.equals(qualifierName)) { - List uncategorizeChainInfoList = new Gson().fromJson(Bytes.toString(cell.getValueArray(), - cell.getValueOffset(), cell.getValueLength()), - new TypeToken>() { - }.getType()); - for (UncategorizeChainInfo chainInfo : uncategorizeChainInfoList) { - result.addRelationship(chainInfo.getCID()); - } - } else { - CategorizedChainInfo categorizedChainInfo = new CategorizedChainInfo( - Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength()) - ); - - for (String cid : categorizedChainInfo.getChildren_Token()) { - result.addRelationship(qualifierName, cid); - } - } - } - } - - return result; - } - public static ChainSpecificMinSummary loadSpecificMinSummary(String key) throws IOException { ChainSpecificMinSummary result = null; Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME)); @@ -358,4 +273,46 @@ public class HBaseUtil { } } } + + public static CallChainTree loadMergedCallChain(String callEntrance) throws IOException { + CallChainTree result = null; + Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_MERGED_CHAIN_DETAIL.TABLE_NAME)); + Get g = new Get(Bytes.toBytes(callEntrance)); + Result r = table.get(g); + if (r.rawCells().length == 0) { + return null; + } + result = new CallChainTree(callEntrance); + for (Cell cell : r.rawCells()) { + if (cell.getValueArray().length > 0) + result.addMergedChainNode(new CallChainTreeNode(Bytes.toString(cell.getValueArray(), + cell.getValueOffset(), cell.getValueLength()))); + } + return result; + } + + public static void saveMergedCallChain(CallChainTree callChainTree) { + // save + // save relationship + } + + public static List loadHasBeenMergeChainIds(String topoId) throws IOException { + List result = new ArrayList(); + Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CALL_CHAIN_TREE_ID_AND_CID_MAPPING.TABLE_NAME)); + Get g = new Get(Bytes.toBytes(topoId)); + Result r = table.get(g); + if (r.rawCells().length == 0) { + return null; + } + for (Cell cell : r.rawCells()) { + if (cell.getValueArray().length > 0) { + List hasBeenMergedCIds = new Gson().fromJson("", + new TypeToken>() { + }.getType()); + result.addAll(hasBeenMergedCIds); + } + + } + return result; + } } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/StringUtil.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/StringUtil.java deleted file mode 100644 index c24615f29..000000000 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/StringUtil.java +++ /dev/null @@ -1,22 +0,0 @@ -package com.ai.cloud.skywalking.analysis.chainbuild.util; - -public class StringUtil { - public static boolean isBlank(String str){ - if(str == null || str == "" || str.trim() == ""){ - return true; - }else{ - return false; - } - } - - public static boolean equal(String str1, String str2){ - if(str1 == null){ - str1 = ""; - } - if(str2 == null){ - str2 = ""; - } - - return str1.trim().equals(str2.trim()); - } -} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/util/SubLevelSpanCostCounter.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/SubLevelSpanCostCounter.java similarity index 87% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/util/SubLevelSpanCostCounter.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/SubLevelSpanCostCounter.java index 90f08ae5a..2c136080d 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/util/SubLevelSpanCostCounter.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/SubLevelSpanCostCounter.java @@ -1,4 +1,4 @@ -package com.ai.cloud.skywalking.analysis.categorize2chain.util; +package com.ai.cloud.skywalking.analysis.chainbuild.util; import java.util.HashMap; import java.util.Map; diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/TokenGenerator.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/TokenGenerator.java index bde561d82..4f4d29d81 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/TokenGenerator.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/TokenGenerator.java @@ -12,13 +12,13 @@ public class TokenGenerator { private TokenGenerator() { //Non } - + public static String generateCID(String originData) { - return "CID_" + generate(originData); + return "CID_" + generate(originData); } - + public static String generateNodeToken(String originData){ - return "C_NID_" + generate(originData); + return "C_NID_" + generate(originData); } private static String generate(String originData) { @@ -41,4 +41,4 @@ public class TokenGenerator { } return result.toString().toUpperCase(); } -} +} \ No newline at end of file diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/VersionIdentifier.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/VersionIdentifier.java index 791da5269..c21452320 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/VersionIdentifier.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/util/VersionIdentifier.java @@ -1,31 +1,25 @@ package com.ai.cloud.skywalking.analysis.chainbuild.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; - } -} + /** + * 根据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; + } +} \ No newline at end of file diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/HBaseTableMetaData.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/HBaseTableMetaData.java index a43cfc634..ca5e88e53 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/HBaseTableMetaData.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/HBaseTableMetaData.java @@ -43,8 +43,6 @@ public class HBaseTableMetaData { public static final String TABLE_NAME = "sw-chain-relationship"; public static final String COLUMN_FAMILY_NAME = "chain-relationship"; - - public static final String UNCATEGORIZE_COLUMN_NAME = "UNCATEGORIZED_CALL_CHAIN"; } /** @@ -101,4 +99,22 @@ public class HBaseTableMetaData { public static final String COLUMN_FAMILY_NAME = "chain_summary"; } + + /** + * 用于存放已经合并的调用链的信息 + * + * @author zhangxin + */ + public final static class TABLE_MERGED_CHAIN_DETAIL { + public static final String TABLE_NAME = "sw-merged-chain-detail"; + + public static final String COLUMN_FAMILY_NAME = "chain_detail"; + } + + + public final static class TABLE_CALL_CHAIN_TREE_ID_AND_CID_MAPPING { + public static final String TABLE_NAME = "sw-topologyId-cid-mapping"; + + public static final String COLUMN_FAMILY_NAME = "sw-topologyId-cid-mapping"; + } }