From fe8911d4aa271683b410bd02c4d702d4604e90df Mon Sep 17 00:00:00 2001 From: wusheng Date: Mon, 2 May 2016 20:12:21 +0800 Subject: [PATCH] =?UTF-8?q?1.=E4=BF=AE=E6=94=B9=E5=A4=A7=E9=87=8F=E7=B1=BB?= =?UTF-8?q?=E5=90=8D=E5=92=8C=E6=96=B9=E6=B3=95=E5=90=8D=202.=E5=A2=9E?= =?UTF-8?q?=E5=8A=A0reduce=E7=9A=84=E5=85=A5=E5=8F=A3=E6=97=A5=E5=BF=97?= =?UTF-8?q?=EF=BC=8C=E4=BB=A5=E5=8F=8A=E5=A4=A7=E6=95=B0=E6=8D=AE=E9=87=8F?= =?UTF-8?q?=E6=89=B9=E9=87=8F=E8=BF=90=E8=A1=8C=E8=AE=A1=E6=95=B0=E5=99=A8?= =?UTF-8?q?=203.=E4=BF=AE=E5=A4=8D=E9=83=A8=E5=88=86=E7=BC=96=E8=AF=91?= =?UTF-8?q?=E5=91=8A=E8=AD=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../chainbuild/ChainBuildReducer.java | 133 +++--- ...maryAction.java => IStatisticsAction.java} | 2 +- ....java => CallChainRelationshipAction.java} | 25 +- ...va => NumberOfCalledStatisticsAction.java} | 8 +- .../entity/CallChainDetailForMysql.java | 12 +- .../chainbuild/entity/CallChainTree.java | 9 +- .../chainbuild/entity/CallChainTreeNode.java | 433 ++++++++++-------- ...> SpecificTimeCallChainTreeContainer.java} | 27 +- .../analysis/chainbuild/po/SummaryType.java | 12 +- .../analysis/mapper/CallChainMapperTest.java | 6 +- 10 files changed, 347 insertions(+), 320 deletions(-) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/{ISummaryAction.java => IStatisticsAction.java} (87%) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/{RelationShipAction.java => CallChainRelationshipAction.java} (55%) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/{DateDetailSummaryAction.java => NumberOfCalledStatisticsAction.java} (83%) rename skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/{SpecificTimeCallTreeMergedChainIdContainer.java => SpecificTimeCallChainTreeContainer.java} (76%) 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 b3e06df32..91e857b5d 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 @@ -1,9 +1,8 @@ package com.ai.cloud.skywalking.analysis.chainbuild; -import com.ai.cloud.skywalking.analysis.chainbuild.action.ISummaryAction; -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 java.io.IOException; +import java.util.Iterator; + import org.apache.hadoop.hbase.util.Bytes; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; @@ -11,81 +10,65 @@ import org.apache.hadoop.mapreduce.Reducer; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; -import java.io.IOException; -import java.util.Iterator; +import com.ai.cloud.skywalking.analysis.chainbuild.action.IStatisticsAction; +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; public class ChainBuildReducer extends Reducer { - private Logger logger = LogManager.getLogger(ChainBuildReducer.class); + private Logger logger = LogManager.getLogger(ChainBuildReducer.class); - @Override - protected void setup(Context context) throws IOException, - InterruptedException { - ConfigInitializer.initialize(); - Config.AnalysisServer.IS_ACCUMULATE_MODE = Boolean.parseBoolean(context.getConfiguration() - .get("skywalking.analysis.mode", "false")); - logger.info("Skywalking analysis mode :[{}]", - Config.AnalysisServer.IS_ACCUMULATE_MODE ? "ACCUMULATE" : "REWRITE"); - } + @Override + protected void setup(Context context) throws IOException, + InterruptedException { + ConfigInitializer.initialize(); + Config.AnalysisServer.IS_ACCUMULATE_MODE = Boolean.parseBoolean(context + .getConfiguration().get("skywalking.analysis.mode", "false")); + logger.info("Skywalking analysis mode :[{}]", + Config.AnalysisServer.IS_ACCUMULATE_MODE ? "ACCUMULATE" + : "REWRITE"); + } - @Override - protected void reduce(Text key, Iterable values, Context context) - throws IOException, InterruptedException { - String reduceKey = Bytes.toString(key.getBytes()); - int index = reduceKey.indexOf(":"); - if (index == -1) { - return; - } - String summaryTypeAndDateStr = reduceKey.substring(0, index - 1); - String entryKey = reduceKey.substring(index + 1); - ISummaryAction summaryAction = SummaryType.chooseSummaryAction(summaryTypeAndDateStr, entryKey); - doReduceAction(summaryAction, values.iterator()); - } + @Override + protected void reduce(Text key, Iterable values, Context context) + throws IOException, InterruptedException { + String reduceKey = Bytes.toString(key.getBytes()); + int index = reduceKey.indexOf(":"); + if (index == -1) { + return; + } + String summaryTypeAndDateStr = reduceKey.substring(0, index - 1); + String entryKey = reduceKey.substring(index + 1); + + logger.debug("begin to reduce for key: {}", reduceKey); + IStatisticsAction summaryAction = SummaryType.chooseSummaryAction( + summaryTypeAndDateStr, entryKey); + doReduceAction(reduceKey, summaryAction, values.iterator()); + } - public void doReduceAction(ISummaryAction summaryAction, Iterator iterator) { - while (iterator.hasNext()) { - String summaryData = iterator.next().toString(); - try { - summaryAction.doAction(summaryData); - } catch (Exception e) { - logger.error( - "Failed to summary call chain, maybe illegal data:" - + summaryData, e); - } - } + public void doReduceAction(String reduceKey, IStatisticsAction summaryAction, + Iterator iterator) { + long dataCounter = 0; + while (iterator.hasNext()) { + String summaryData = iterator.next().toString(); + try { + summaryAction.doAction(summaryData); + } catch (Exception e) { + logger.error( + "Failed to summary call chain, maybe illegal data:" + + summaryData, e); + }finally{ + dataCounter++; + } + if(dataCounter % 1000 == 0){ + logger.debug("reduce for key: {}, count: {}", reduceKey, dataCounter); + } + } - try { - summaryAction.doSave(); - } catch (Exception e) { - logger.error("Failed to save summaryresult/chainTree.", e); - } - } -// -// public void doReduceAction(String key, SummaryType summaryType, String summaryDateStr, Iterator chainInfoIterator) -// throws IOException, InterruptedException { -// CallChainTree chainTree = CallChainTree.load(key); -// SpecificTimeCallTreeMergedChainIdContainer container -// = new SpecificTimeCallTreeMergedChainIdContainer(chainTree.getTreeToken()); -// while (chainInfoIterator.hasNext()) { -// String chainNodeData = chainInfoIterator.next().toString(); -//// ChainInfo chainInfo = null; -// ChainNode chainNode = null; -// try { -//// chainInfo = new Gson().fromJson(callChainData, ChainInfo.class); -// chainNode = new Gson().fromJson(chainNodeData, ChainNode.class); -// // container.addMergedChainNodeIdIfNotContain(chainNode); -// //container.addMergedChainIfNotContain(chainInfo); -// //chainTree.summary(chainInfo, summaryType); -// } catch (Exception e) { -// logger.error( -// "Failed to summary call chain, maybe illegal data:" -// + chainNodeData, e); -// } -// } -// try { -// container.saveToHBase(summaryType); -// chainTree.saveToHBase(); -// } catch (Exception e) { -// logger.error("Failed to save summaryresult/chainTree.", e); -// } -// } + try { + summaryAction.doSave(); + } 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/action/ISummaryAction.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/IStatisticsAction.java similarity index 87% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/ISummaryAction.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/IStatisticsAction.java index 24d12a6a3..6091de928 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/ISummaryAction.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/IStatisticsAction.java @@ -4,7 +4,7 @@ package com.ai.cloud.skywalking.analysis.chainbuild.action; import java.io.IOException; import java.sql.SQLException; -public interface ISummaryAction { +public interface IStatisticsAction { void doAction(String summaryData) throws IOException; void doSave() throws InterruptedException, SQLException, IOException; diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/RelationShipAction.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/CallChainRelationshipAction.java similarity index 55% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/RelationShipAction.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/CallChainRelationshipAction.java index e5c112dc7..325655c85 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/RelationShipAction.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/CallChainRelationshipAction.java @@ -1,27 +1,26 @@ package com.ai.cloud.skywalking.analysis.chainbuild.action.impl; -import com.ai.cloud.skywalking.analysis.chainbuild.action.ISummaryAction; -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.ChainNode; -import com.ai.cloud.skywalking.analysis.chainbuild.po.SpecificTimeCallTreeMergedChainIdContainer; -import com.google.gson.Gson; - import java.io.IOException; import java.sql.SQLException; -public class RelationShipAction implements ISummaryAction { +import com.ai.cloud.skywalking.analysis.chainbuild.action.IStatisticsAction; +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.SpecificTimeCallChainTreeContainer; +import com.google.gson.Gson; + +public class CallChainRelationshipAction implements IStatisticsAction { private CallChainTree chainTree; - private SpecificTimeCallTreeMergedChainIdContainer container; - public RelationShipAction(String entryKey) throws IOException { - chainTree = CallChainTree.load(entryKey); - container = new SpecificTimeCallTreeMergedChainIdContainer(chainTree.getTreeToken()); + private SpecificTimeCallChainTreeContainer container; + public CallChainRelationshipAction(String entryKey) throws IOException { + chainTree = CallChainTree.create(entryKey); + container = new SpecificTimeCallChainTreeContainer(chainTree.getTreeToken()); } @Override public void doAction(String summaryData) throws IOException { ChainInfo chainInfo = new Gson().fromJson(summaryData, ChainInfo.class); - container.addMergedChainIfNotContain(chainInfo); + container.addChainIfNew(chainInfo); } @Override diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/DateDetailSummaryAction.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/NumberOfCalledStatisticsAction.java similarity index 83% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/DateDetailSummaryAction.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/NumberOfCalledStatisticsAction.java index 3d32327d2..f96f76082 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/DateDetailSummaryAction.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/action/impl/NumberOfCalledStatisticsAction.java @@ -1,6 +1,6 @@ package com.ai.cloud.skywalking.analysis.chainbuild.action.impl; -import com.ai.cloud.skywalking.analysis.chainbuild.action.ISummaryAction; +import com.ai.cloud.skywalking.analysis.chainbuild.action.IStatisticsAction; import com.ai.cloud.skywalking.analysis.chainbuild.entity.CallChainTree; import com.ai.cloud.skywalking.analysis.chainbuild.entity.CallChainTreeNode; import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode; @@ -10,13 +10,13 @@ import com.google.gson.Gson; import java.io.IOException; import java.sql.SQLException; -public class DateDetailSummaryAction implements ISummaryAction { +public class NumberOfCalledStatisticsAction implements IStatisticsAction { private CallChainTree callChainTree; private String summaryDate; private SummaryType summaryType; - public DateDetailSummaryAction(String entryKey, String summaryDate) throws IOException { - callChainTree = CallChainTree.load(entryKey); + public NumberOfCalledStatisticsAction(String entryKey, String summaryDate) throws IOException { + callChainTree = CallChainTree.create(entryKey); this.summaryDate = summaryDate; } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainDetailForMysql.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainDetailForMysql.java index 206d1e372..65d826ff1 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainDetailForMysql.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/entity/CallChainDetailForMysql.java @@ -1,17 +1,15 @@ package com.ai.cloud.skywalking.analysis.chainbuild.entity; -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; import java.util.Collection; import java.util.HashMap; import java.util.Map; +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.google.gson.Gson; + public class CallChainDetailForMysql { private String chainToken; private String treeToken; 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 1ee88e2a0..72f0ff524 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 @@ -4,13 +4,12 @@ import java.io.IOException; import java.util.HashMap; 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; +import com.ai.cloud.skywalking.analysis.chainbuild.po.SummaryType; +import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator; + public class CallChainTree { private Logger logger = LogManager.getLogger(CallChainTree.class); @@ -31,7 +30,7 @@ public class CallChainTree { logger.info("CallEntrance:[{}] == TreeToken[{}]",callEntrance, treeToken); } - public static CallChainTree load(String callEntrance) throws IOException { + public static CallChainTree create(String callEntrance) throws IOException { CallChainTree chain = new CallChainTree(callEntrance); return chain; } 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 2ff895796..b2a737272 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 @@ -21,220 +21,269 @@ import java.util.*; * @author wusheng */ public class CallChainTreeNode { - private Logger logger = LogManager.getLogger(CallChainTreeNode.class); + private Logger logger = LogManager.getLogger(CallChainTreeNode.class); - @Expose - private String traceLevelId; - @Expose - private String viewPointId; + @Expose + private String traceLevelId; + @Expose + private String viewPointId; - /** - * key: treeId + 小时 - * value: 当前树的当前小时范围内的,所有分钟和节点的统计数据 - */ - private Map chainNodeSpecificMinSummaryContainer; + /** + * key: treeId + 小时 value: 当前树的当前小时范围内的,所有分钟和节点的统计数据 + */ + private Map chainNodeSpecificMinSummaryContainer; - /** - * key: treeId + 天 - * value: 当前树的当前天范围内的,所有小时和节点的统计数据 - */ - private Map chainNodeSpecificHourSummaryContainer; + /** + * key: treeId + 天 value: 当前树的当前天范围内的,所有小时和节点的统计数据 + */ + private Map chainNodeSpecificHourSummaryContainer; - /** - * key: treeId + 月 - * value: 当前树的当前月范围内的,所有天和节点的统计数据 - */ - private Map chainNodeSpecificDaySummaryContainer; + /** + * key: treeId + 月 value: 当前树的当前月范围内的,所有天和节点的统计数据 + */ + private Map chainNodeSpecificDaySummaryContainer; - /** - * key: treeId + 年 - * value: 当前树的当前年范围内的,所有月份和节点的统计数据 - */ - private Map chainNodeSpecificMonthSummaryContainer; + /** + * key: treeId + 年 value: 当前树的当前年范围内的,所有月份和节点的统计数据 + */ + private Map chainNodeSpecificMonthSummaryContainer; - public CallChainTreeNode(ChainNode node) { - this.traceLevelId = node.getTraceLevelId(); - chainNodeSpecificMinSummaryContainer = new HashMap(); - chainNodeSpecificHourSummaryContainer = new HashMap(); - chainNodeSpecificDaySummaryContainer = new HashMap(); - chainNodeSpecificMonthSummaryContainer = new HashMap(); - this.viewPointId = node.getViewPoint(); - } + public CallChainTreeNode(ChainNode node) { + this.traceLevelId = node.getTraceLevelId(); + chainNodeSpecificMinSummaryContainer = new HashMap(); + chainNodeSpecificHourSummaryContainer = new HashMap(); + chainNodeSpecificDaySummaryContainer = new HashMap(); + chainNodeSpecificMonthSummaryContainer = new HashMap(); + this.viewPointId = node.getViewPoint(); + } - public void summary(String treeId, ChainNode node, SummaryType summaryType, String summaryDate) throws IOException { - Calendar calendar = Calendar.getInstance(); - calendar.setTime(new Date(node.getStartDate())); - switch (summaryType) { - case HOUR: { - summaryMinResult(treeId, node, calendar, summaryDate); - break; - } - case DAY: { - summaryHourResult(treeId, node, calendar, summaryDate); - break; - } - case MONTH: { - summaryDayResult(treeId, node, calendar, summaryDate); - break; - } - case YEAR: { - summaryMonthResult(treeId, node, calendar, summaryDate); - break; - } - } + /** + * 针对节点、汇总类型,进行数据汇总。
+ * 已汇总数据加载使用lazy load模式,只有此次汇总过程中触发的数据节点统计数量,才会被读取
+ * + * @param treeId + * @param node + * @param summaryType + * @param summaryDate + * @throws IOException + */ + public void summary(String treeId, ChainNode node, SummaryType summaryType, + String summaryDate) throws IOException { + Calendar calendar = Calendar.getInstance(); + calendar.setTime(new Date(node.getStartDate())); + switch (summaryType) { + case HOUR: { + summaryMinResult(treeId, node, calendar, summaryDate); + break; + } + case DAY: { + summaryHourResult(treeId, node, calendar, summaryDate); + break; + } + case MONTH: { + summaryDayResult(treeId, node, calendar, summaryDate); + break; + } + case YEAR: { + summaryMonthResult(treeId, node, calendar, summaryDate); + break; + } + default: { + logger.error("unknown summary type :{}", summaryType); + break; + } + } - } + } - private void summaryMonthResult(String treeId, ChainNode node, Calendar calendar, String summaryDate) throws IOException { - String keyOfMonthSummaryTable = generateRowKey(treeId, summaryDate); - ChainNodeSpecificMonthSummary monthSummary = chainNodeSpecificMonthSummaryContainer.get(keyOfMonthSummaryTable); - if (monthSummary == null) { - if (Config.AnalysisServer.IS_ACCUMULATE_MODE) { - monthSummary = HBaseUtil.loadSpecificMonthSummary(keyOfMonthSummaryTable, getTreeNodeId()); - } else { - monthSummary = new ChainNodeSpecificMonthSummary(); - } - chainNodeSpecificMonthSummaryContainer.put(keyOfMonthSummaryTable, monthSummary); - } - monthSummary.summary(String.valueOf(calendar.get(Calendar.MONTH) + 1), node); - } + private void summaryMonthResult(String treeId, ChainNode node, + Calendar calendar, String summaryDate) throws IOException { + String keyOfMonthSummaryTable = generateRowKey(treeId, summaryDate); + ChainNodeSpecificMonthSummary monthSummary = chainNodeSpecificMonthSummaryContainer + .get(keyOfMonthSummaryTable); + if (monthSummary == null) { + if (Config.AnalysisServer.IS_ACCUMULATE_MODE) { + monthSummary = HBaseUtil.loadSpecificMonthSummary( + keyOfMonthSummaryTable, getTreeNodeId()); + } else { + monthSummary = new ChainNodeSpecificMonthSummary(); + } + chainNodeSpecificMonthSummaryContainer.put(keyOfMonthSummaryTable, + monthSummary); + } + monthSummary.summary(String.valueOf(calendar.get(Calendar.MONTH) + 1), + node); + } - private void summaryDayResult(String treeId, ChainNode node, Calendar calendar, String summaryDate) throws IOException { - String keyOfDaySummaryTable = generateRowKey(treeId, summaryDate); - ChainNodeSpecificDaySummary daySummary = chainNodeSpecificDaySummaryContainer.get(keyOfDaySummaryTable); - if (daySummary == null) { - if (Config.AnalysisServer.IS_ACCUMULATE_MODE) { - daySummary = HBaseUtil.loadSpecificDaySummary(keyOfDaySummaryTable, getTreeNodeId()); - } else { - daySummary = new ChainNodeSpecificDaySummary(); - } - chainNodeSpecificDaySummaryContainer.put(keyOfDaySummaryTable, daySummary); - } - daySummary.summary(String.valueOf(calendar.get(Calendar.DAY_OF_MONTH)), node); - } + private void summaryDayResult(String treeId, ChainNode node, + Calendar calendar, String summaryDate) throws IOException { + String keyOfDaySummaryTable = generateRowKey(treeId, summaryDate); + ChainNodeSpecificDaySummary daySummary = chainNodeSpecificDaySummaryContainer + .get(keyOfDaySummaryTable); + if (daySummary == null) { + if (Config.AnalysisServer.IS_ACCUMULATE_MODE) { + daySummary = HBaseUtil.loadSpecificDaySummary( + keyOfDaySummaryTable, getTreeNodeId()); + } else { + daySummary = new ChainNodeSpecificDaySummary(); + } + chainNodeSpecificDaySummaryContainer.put(keyOfDaySummaryTable, + daySummary); + } + daySummary.summary(String.valueOf(calendar.get(Calendar.DAY_OF_MONTH)), + node); + } - private void summaryHourResult(String treeId, ChainNode node, Calendar calendar, String summaryDate) throws IOException { - String keyOfHourSummaryTable = generateRowKey(treeId, summaryDate); - ChainNodeSpecificHourSummary hourSummary = chainNodeSpecificHourSummaryContainer.get(keyOfHourSummaryTable); - if (hourSummary == null) { - if (Config.AnalysisServer.IS_ACCUMULATE_MODE) { - hourSummary = HBaseUtil.loadSpecificHourSummary(keyOfHourSummaryTable, getTreeNodeId()); - } else { - hourSummary = new ChainNodeSpecificHourSummary(); - } - chainNodeSpecificHourSummaryContainer.put(keyOfHourSummaryTable, hourSummary); - } - hourSummary.summary(String.valueOf(calendar.get(Calendar.HOUR)), node); - } + private void summaryHourResult(String treeId, ChainNode node, + Calendar calendar, String summaryDate) throws IOException { + String keyOfHourSummaryTable = generateRowKey(treeId, summaryDate); + ChainNodeSpecificHourSummary hourSummary = chainNodeSpecificHourSummaryContainer + .get(keyOfHourSummaryTable); + if (hourSummary == null) { + if (Config.AnalysisServer.IS_ACCUMULATE_MODE) { + hourSummary = HBaseUtil.loadSpecificHourSummary( + keyOfHourSummaryTable, getTreeNodeId()); + } else { + hourSummary = new ChainNodeSpecificHourSummary(); + } + chainNodeSpecificHourSummaryContainer.put(keyOfHourSummaryTable, + hourSummary); + } + hourSummary.summary(String.valueOf(calendar.get(Calendar.HOUR)), node); + } - /** - * 按分钟维度进行汇总
- * chainNodeContainer以treeId和时间(精确到分钟)为key,value为当前时间范围内的所有分钟的汇总数据 - */ - private void summaryMinResult(String treeId, ChainNode node, Calendar calendar, String summaryDate) throws IOException { - String keyOfMinSummaryTable = generateRowKey(treeId, summaryDate); - ChainNodeSpecificMinSummary minSummary = chainNodeSpecificMinSummaryContainer.get(keyOfMinSummaryTable); - if (minSummary == null) { - if (Config.AnalysisServer.IS_ACCUMULATE_MODE) { - minSummary = HBaseUtil.loadSpecificMinSummary(keyOfMinSummaryTable, getTreeNodeId()); - } else { - minSummary = new ChainNodeSpecificMinSummary(); - } - chainNodeSpecificMinSummaryContainer.put(keyOfMinSummaryTable, minSummary); - } - minSummary.summary(String.valueOf(calendar.get(Calendar.MINUTE)), node); - } + /** + * 按分钟维度进行汇总
+ * chainNodeContainer以treeId和时间(精确到分钟)为key,value为当前时间范围内的所有分钟的汇总数据 + */ + private void summaryMinResult(String treeId, ChainNode node, + Calendar calendar, String summaryDate) throws IOException { + String keyOfMinSummaryTable = generateRowKey(treeId, summaryDate); + ChainNodeSpecificMinSummary minSummary = chainNodeSpecificMinSummaryContainer + .get(keyOfMinSummaryTable); + if (minSummary == null) { + if (Config.AnalysisServer.IS_ACCUMULATE_MODE) { + minSummary = HBaseUtil.loadSpecificMinSummary( + keyOfMinSummaryTable, getTreeNodeId()); + } else { + minSummary = new ChainNodeSpecificMinSummary(); + } + chainNodeSpecificMinSummaryContainer.put(keyOfMinSummaryTable, + minSummary); + } + minSummary.summary(String.valueOf(calendar.get(Calendar.MINUTE)), node); + } - private String generateRowKey(String treeId, String dateKey) { - return treeId + "/" + dateKey; - } + private String generateRowKey(String treeId, String dateKey) { + return treeId + "/" + dateKey; + } + @Override + public String toString() { + return new GsonBuilder().excludeFieldsWithoutExposeAnnotation() + .create().toJson(this); + } - @Override - public String toString() { - return new GsonBuilder().excludeFieldsWithoutExposeAnnotation().create().toJson(this); - } + /** + * 存储入库时
+ * hbase的key 为 treeId + 小时
+ * 列族中,列为节点id,规则为:traceLevelId + "@" + viewPointId
+ * 列的值,为当前节点按小时内各分钟的汇总
+ * + * @throws IOException + * @throws InterruptedException + */ + public void saveSummaryResultToHBase(SummaryType summaryType) + throws IOException, InterruptedException { + switch (summaryType) { + case HOUR: { + batchSaveMinSummaryResult(); + break; + } + case DAY: { + batchSaveHourSummaryResult(); + break; + } + case MONTH: { + batchSaveDaySummaryResult(); + break; + } + case YEAR: { + batchSaveMonthSummaryResult(); + break; + } + default: { + logger.error("unknown summary type :{}", summaryType); + break; + } + } + } - /** - * 存储入库时
- * hbase的key 为 treeId + 小时
- * 列族中,列为节点id,规则为:traceLevelId + "@" + viewPointId
- * 列的值,为当前节点按小时内各分钟的汇总
- * - * @throws IOException - * @throws InterruptedException - */ - public void saveSummaryResultToHBase(SummaryType summaryType) throws IOException, InterruptedException { - switch (summaryType) { - case HOUR: { - batchSaveMinSummaryResult(); - break; - } - case DAY: { - batchSaveHourSummaryResult(); - break; - } - case MONTH: { - batchSaveDaySummaryResult(); - break; - } - case YEAR: { - batchSaveMonthSummaryResult(); - break; - } - } - } + private void batchSaveMonthSummaryResult() throws IOException, + InterruptedException { + List puts = new ArrayList(); + for (Map.Entry entry : chainNodeSpecificMonthSummaryContainer + .entrySet()) { + Put put = new Put(entry.getKey().getBytes()); + put.addColumn( + HBaseTableMetaData.TABLE_CHAIN_ONE_MONTH_SUMMARY.COLUMN_FAMILY_NAME + .getBytes(), getTreeNodeId().getBytes(), entry + .getValue().toString().getBytes()); + puts.add(put); + } - private void batchSaveMonthSummaryResult() throws IOException, InterruptedException { - List puts = new ArrayList(); - for (Map.Entry entry : chainNodeSpecificMonthSummaryContainer.entrySet()) { - Put put = new Put(entry.getKey().getBytes()); - put.addColumn(HBaseTableMetaData.TABLE_CHAIN_ONE_MONTH_SUMMARY.COLUMN_FAMILY_NAME.getBytes() - , getTreeNodeId().getBytes(), entry.getValue().toString().getBytes()); - puts.add(put); - } + HBaseUtil.batchSaveMonthSummaryResult(puts); + } - HBaseUtil.batchSaveMonthSummaryResult(puts); - } + private void batchSaveDaySummaryResult() throws IOException, + InterruptedException { + List puts = new ArrayList(); + for (Map.Entry entry : chainNodeSpecificDaySummaryContainer + .entrySet()) { + Put put = new Put(entry.getKey().getBytes()); + put.addColumn( + HBaseTableMetaData.TABLE_CHAIN_ONE_DAY_SUMMARY.COLUMN_FAMILY_NAME + .getBytes(), getTreeNodeId().getBytes(), entry + .getValue().toString().getBytes()); + puts.add(put); + } + HBaseUtil.batchSaveDaySummaryResult(puts); + } - private void batchSaveDaySummaryResult() throws IOException, InterruptedException { - List puts = new ArrayList(); - for (Map.Entry entry : chainNodeSpecificDaySummaryContainer.entrySet()) { - Put put = new Put(entry.getKey().getBytes()); - put.addColumn(HBaseTableMetaData.TABLE_CHAIN_ONE_DAY_SUMMARY.COLUMN_FAMILY_NAME.getBytes() - , getTreeNodeId().getBytes(), entry.getValue().toString().getBytes()); - puts.add(put); - } + private void batchSaveHourSummaryResult() throws IOException, + InterruptedException { + List puts = new ArrayList(); + for (Map.Entry entry : chainNodeSpecificHourSummaryContainer + .entrySet()) { + Put put = new Put(entry.getKey().getBytes()); + put.addColumn( + HBaseTableMetaData.TABLE_CHAIN_ONE_HOUR_SUMMARY.COLUMN_FAMILY_NAME + .getBytes(), getTreeNodeId().getBytes(), entry + .getValue().toString().getBytes()); + puts.add(put); + } - HBaseUtil.batchSaveDaySummaryResult(puts); - } + HBaseUtil.batchSaveHourSummaryResult(puts); + } - private void batchSaveHourSummaryResult() throws IOException, InterruptedException { - List puts = new ArrayList(); - for (Map.Entry entry : chainNodeSpecificHourSummaryContainer.entrySet()) { - Put put = new Put(entry.getKey().getBytes()); - put.addColumn(HBaseTableMetaData.TABLE_CHAIN_ONE_HOUR_SUMMARY.COLUMN_FAMILY_NAME.getBytes() - , getTreeNodeId().getBytes(), entry.getValue().toString().getBytes()); - puts.add(put); - } + private void batchSaveMinSummaryResult() throws IOException, + InterruptedException { + List puts = new ArrayList(); + for (Map.Entry entry : chainNodeSpecificMinSummaryContainer + .entrySet()) { + Put put = new Put(entry.getKey().getBytes()); + put.addColumn( + HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY.COLUMN_FAMILY_NAME + .getBytes(), getTreeNodeId().getBytes(), entry + .getValue().toString().getBytes()); + puts.add(put); + } - HBaseUtil.batchSaveHourSummaryResult(puts); - } + HBaseUtil.batchSaveMinSummaryResult(puts); + } - private void batchSaveMinSummaryResult() throws IOException, InterruptedException { - List puts = new ArrayList(); - for (Map.Entry entry : chainNodeSpecificMinSummaryContainer.entrySet()) { - Put put = new Put(entry.getKey().getBytes()); - put.addColumn(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY.COLUMN_FAMILY_NAME.getBytes() - , getTreeNodeId().getBytes(), entry.getValue().toString().getBytes()); - puts.add(put); - } - - HBaseUtil.batchSaveMinSummaryResult(puts); - } - - public String getTreeNodeId() { - return traceLevelId + "@" + viewPointId; - } + public String getTreeNodeId() { + return traceLevelId + "@" + viewPointId; + } } 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/SpecificTimeCallChainTreeContainer.java similarity index 76% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/SpecificTimeCallTreeMergedChainIdContainer.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chainbuild/po/SpecificTimeCallChainTreeContainer.java index 941242dcb..6d048a853 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/SpecificTimeCallChainTreeContainer.java @@ -10,25 +10,25 @@ import java.io.IOException; import java.sql.SQLException; import java.util.*; -public class SpecificTimeCallTreeMergedChainIdContainer { +public class SpecificTimeCallChainTreeContainer { private String treeToken; private Map> hasBeenMergedChainIds; - private Map callChainDetailMap; + private Map newChain4DB; // 本次Reduce合并过的调用链 - private Map combineChains; + private Map newChains; - public SpecificTimeCallTreeMergedChainIdContainer(String treeToken) { + public SpecificTimeCallChainTreeContainer(String treeToken) { this.treeToken = treeToken; hasBeenMergedChainIds = new HashMap>(); - combineChains = new HashMap(); - callChainDetailMap = new HashMap(); + newChains = new HashMap(); + newChain4DB = new HashMap(); } - public void addMergedChainIfNotContain(ChainInfo chainInfo) throws IOException { + public void addChainIfNew(ChainInfo chainInfo) throws IOException { String key = generateKey(chainInfo.getStartDate()); List cIds = hasBeenMergedChainIds.get(key); if (cIds == null) { @@ -38,11 +38,10 @@ public class SpecificTimeCallTreeMergedChainIdContainer { if (!cIds.contains(chainInfo.getCID())) { cIds.add(chainInfo.getCID()); - combineChains.put(chainInfo.getCID(), chainInfo); + newChains.put(chainInfo.getCID(), chainInfo); - // if (chainInfo.getChainStatus() == ChainInfo.ChainStatus.NORMAL) { - callChainDetailMap.put(chainInfo.getCID(), new CallChainDetailForMysql(chainInfo, treeToken)); + newChain4DB.put(chainInfo.getCID(), new CallChainDetailForMysql(chainInfo, treeToken)); } } } @@ -55,13 +54,13 @@ public class SpecificTimeCallTreeMergedChainIdContainer { } public void saveToHBase() throws IOException, InterruptedException, SQLException { - batchSaveCurrentHasBeenMergedChainInfo(); + batchSaveNewChainsInfo(); batchSaveMergedChainId(); batchSaveToMysql(); } private void batchSaveToMysql() throws SQLException { - for (Map.Entry entry : callChainDetailMap.entrySet()) { + for (Map.Entry entry : newChain4DB.entrySet()) { entry.getValue().saveToMysql(); } } @@ -90,9 +89,9 @@ public class SpecificTimeCallTreeMergedChainIdContainer { * @throws IOException * @throws InterruptedException */ - private void batchSaveCurrentHasBeenMergedChainInfo() throws IOException, InterruptedException { + private void batchSaveNewChainsInfo() throws IOException, InterruptedException { List chainInfoPuts = new ArrayList(); - for (Map.Entry entry : combineChains.entrySet()) { + for (Map.Entry entry : newChains.entrySet()) { Put put = new Put(entry.getKey().getBytes()); entry.getValue().saveToHBase(put); chainInfoPuts.add(put); 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 index 5559f5eaa..19ea679fd 100644 --- 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 @@ -1,8 +1,8 @@ package com.ai.cloud.skywalking.analysis.chainbuild.po; -import com.ai.cloud.skywalking.analysis.chainbuild.action.ISummaryAction; -import com.ai.cloud.skywalking.analysis.chainbuild.action.impl.DateDetailSummaryAction; -import com.ai.cloud.skywalking.analysis.chainbuild.action.impl.RelationShipAction; +import com.ai.cloud.skywalking.analysis.chainbuild.action.IStatisticsAction; +import com.ai.cloud.skywalking.analysis.chainbuild.action.impl.NumberOfCalledStatisticsAction; +import com.ai.cloud.skywalking.analysis.chainbuild.action.impl.CallChainRelationshipAction; import java.io.IOException; @@ -19,7 +19,7 @@ public enum SummaryType { return value; } - public static ISummaryAction chooseSummaryAction(String summaryTypeAndDateStr, String entryKey) throws IOException { + public static IStatisticsAction chooseSummaryAction(String summaryTypeAndDateStr, String entryKey) throws IOException { char valueChar = summaryTypeAndDateStr.charAt(0); // HOUR : 2016-05-02/12 // DAY : 2016-05-02 @@ -42,11 +42,11 @@ public enum SummaryType { type = YEAR; break; case 'R': - return new RelationShipAction(entryKey); + return new CallChainRelationshipAction(entryKey); default: throw new RuntimeException("Can not find the summary type[" + valueChar + "]"); } - DateDetailSummaryAction summaryAction = new DateDetailSummaryAction(entryKey, summaryDateStr); + NumberOfCalledStatisticsAction summaryAction = new NumberOfCalledStatisticsAction(entryKey, summaryDateStr); summaryAction.setSummaryType(type); return summaryAction; } diff --git a/skywalking-analysis/src/test/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapperTest.java b/skywalking-analysis/src/test/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapperTest.java index a58ded1dd..18b3bbe14 100644 --- a/skywalking-analysis/src/test/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapperTest.java +++ b/skywalking-analysis/src/test/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapperTest.java @@ -3,7 +3,7 @@ package com.ai.cloud.skywalking.analysis.mapper; import com.ai.cloud.skywalking.analysis.chainbuild.ChainBuildMapper; import com.ai.cloud.skywalking.analysis.chainbuild.ChainBuildReducer; -import com.ai.cloud.skywalking.analysis.chainbuild.action.ISummaryAction; +import com.ai.cloud.skywalking.analysis.chainbuild.action.IStatisticsAction; import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo; import com.ai.cloud.skywalking.analysis.chainbuild.po.SummaryType; import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator; @@ -62,9 +62,9 @@ public class CallChainMapperTest { String summaryTypeAndDateStr = reduceKey.substring(0, index - 1); String entryKey = reduceKey.substring(index + 1); - ISummaryAction summaryAction = SummaryType.chooseSummaryAction(summaryTypeAndDateStr, entryKey); + IStatisticsAction summaryAction = SummaryType.chooseSummaryAction(summaryTypeAndDateStr, entryKey); - new ChainBuildReducer().doReduceAction(summaryAction, chainNodeInfo.iterator()); + new ChainBuildReducer().doReduceAction(reduceKey, summaryAction, chainNodeInfo.iterator()); } public static List selectByTraceId(String traceId) throws IOException {