From 89dbf1e5d6fdea4b8dbd30a38c7b01bb19f7104b Mon Sep 17 00:00:00 2001 From: ascrutae Date: Sun, 1 May 2016 11:31:22 +0800 Subject: [PATCH] =?UTF-8?q?=E5=B9=B4=E6=9C=88=E6=97=A5=E6=97=B6=E6=8A=A5?= =?UTF-8?q?=E8=A1=A8=E5=88=86=E6=89=B9=E6=AC=A1=E7=BB=9F=E8=AE=A1=EF=BC=8C?= =?UTF-8?q?=E9=80=9A=E8=BF=87=E4=BF=AE=E6=94=B9MapReduce=E7=9A=84Key?= =?UTF-8?q?=E7=9A=84=E4=B8=8D=E5=90=8C=E7=BB=9F=E8=AE=A1=E4=B8=8D=E5=90=8C?= =?UTF-8?q?=E7=9A=84=E5=80=BC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../analysis/chainbuild/ChainBuildMapper.java | 9 +++++- .../chainbuild/ChainBuildReducer.java | 15 ++++++--- .../chainbuild/entity/CallChainTree.java | 5 +-- .../chainbuild/entity/CallChainTreeNode.java | 20 ++++++++---- ...ficTimeCallTreeMergedChainIdContainer.java | 15 +++++---- .../analysis/chainbuild/po/SummaryType.java | 31 +++++++++++++++++++ 6 files changed, 76 insertions(+), 19 deletions(-) create mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/SummaryType.java 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 8ac5a17d5..5cdaa2941 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 @@ -5,6 +5,7 @@ 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.po.SummaryType; 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; @@ -57,7 +58,13 @@ public class ChainBuildMapper extends TableMapper { + "] to chain with cid[" + chainInfo.getCID() + "]."); if (chainInfo.getCallEntrance() != null && chainInfo.getCallEntrance().length() > 0) { context.write( - new Text(chainInfo.getCallEntrance()), new Text(new Gson().toJson(chainInfo))); + new Text(chainInfo.getCallEntrance() + "@#!" + SummaryType.MINUTER.getValue()), new Text(new Gson().toJson(chainInfo))); + context.write( + new Text(chainInfo.getCallEntrance() + "@#!" + SummaryType.HOUR.getValue()), new Text(new Gson().toJson(chainInfo))); + context.write( + new Text(chainInfo.getCallEntrance() + "@#!" + SummaryType.DAY.getValue()), new Text(new Gson().toJson(chainInfo))); + context.write( + new Text(chainInfo.getCallEntrance() + "@#!" + SummaryType.MONTH.getValue()), new Text(new Gson().toJson(chainInfo))); } } catch (Exception e) { logger.error("Failed to mapper call chain[" + key.toString() + "]", 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 index 0deba4c67..a8d146c6d 100644 --- 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 @@ -3,6 +3,7 @@ package com.ai.cloud.skywalking.analysis.chainbuild; import com.ai.cloud.skywalking.analysis.chainbuild.entity.CallChainTree; import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo; import com.ai.cloud.skywalking.analysis.chainbuild.po.SpecificTimeCallTreeMergedChainIdContainer; +import com.ai.cloud.skywalking.analysis.chainbuild.po.SummaryType; import com.ai.cloud.skywalking.analysis.config.Config; import com.ai.cloud.skywalking.analysis.config.ConfigInitializer; import com.google.gson.Gson; @@ -32,10 +33,16 @@ public class ChainBuildReducer extends Reducer { @Override protected void reduce(Text key, Iterable values, Context context) throws IOException, InterruptedException { - doReduceAction(Bytes.toString(key.getBytes()), values.iterator()); + String reduceKey = Bytes.toString(key.getBytes()); + int index = reduceKey.indexOf("@#!"); + if (index == -1){ + return; + } + + doReduceAction(reduceKey.substring(0, index - 1), SummaryType.convert(reduceKey.substring(index + 3)), values.iterator()); } - public void doReduceAction(String key, Iterator chainInfoIterator) + public void doReduceAction(String key, SummaryType summaryType, Iterator chainInfoIterator) throws IOException, InterruptedException { CallChainTree chainTree = CallChainTree.load(key); SpecificTimeCallTreeMergedChainIdContainer container @@ -46,7 +53,7 @@ public class ChainBuildReducer extends Reducer { try { chainInfo = new Gson().fromJson(callChainData, ChainInfo.class); container.addMergedChainIfNotContain(chainInfo); - chainTree.summary(chainInfo); + chainTree.summary(chainInfo, summaryType); } catch (Exception e) { logger.error( "Failed to summary call chain, maybe illegal data:" @@ -54,7 +61,7 @@ public class ChainBuildReducer extends Reducer { } } try { - container.saveToHBase(); + container.saveToHBase(summaryType); chainTree.saveToHbase(); } catch (Exception e) { logger.error("Failed to save summaryresult/chainTree.", e); diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainTree.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainTree.java index cfb4d938d..317a63c42 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainTree.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainTree.java @@ -6,6 +6,7 @@ import java.util.Map; 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.SummaryType; import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -35,7 +36,7 @@ public class CallChainTree { return chain; } - public void summary(ChainInfo chainInfo) throws IOException { + public void summary(ChainInfo chainInfo, SummaryType summaryType) throws IOException { for (ChainNode node : chainInfo.getNodes()) { CallChainTreeNode newCallChainTreeNode = new CallChainTreeNode(node); CallChainTreeNode callChainTreeNode = nodes.get(newCallChainTreeNode.getTreeNodeId()); @@ -43,7 +44,7 @@ public class CallChainTree { callChainTreeNode = newCallChainTreeNode; nodes.put(newCallChainTreeNode.getTreeNodeId(), callChainTreeNode); } - callChainTreeNode.summary(treeToken, node); + callChainTreeNode.summary(treeToken, node, summaryType); } } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainTreeNode.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainTreeNode.java index 557aa0d60..0cebefa0a 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainTreeNode.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainTreeNode.java @@ -1,6 +1,7 @@ package com.ai.cloud.skywalking.analysis.chainbuild.entity; import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; +import com.ai.cloud.skywalking.analysis.chainbuild.po.SummaryType; import com.ai.cloud.skywalking.analysis.chainbuild.util.HBaseUtil; import com.ai.cloud.skywalking.analysis.config.Config; import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData; @@ -60,14 +61,21 @@ public class CallChainTreeNode { this.viewPointId = node.getViewPoint(); } - public void summary(String treeId, ChainNode node) throws IOException { + public void summary(String treeId, ChainNode node, SummaryType summaryType) throws IOException { Calendar calendar = Calendar.getInstance(); calendar.setTime(new Date(node.getStartDate())); - - summaryMinResult(treeId, node, calendar); - summaryHourResult(treeId, node, calendar); - summaryDayResult(treeId, node, calendar); - summaryMonthResult(treeId, node, calendar); + if (summaryType == SummaryType.MINUTER) { + summaryMinResult(treeId, node, calendar); + } + if (summaryType == SummaryType.DAY) { + summaryHourResult(treeId, node, calendar); + } + if (summaryType == SummaryType.DAY) { + summaryDayResult(treeId, node, calendar); + } + if (summaryType == SummaryType.MONTH) { + summaryMonthResult(treeId, node, calendar); + } } private void summaryMonthResult(String treeId, ChainNode node, Calendar calendar) throws IOException { diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/SpecificTimeCallTreeMergedChainIdContainer.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/SpecificTimeCallTreeMergedChainIdContainer.java index 7ab2c0d87..4ad3dc400 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/SpecificTimeCallTreeMergedChainIdContainer.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/SpecificTimeCallTreeMergedChainIdContainer.java @@ -42,9 +42,9 @@ public class SpecificTimeCallTreeMergedChainIdContainer { // if (chainInfo.getChainStatus() == ChainInfo.ChainStatus.NORMAL) { - callChainDetailMap.put(chainInfo.getCID(), new CallChainDetailForMysql(chainInfo,treeToken)); + callChainDetailMap.put(chainInfo.getCID(), new CallChainDetailForMysql(chainInfo, treeToken)); } - }else{ + } else { } } @@ -56,14 +56,17 @@ public class SpecificTimeCallTreeMergedChainIdContainer { return treeToken + "@" + calendar.get(Calendar.YEAR) + "-" + calendar.get(Calendar.MONTH); } - public void saveToHBase() throws IOException, InterruptedException, SQLException { + public void saveToHBase(SummaryType summaryType) throws IOException, InterruptedException, SQLException { batchSaveCurrentHasBeenMergedChainInfo(); batchSaveMergedChainId(); - batchSaveToMysql(); + batchSaveToMysql(summaryType); } - private void batchSaveToMysql() throws SQLException { - for (Map.Entry entry : callChainDetailMap.entrySet()){ + private void batchSaveToMysql(SummaryType summaryType) throws SQLException { + if (summaryType != SummaryType.MONTH) { + return; + } + for (Map.Entry entry : callChainDetailMap.entrySet()) { entry.getValue().saveToMysql(); } } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/SummaryType.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/SummaryType.java new file mode 100644 index 000000000..0b4d6a3ac --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/SummaryType.java @@ -0,0 +1,31 @@ +package com.ai.cloud.skywalking.analysis.chainbuild.po; + +public enum SummaryType { + MINUTER('m'), HOUR('H'), DAY('D'), MONTH('M'); + + private char value; + + SummaryType(char value) { + this.value = value; + } + + public static SummaryType convert(String value) { + char valueChar = value.charAt(0); + switch (valueChar) { + case 'm': + return MINUTER; + case 'H': + return HOUR; + case 'D': + return DAY; + case 'M': + return MONTH; + default: + throw new RuntimeException("Can not find the summary type[" + value + "]"); + } + } + + public char getValue() { + return value; + } +}