diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java index 9f1432ecc..a752fb298 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java @@ -49,7 +49,7 @@ public class Categorize2ChainMapper extends TableMapper { chainInfo = spanToChainInfo(Bytes.toString(key.get()), spanList); logger.info("Success convert span to chain info...." - + chainInfo.getCID()); + + chainInfo.getCID() + " TraceId : " + Bytes.toString(key.get())); context.write( new Text(chainInfo.getUserId() + ":" + chainInfo.getEntranceNodeToken()), chainInfo); 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 index 378648cde..f52af2e7f 100644 --- 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 @@ -31,8 +31,8 @@ public class Categorize2ChainReducer extends Reducer chainInfoIterator) throws IOException, InterruptedException { int totalCount = 0; try { - ChainRelationship chainRelate = HBaseUtil.selectCallChainRelationship(key.toString()); - ChainSummary summary = new ChainSummary(); + ChainRelationship chainRelate = HBaseUtil.loadCallChainRelationship(key.toString()); + ChainSummaryWithoutRelationship summary = new ChainSummaryWithoutRelationship(); while (chainInfoIterator.hasNext()) { ChainInfo chainInfo = chainInfoIterator.next(); try { diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/CategorizedChainInfo.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/CategorizedChainInfo.java index 6304d5dd8..50b7b7246 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/CategorizedChainInfo.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/CategorizedChainInfo.java @@ -7,7 +7,9 @@ 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; @@ -18,7 +20,7 @@ public class CategorizedChainInfo { private List children_Token; public CategorizedChainInfo(ChainInfo chainInfo) { - cid = chainInfo.getCID(); + cid = chainInfo.getCID(); StringBuilder stringBuilder = new StringBuilder(); boolean flag = false; @@ -26,7 +28,7 @@ public class CategorizedChainInfo { if (flag) { stringBuilder.append(";"); } - stringBuilder.append(chainNode.getNodeToken()); + stringBuilder.append(chainNode.getTraceLevelId() + "-" + chainNode.getNodeToken()); flag = true; } @@ -52,6 +54,10 @@ public class CategorizedChainInfo { 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()); } @@ -64,4 +70,8 @@ public class CategorizedChainInfo { 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/ChainNodeSpecificTimeWindowSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainNodeSpecificTimeWindowSummary.java index 165e52123..f67b71ea9 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainNodeSpecificTimeWindowSummary.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainNodeSpecificTimeWindowSummary.java @@ -13,12 +13,14 @@ 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) { + public static ChainNodeSpecificTimeWindowSummary newInstance(String traceLevelId, String nodeToken) { ChainNodeSpecificTimeWindowSummary cns = new ChainNodeSpecificTimeWindowSummary(); cns.traceLevelId = traceLevelId; + cns.nodeToken = nodeToken; return cns; } @@ -32,6 +34,7 @@ public class ChainNodeSpecificTimeWindowSummary { summerValueMap = new Gson().fromJson(jsonObject.get("summerValueMap").toString(), new TypeToken>() { }.getType()); + nodeToken = jsonObject.get("nodeToken").getAsString(); } public String getTraceLevelId() { @@ -57,4 +60,12 @@ public class ChainNodeSpecificTimeWindowSummary { 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/ChainNodeSpecificTimeWindowSummaryValue.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainNodeSpecificTimeWindowSummaryValue.java index 8bf09d8dc..b3e91788f 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainNodeSpecificTimeWindowSummaryValue.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainNodeSpecificTimeWindowSummaryValue.java @@ -19,18 +19,10 @@ public class ChainNodeSpecificTimeWindowSummaryValue { return totalCall; } - public void setTotalCall(long totalCall) { - this.totalCall = totalCall; - } - public long getTotalCostTime() { return totalCostTime; } - public void setTotalCostTime(long totalCostTime) { - this.totalCostTime = totalCostTime; - } - public long getCorrectNumber() { return correctNumber; } @@ -39,10 +31,6 @@ public class ChainNodeSpecificTimeWindowSummaryValue { return humanInterruptionNumber; } - public void setCorrectNumber(long correctNumber) { - this.correctNumber = correctNumber; - } - public void summary(ChainNode node) { totalCall++; if (node.getStatus() == ChainNode.NodeStatus.NORMAL) { @@ -53,4 +41,11 @@ public class ChainNodeSpecificTimeWindowSummaryValue { } 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/ChainRelationship.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainRelationship.java index 3f7faa273..b4478f162 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainRelationship.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainRelationship.java @@ -50,6 +50,7 @@ public class ChainRelationship { 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; diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainSpecificTimeWindowSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainSpecificTimeWindowSummary.java index 1e1b435b8..aeaeb780f 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainSpecificTimeWindowSummary.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainSpecificTimeWindowSummary.java @@ -15,7 +15,7 @@ public class ChainSpecificTimeWindowSummary { private static Logger logger = LoggerFactory.getLogger(ChainSpecificTimeWindowSummary.class.getName()); /** - * key : cid + 时间窗口 + * key : cid + uid + 时间窗口 */ private Map chainNodeSummaryResultMap; @@ -23,12 +23,12 @@ public class ChainSpecificTimeWindowSummary { chainNodeSummaryResultMap = new HashMap(); } - public static ChainSpecificTimeWindowSummary load(String cid_time) { + public static ChainSpecificTimeWindowSummary load(String cid_uid_time) { ChainSpecificTimeWindowSummary result = null; try { - result = HBaseUtil.selectChainSummaryResult(cid_time); + result = HBaseUtil.selectChainSummaryResult(cid_uid_time); } catch (IOException e) { - logger.error("Failed to load the key[" + cid_time + "] summary result.", e); + logger.error("Failed to load the key[" + cid_uid_time + "] summary result.", e); } if (result == null) { @@ -45,7 +45,7 @@ public class ChainSpecificTimeWindowSummary { String tlid = node.getTraceLevelId(); ChainNodeSpecificTimeWindowSummary chainNodeSummaryResult = chainNodeSummaryResultMap.get(tlid); if (chainNodeSummaryResult == null) { - chainNodeSummaryResult = ChainNodeSpecificTimeWindowSummary.newInstance(tlid); + chainNodeSummaryResult = ChainNodeSpecificTimeWindowSummary.newInstance(tlid, node.getNodeToken()); chainNodeSummaryResultMap.put(tlid, chainNodeSummaryResult); } chainNodeSummaryResult.summary(node); diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainSummaryWithoutRelationship.java similarity index 95% rename from skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainSummary.java rename to skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainSummaryWithoutRelationship.java index c15346306..2c43f2e0c 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainSummary.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/ChainSummaryWithoutRelationship.java @@ -13,13 +13,13 @@ import java.sql.Timestamp; import java.text.SimpleDateFormat; import java.util.*; -public class ChainSummary { +public class ChainSummaryWithoutRelationship { - private static Logger logger = LoggerFactory.getLogger(ChainSummary.class.getName()); + private static Logger logger = LoggerFactory.getLogger(ChainSummaryWithoutRelationship.class.getName()); private Map loadedChainSpecificTimeWindowSummary; private Map updateChainInfo; - public ChainSummary() { + public ChainSummaryWithoutRelationship() { loadedChainSpecificTimeWindowSummary = new HashMap(); updateChainInfo = new HashMap(); } 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/categorize2chain/DBCallChainInfoDao.java index 760829068..d0cd5510d 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/categorize2chain/DBCallChainInfoDao.java @@ -21,7 +21,7 @@ public class DBCallChainInfoDao { connection = DriverManager.getConnection(Config.MySql.URL, Config.MySql.USERNAME, Config.MySql.PASSWORD); } catch (ClassNotFoundException e) { - logger.error("Failed to find jdbc driver class[" + logger.error("Failed to searchRelationship jdbc driver class[" + Config.MySql.DRIVER_CLASS + "]", e); System.exit(-1); } catch (SQLException e) { 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/categorize2chain/SpanEntry.java index a0f2a1be2..fa3636dcb 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/categorize2chain/SpanEntry.java @@ -102,7 +102,7 @@ public class SpanEntry { if (serverSpan != null) { if (serverSpan.getExceptionStack() != null && serverSpan.getExceptionStack().length() > 0) { - if(clientSpan.getStatusCode() == 1){ + if(clientSpan != null && clientSpan.getStatusCode() == 1){ return ChainNode.NodeStatus.ABNORMAL; }else{ return ChainNode.NodeStatus.HUMAN_INTERRUPTION; 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/categorize2chain/filter/SpanNodeProcessChain.java index fa917c7ca..d3bb6a3d3 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/categorize2chain/filter/SpanNodeProcessChain.java @@ -25,7 +25,7 @@ public class SpanNodeProcessChain { try { properties.load(SpanNodeProcessChain.class.getResourceAsStream("/viewpointfilter.conf")); } catch (IOException e) { - logger.error("Failed to find config file[viewpointfilter.conf]", e); + logger.error("Failed to searchRelationship config file[viewpointfilter.conf]", e); System.exit(-1); } @@ -41,7 +41,7 @@ public class SpanNodeProcessChain { tmpSpanNodeFilter.setNextProcessChain(currentFilter); currentFilter = tmpSpanNodeFilter; } catch (ClassNotFoundException e) { - logger.error("Filed to find class[" + Config.Filter.FILTER_PACKAGE_NAME + "." + filters[i] + "]", e); + logger.error("Filed to searchRelationship class[" + Config.Filter.FILTER_PACKAGE_NAME + "." + filters[i] + "]", e); System.exit(-1); } catch (InstantiationException e) { logger.error("Can not instance class[" + Config.Filter.FILTER_PACKAGE_NAME + "." + filters[i] + "]", e); 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/categorize2chain/util/HBaseUtil.java index f6fb6f7ca..47544b93e 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/categorize2chain/util/HBaseUtil.java @@ -2,6 +2,8 @@ package com.ai.cloud.skywalking.analysis.categorize2chain.util; import com.ai.cloud.skywalking.analysis.categorize2chain.*; import com.ai.cloud.skywalking.analysis.categorize2chain.model.ChainInfo; +import com.ai.cloud.skywalking.analysis.chain2summary.ChainRelationship4Search; +import com.ai.cloud.skywalking.analysis.chain2summary.model.*; import com.ai.cloud.skywalking.analysis.config.Config; import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData; import com.google.gson.Gson; @@ -14,7 +16,9 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.IOException; +import java.util.ArrayList; import java.util.List; +import java.util.Map; public class HBaseUtil { private static Logger logger = LoggerFactory.getLogger(HBaseUtil.class.getName()); @@ -26,13 +30,29 @@ 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_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_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_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); } catch (IOException e) { logger.error("Create tables failed", e); @@ -61,7 +81,7 @@ public class HBaseUtil { connection = ConnectionFactory.createConnection(configuration); } } - + public static boolean saveCidTidMapping(String traceId, ChainInfo chainInfo) { Table table = null; @@ -90,7 +110,7 @@ public class HBaseUtil { } - public static ChainRelationship selectCallChainRelationship(String key) throws IOException { + 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)); @@ -167,4 +187,169 @@ 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)); + Get g = new Get(Bytes.toBytes(key)); + Result r = table.get(g); + + if (r.rawCells().length == 0) { + return null; + } + result = new ChainSpecificMinSummary(); + for (Cell cell : r.rawCells()) { + if (cell.getValueArray().length > 0) + result.addNodeSummaryResult(new ChainNodeSpecificMinSummary(Bytes.toString(cell.getValueArray(), + cell.getValueOffset(), cell.getValueLength()))); + } + return result; + } + + public static ChainSpecificHourSummary loadSpecificHourSummary(String key) throws IOException { + ChainSpecificHourSummary result = null; + Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CHAIN_ONE_HOUR_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME)); + Get g = new Get(Bytes.toBytes(key)); + Result r = table.get(g); + + if (r.rawCells().length == 0) { + return null; + } + result = new ChainSpecificHourSummary(); + for (Cell cell : r.rawCells()) { + if (cell.getValueArray().length > 0) + result.addNodeSummaryResult(new ChainNodeSpecificHourSummary(Bytes.toString(cell.getValueArray(), + cell.getValueOffset(), cell.getValueLength()))); + } + return result; + } + + public static ChainSpecificDaySummary loadSpecificDaySummary(String key) throws IOException { + ChainSpecificDaySummary result = null; + Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CHAIN_ONE_DAY_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME)); + Get g = new Get(Bytes.toBytes(key)); + Result r = table.get(g); + + if (r.rawCells().length == 0) { + return null; + } + result = new ChainSpecificDaySummary(); + for (Cell cell : r.rawCells()) { + if (cell.getValueArray().length > 0) + result.addNodeSummaryResult(new ChainNodeSpecificDaySummary(Bytes.toString(cell.getValueArray(), + cell.getValueOffset(), cell.getValueLength()))); + } + return result; + } + + public static ChainSpecificMonthSummary loadSpecificMonthSummary(String key) throws IOException { + ChainSpecificMonthSummary result = null; + Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CHAIN_ONE_MONTH_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME)); + Get g = new Get(Bytes.toBytes(key)); + Result r = table.get(g); + + if (r.rawCells().length == 0) { + return null; + } + result = new ChainSpecificMonthSummary(); + for (Cell cell : r.rawCells()) { + if (cell.getValueArray().length > 0) + result.addNodeSummaryResult(new ChainNodeSpecificMonthSummary(Bytes.toString(cell.getValueArray(), + cell.getValueOffset(), cell.getValueLength()))); + } + return result; + } + + public static void batchSaveSpecificMinSummary(Map minSummary) throws IOException, InterruptedException { + List puts = new ArrayList(); + for (Map.Entry entry : minSummary.entrySet()) { + Put put = new Put(entry.getKey().getBytes()); + entry.getValue().save(put); + puts.add(put); + } + + batchSavePuts(puts, HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME); + } + + public static void batchSaveSpecificDaySummary(Map daySummaryMap) throws IOException, InterruptedException { + List puts = new ArrayList(); + for (Map.Entry entry : daySummaryMap.entrySet()) { + Put put = new Put(entry.getKey().getBytes()); + entry.getValue().save(put); + puts.add(put); + } + + batchSavePuts(puts, HBaseTableMetaData.TABLE_CHAIN_ONE_DAY_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME); + } + + public static void batchSaveSpecificHourSummary(Map hourSummaryMap) throws IOException, InterruptedException { + List puts = new ArrayList(); + for (Map.Entry entry : hourSummaryMap.entrySet()) { + Put put = new Put(entry.getKey().getBytes()); + entry.getValue().save(put); + puts.add(put); + } + + batchSavePuts(puts, HBaseTableMetaData.TABLE_CHAIN_ONE_HOUR_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME); + } + + public static void batchSaveSpecificMonthSummary(Map monthSummaryMap) throws IOException, InterruptedException { + List puts = new ArrayList(); + for (Map.Entry entry : monthSummaryMap.entrySet()) { + Put put = new Put(entry.getKey().getBytes()); + entry.getValue().save(put); + puts.add(put); + } + + batchSavePuts(puts, HBaseTableMetaData.TABLE_CHAIN_ONE_MONTH_SUMMARY_INCLUDE_RELATIONSHIP.TABLE_NAME); + } + + private static void batchSavePuts(List puts, String tableName) throws IOException, InterruptedException { + Object[] resultArrays = new Object[puts.size()]; + Table table = connection.getTable(TableName.valueOf(tableName)); + table.batch(puts, resultArrays); + for (Object result : resultArrays) { + if (result == null) { + logger.error("Failed to save chain specific Summary"); + } + } + } } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Chain2SummaryMapper.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Chain2SummaryMapper.java new file mode 100644 index 000000000..96b9a3636 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Chain2SummaryMapper.java @@ -0,0 +1,41 @@ +package com.ai.cloud.skywalking.analysis.chain2summary; + +import com.ai.cloud.skywalking.analysis.config.ConfigInitializer; +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 java.io.IOException; + +public class Chain2SummaryMapper extends TableMapper { + + private Logger logger = LoggerFactory + .getLogger(Chain2SummaryMapper.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 { + try { + ChainSpecificTimeSummary summary = new ChainSpecificTimeSummary(Bytes.toString(key.get())); + for (Cell cell : value.rawCells()) { + summary.addChainNodeSummaryResult(Bytes.toString(cell.getValueArray(), + cell.getValueOffset(), cell.getValueLength())); + } + context.write(new Text(summary.buildMapperKey()), summary); + } catch (Exception e) { + logger.error("Failed to mapper call chain[" + key.toString() + "]", + e); + } + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Chain2SummaryReducer.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Chain2SummaryReducer.java new file mode 100644 index 000000000..66e176849 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Chain2SummaryReducer.java @@ -0,0 +1,31 @@ +package com.ai.cloud.skywalking.analysis.chain2summary; + +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 java.io.IOException; +import java.util.Iterator; + +public class Chain2SummaryReducer extends Reducer { + + @Override + protected void setup(Context context) throws IOException, InterruptedException { + ConfigInitializer.initialize(); + } + + @Override + protected void reduce(Text key, Iterable values, Context context) throws IOException, InterruptedException { + ChainRelationship4Search chainRelationship = ChainRelationship4Search.load(Bytes.toString(key.getBytes())); + Iterator summaryIterator = values.iterator(); + Summary summary = new Summary(); + while (summaryIterator.hasNext()) { + ChainSpecificTimeSummary timeSummary = summaryIterator.next(); + summary.summary(timeSummary, chainRelationship); + } + + summary.saveToHBase(); + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/ChainRelationship4Search.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/ChainRelationship4Search.java new file mode 100644 index 000000000..06736e117 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/ChainRelationship4Search.java @@ -0,0 +1,32 @@ +package com.ai.cloud.skywalking.analysis.chain2summary; + +import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil; + +import java.io.IOException; +import java.util.HashMap; +import java.util.Map; + +public class ChainRelationship4Search { + + // key: 正常链路ID value: 正常链路ID + // key: 异常链路ID value: 正常链路ID + // key: 未分类链路ID value: 未分类链路ID + private Map chainRelationshipMap = new HashMap(); + + public static ChainRelationship4Search load(String rowkey) throws IOException { + ChainRelationship4Search chainRelationship4Search = HBaseUtil.queryChainRelationship(rowkey); + return chainRelationship4Search; + } + + public void addRelationship(String cid) { + chainRelationshipMap.put(cid, cid); + } + + public void addRelationship(String normalCID, String abnormalCID) { + chainRelationshipMap.put(normalCID, abnormalCID); + } + + public String searchRelationship(String cid) { + return chainRelationshipMap.get(cid); + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/ChainSpecificTimeSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/ChainSpecificTimeSummary.java new file mode 100644 index 000000000..2071ea594 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/ChainSpecificTimeSummary.java @@ -0,0 +1,108 @@ +package com.ai.cloud.skywalking.analysis.chain2summary; + +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummary; +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.io.Writable; + +import java.io.DataInput; +import java.io.DataOutput; +import java.io.IOException; +import java.text.ParseException; +import java.text.SimpleDateFormat; +import java.util.Calendar; +import java.util.Date; +import java.util.HashMap; +import java.util.Map; + +public class ChainSpecificTimeSummary implements Writable { + private String cId; + private String userId; + private String entranceNodeToken; + //key : TraceLevelId + private Map summaryMap; + private long summaryTimestamp; + + public ChainSpecificTimeSummary(String rowKey) throws ParseException { + String[] splitValue = rowKey.split("-"); + this.cId = splitValue[0]; + this.userId = splitValue[1]; + this.summaryTimestamp = new SimpleDateFormat("yyyy/MM/dd HH:mm:ss").parse(splitValue[2]).getTime(); + summaryMap = new HashMap(); + } + + @Override + public void write(DataOutput out) throws IOException { + out.write(new Gson().toJson(this).getBytes()); + } + + @Override + public void readFields(DataInput in) throws IOException { + JsonObject jsonObject = (JsonObject) new JsonParser().parse(in.readLine()); + cId = jsonObject.get("cId").getAsString(); + userId = jsonObject.get("userId").getAsString(); + entranceNodeToken = jsonObject.get("entranceNodeToken").getAsString(); + summaryMap = new Gson().fromJson(jsonObject.get("summaryMap").toString(), + new TypeToken>() { + }.getType()); + } + + public void addChainNodeSummaryResult(String summaryResult) { + ChainNodeSpecificTimeWindowSummary chainNodeSpecificTimeWindowSummary = new Gson(). + fromJson(summaryResult, ChainNodeSpecificTimeWindowSummary.class); + + if ("0".equals(chainNodeSpecificTimeWindowSummary.getTraceLevelId())) { + this.entranceNodeToken = chainNodeSpecificTimeWindowSummary.getNodeToken(); + } + + summaryMap.put(chainNodeSpecificTimeWindowSummary.getTraceLevelId(), chainNodeSpecificTimeWindowSummary); + } + + public String buildMapperKey() { + return userId + ":" + entranceNodeToken; + } + + + public String getcId() { + return cId; + } + + public String getHourKey() { + Calendar calendar = Calendar.getInstance(); + calendar.setTime(new Date(summaryTimestamp)); + return calendar.get(Calendar.YEAR) + "-" + calendar.get(Calendar.MONTH) + calendar.get(Calendar.DAY_OF_MONTH) + + " " + calendar.get(Calendar.HOUR); + } + + public String getDayKey() { + Calendar calendar = Calendar.getInstance(); + calendar.setTime(new Date(summaryTimestamp)); + return calendar.get(Calendar.YEAR) + "-" + calendar.get(Calendar.MONTH) + calendar.get(Calendar.DAY_OF_MONTH); + } + + public String getMonthKey() { + Calendar calendar = Calendar.getInstance(); + calendar.setTime(new Date(summaryTimestamp)); + return calendar.get(Calendar.YEAR) + "-" + calendar.get(Calendar.MONTH); + } + + public String getYearKey() { + Calendar calendar = Calendar.getInstance(); + calendar.setTime(new Date(summaryTimestamp)); + return String.valueOf(calendar.get(Calendar.YEAR)); + } + + public String getUserId() { + return userId; + } + + public Map getSummaryMap() { + return summaryMap; + } + + public long getSummaryTimestamp() { + return summaryTimestamp; + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/ChainSummaryWithRelationship.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/ChainSummaryWithRelationship.java new file mode 100644 index 000000000..45f8023af --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/ChainSummaryWithRelationship.java @@ -0,0 +1,104 @@ +package com.ai.cloud.skywalking.analysis.chain2summary; + +import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil; +import com.ai.cloud.skywalking.analysis.chain2summary.model.ChainSpecificDaySummary; +import com.ai.cloud.skywalking.analysis.chain2summary.model.ChainSpecificHourSummary; +import com.ai.cloud.skywalking.analysis.chain2summary.model.ChainSpecificMinSummary; +import com.ai.cloud.skywalking.analysis.chain2summary.model.ChainSpecificMonthSummary; + +import java.io.IOException; +import java.util.HashMap; +import java.util.Map; + +public class ChainSummaryWithRelationship { + + private String cid; + // key : cid + userId + 小时 + private Map minSummary; + // key : cid + userId + 天 + private Map hourSummary; + // key : cid + userId + 月 + private Map daySummary; + // key : cid + userId + 年 + private Map monthSummary; + + public ChainSummaryWithRelationship(String cid) { + this.cid = cid; + minSummary = new HashMap(); + hourSummary = new HashMap(); + daySummary = new HashMap(); + monthSummary = new HashMap(); + } + + public void saveToHBase() throws IOException, InterruptedException { + HBaseUtil.batchSaveSpecificMinSummary(minSummary); + HBaseUtil.batchSaveSpecificHourSummary(hourSummary); + HBaseUtil.batchSaveSpecificDaySummary(daySummary); + HBaseUtil.batchSaveSpecificMonthSummary(monthSummary); + } + + public void summary(ChainSpecificTimeSummary timeSummary) throws IOException { + loadSummaryIfNecessary(cid,timeSummary); + // + minSummary.get(buildMinSummaryRowKey(cid, timeSummary)).summary(timeSummary); + hourSummary.get(buildHourSummaryRowKey(cid, timeSummary)).summary(timeSummary); + daySummary.get(buildDaySummaryRowKey(cid, timeSummary)).summary(timeSummary); + monthSummary.get(buildMonthSummaryRowKey(cid, timeSummary)).summary(timeSummary); + } + + + private void loadSummaryIfNecessary(String cid, ChainSpecificTimeSummary timeSummary) throws IOException { + loadMinSummaryIfNecessary(cid, timeSummary); + loadHourSummaryIfNecessary(cid, timeSummary); + loadDaySummaryIfNecessary(cid, timeSummary); + loadMonthSummaryIfNecessary(cid, timeSummary); + } + + private void loadMonthSummaryIfNecessary(String cid, ChainSpecificTimeSummary timeSummary) throws IOException { + String month_RowKey = buildMonthSummaryRowKey(cid, timeSummary); + if (!monthSummary.containsKey(month_RowKey)) { + monthSummary.put(month_RowKey, HBaseUtil.loadSpecificMonthSummary(month_RowKey)); + } + } + + private void loadDaySummaryIfNecessary(String cid, ChainSpecificTimeSummary timeSummary) throws IOException { + String day_RowKey = buildDaySummaryRowKey(cid, timeSummary); + if (!daySummary.containsKey(day_RowKey)) { + daySummary.put(day_RowKey, HBaseUtil.loadSpecificDaySummary(day_RowKey)); + } + } + + private void loadHourSummaryIfNecessary(String cid, ChainSpecificTimeSummary timeSummary) throws IOException { + String hour_RowKey = buildHourSummaryRowKey(cid, timeSummary); + if (!hourSummary.containsKey(hour_RowKey)) { + hourSummary.put(hour_RowKey, HBaseUtil.loadSpecificHourSummary(hour_RowKey)); + } + } + + private void loadMinSummaryIfNecessary(String cid, ChainSpecificTimeSummary timeSummary) throws IOException { + String min_RowKey = buildMinSummaryRowKey(cid, timeSummary); + if (!minSummary.containsKey(min_RowKey)) { + minSummary.put(min_RowKey, HBaseUtil.loadSpecificMinSummary(min_RowKey)); + } + } + + // 月统计是以年作为RowKey的 + private static String buildMonthSummaryRowKey(String cid, ChainSpecificTimeSummary timeSummary) { + return cid + "-" + timeSummary.getUserId() + "-" + timeSummary.getYearKey(); + } + + // 天统计是以月作为RowKey的 + private static String buildDaySummaryRowKey(String cid, ChainSpecificTimeSummary timeSummary) { + return cid + "-" + timeSummary.getUserId() + "-" + timeSummary.getMonthKey(); + } + + // 小时统计是以天作为RowKey的 + private static String buildHourSummaryRowKey(String cid, ChainSpecificTimeSummary timeSummary) { + return cid + "-" + timeSummary.getUserId() + "-" + timeSummary.getDayKey(); + } + + // 分钟统计是以小时作为RowKey的 + private static String buildMinSummaryRowKey(String cid, ChainSpecificTimeSummary timeSummary) { + return cid + "-" + timeSummary.getUserId() + "-" + timeSummary.getHourKey(); + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Summary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Summary.java new file mode 100644 index 000000000..8ff90619a --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/Summary.java @@ -0,0 +1,28 @@ +package com.ai.cloud.skywalking.analysis.chain2summary; + +import java.io.IOException; +import java.util.Map; + +public class Summary { + private Map stringChainSummaryWithRelationshipMap; + + public void summary(ChainSpecificTimeSummary timeSummary, ChainRelationship4Search chainRelationship) throws IOException { + String cid = chainRelationship.searchRelationship(timeSummary.getcId()); + if (cid == null || cid.length() == 0) { + cid = timeSummary.getcId(); + } + + if (!stringChainSummaryWithRelationshipMap.containsKey(cid)) { + stringChainSummaryWithRelationshipMap.put(cid, new ChainSummaryWithRelationship(cid)); + } + + ChainSummaryWithRelationship chainSummaryWithRelationship = stringChainSummaryWithRelationshipMap.get(cid); + chainSummaryWithRelationship.summary(timeSummary); + } + + public void saveToHBase() throws IOException, InterruptedException { + for (Map.Entry entry : stringChainSummaryWithRelationshipMap.entrySet()) { + entry.getValue().saveToHBase(); + } + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificDaySummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificDaySummary.java new file mode 100644 index 000000000..990593b08 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificDaySummary.java @@ -0,0 +1,42 @@ +package com.ai.cloud.skywalking.analysis.chain2summary.model; + +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummary; +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummaryValue; +import com.google.gson.Gson; +import com.google.gson.JsonObject; +import com.google.gson.JsonParser; +import com.google.gson.reflect.TypeToken; + +import java.util.Calendar; +import java.util.Date; +import java.util.Map; + +public class ChainNodeSpecificDaySummary { + private String traceLevelId; + // key: 天 + private Map summerValueMap; + + public ChainNodeSpecificDaySummary(String originData) { + JsonObject jsonObject = (JsonObject) new JsonParser().parse(originData); + traceLevelId = jsonObject.get("traceLevelId").getAsString(); + summerValueMap = new Gson().fromJson(jsonObject.get("summerValueMap").toString(), + new TypeToken>() { + }.getType()); + } + + public String getTraceLevelId() { + return traceLevelId; + } + + public void summary(long summaryTimestamp,ChainNodeSpecificTimeWindowSummary value) { + for (Map.Entry entry : value.getSummerValueMap().entrySet()) { + summerValueMap.get(generateSummaryValueMapKey(summaryTimestamp)).accumulate(entry.getValue()); + } + } + + private String generateSummaryValueMapKey(long timeStamp) { + Calendar calendar = Calendar.getInstance(); + calendar.setTime(new Date(timeStamp)); + return String.valueOf(calendar.get(Calendar.DAY_OF_MONTH)); + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificHourSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificHourSummary.java new file mode 100644 index 000000000..2b288fcc9 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificHourSummary.java @@ -0,0 +1,42 @@ +package com.ai.cloud.skywalking.analysis.chain2summary.model; + +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummary; +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummaryValue; +import com.google.gson.Gson; +import com.google.gson.JsonObject; +import com.google.gson.JsonParser; +import com.google.gson.reflect.TypeToken; + +import java.util.Calendar; +import java.util.Date; +import java.util.Map; + +public class ChainNodeSpecificHourSummary { + private String traceLevelId; + // key : 小时 + private Map summerValueMap; + + public ChainNodeSpecificHourSummary(String originData) { + JsonObject jsonObject = (JsonObject) new JsonParser().parse(originData); + traceLevelId = jsonObject.get("traceLevelId").getAsString(); + summerValueMap = new Gson().fromJson(jsonObject.get("summerValueMap").toString(), + new TypeToken>() { + }.getType()); + } + + public String getTraceLevelId() { + return traceLevelId; + } + + public void summary(long summaryTimestamp, ChainNodeSpecificTimeWindowSummary value) { + for (Map.Entry entry : value.getSummerValueMap().entrySet()) { + summerValueMap.get(generateSummaryValueMapKey(summaryTimestamp)).accumulate(entry.getValue()); + } + } + + private String generateSummaryValueMapKey(long timeStamp) { + Calendar calendar = Calendar.getInstance(); + calendar.setTime(new Date(timeStamp)); + return String.valueOf(calendar.get(Calendar.HOUR)); + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificMinSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificMinSummary.java new file mode 100644 index 000000000..268f62cc7 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificMinSummary.java @@ -0,0 +1,36 @@ +package com.ai.cloud.skywalking.analysis.chain2summary.model; + +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummary; +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummaryValue; +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 ChainNodeSpecificMinSummary { + + private String traceLevelId; + // key: 分钟 value: 统计结果 + private Map summerValueMap; + + + public ChainNodeSpecificMinSummary(String originData) { + JsonObject jsonObject = (JsonObject) new JsonParser().parse(originData); + traceLevelId = jsonObject.get("traceLevelId").getAsString(); + summerValueMap = new Gson().fromJson(jsonObject.get("summerValueMap").toString(), + new TypeToken>() { + }.getType()); + } + + public String getTraceLevelId() { + return traceLevelId; + } + + public void summary(ChainNodeSpecificTimeWindowSummary value) { + for (Map.Entry entry :value.getSummerValueMap().entrySet()){ + summerValueMap.get(entry.getKey()).accumulate(entry.getValue()); + } + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificMonthSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificMonthSummary.java new file mode 100644 index 000000000..44e5dc3fd --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainNodeSpecificMonthSummary.java @@ -0,0 +1,42 @@ +package com.ai.cloud.skywalking.analysis.chain2summary.model; + +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummary; +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummaryValue; +import com.google.gson.Gson; +import com.google.gson.JsonObject; +import com.google.gson.JsonParser; +import com.google.gson.reflect.TypeToken; + +import java.util.Calendar; +import java.util.Date; +import java.util.Map; + +public class ChainNodeSpecificMonthSummary { + + private Map summerValueMap; + private String traceLevelId; + + public ChainNodeSpecificMonthSummary(String originData) { + JsonObject jsonObject = (JsonObject) new JsonParser().parse(originData); + traceLevelId = jsonObject.get("traceLevelId").getAsString(); + summerValueMap = new Gson().fromJson(jsonObject.get("summerValueMap").toString(), + new TypeToken>() { + }.getType()); + } + + public void summary(long summaryTimestamp, ChainNodeSpecificTimeWindowSummary value) { + for (Map.Entry entry : value.getSummerValueMap().entrySet()) { + summerValueMap.get(generateSummaryValueMapKey(summaryTimestamp)).accumulate(entry.getValue()); + } + } + + public String getTraceLevelId() { + return traceLevelId; + } + + private String generateSummaryValueMapKey(long timeStamp) { + Calendar calendar = Calendar.getInstance(); + calendar.setTime(new Date(timeStamp)); + return String.valueOf(calendar.get(Calendar.MONTH)); + } +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainSpecificDaySummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainSpecificDaySummary.java new file mode 100644 index 000000000..6730f274f --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainSpecificDaySummary.java @@ -0,0 +1,37 @@ +package com.ai.cloud.skywalking.analysis.chain2summary.model; + +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummary; +import com.ai.cloud.skywalking.analysis.chain2summary.ChainSpecificTimeSummary; +import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData; +import org.apache.hadoop.hbase.client.Put; + +import java.util.HashMap; +import java.util.Map; + +public class ChainSpecificDaySummary { + private Map chainNodeSpecificHourSummaryMap; + + public ChainSpecificDaySummary() { + chainNodeSpecificHourSummaryMap = new HashMap(); + } + + public void addNodeSummaryResult(ChainNodeSpecificDaySummary chainNodeSpecificHourSummary) { + chainNodeSpecificHourSummaryMap.put(chainNodeSpecificHourSummary.getTraceLevelId(), chainNodeSpecificHourSummary); + } + + public void summary(ChainSpecificTimeSummary timeSummary) { + Map chainNodeSpecificTimeWindowSummaryMap = timeSummary.getSummaryMap(); + + for (Map.Entry entry : chainNodeSpecificTimeWindowSummaryMap.entrySet()) { + chainNodeSpecificHourSummaryMap.get(entry.getKey()).summary(timeSummary.getSummaryTimestamp(), entry.getValue()); + } + } + + public void save(Put put) { + for (Map.Entry entry : chainNodeSpecificHourSummaryMap.entrySet()) { + put.addColumn(HBaseTableMetaData.TABLE_CHAIN_ONE_DAY_SUMMARY_INCLUDE_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/chain2summary/model/ChainSpecificHourSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainSpecificHourSummary.java new file mode 100644 index 000000000..749dfd63d --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainSpecificHourSummary.java @@ -0,0 +1,38 @@ +package com.ai.cloud.skywalking.analysis.chain2summary.model; + +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummary; +import com.ai.cloud.skywalking.analysis.chain2summary.ChainSpecificTimeSummary; +import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData; +import org.apache.hadoop.hbase.client.Put; + +import java.util.HashMap; +import java.util.Map; + +public class ChainSpecificHourSummary { + // key : TraceLevelId + private Map chainNodeSpecificHourSummaryMap; + + public ChainSpecificHourSummary() { + chainNodeSpecificHourSummaryMap = new HashMap(); + } + + public void addNodeSummaryResult(ChainNodeSpecificHourSummary chainNodeSpecificHourSummary) { + chainNodeSpecificHourSummaryMap.put(chainNodeSpecificHourSummary.getTraceLevelId(), chainNodeSpecificHourSummary); + } + + public void summary(ChainSpecificTimeSummary timeSummary) { + Map chainNodeSpecificTimeWindowSummaryMap = timeSummary.getSummaryMap(); + + for (Map.Entry entry : chainNodeSpecificTimeWindowSummaryMap.entrySet()){ + chainNodeSpecificHourSummaryMap.get(entry.getKey()).summary(timeSummary.getSummaryTimestamp(),entry.getValue()); + } + } + + public void save(Put put) { + for (Map.Entry entry : chainNodeSpecificHourSummaryMap.entrySet()) { + put.addColumn(HBaseTableMetaData.TABLE_CHAIN_ONE_HOUR_SUMMARY_INCLUDE_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/chain2summary/model/ChainSpecificMinSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainSpecificMinSummary.java new file mode 100644 index 000000000..efef5b0b3 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainSpecificMinSummary.java @@ -0,0 +1,39 @@ +package com.ai.cloud.skywalking.analysis.chain2summary.model; + +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummary; +import com.ai.cloud.skywalking.analysis.chain2summary.ChainSpecificTimeSummary; +import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData; +import org.apache.hadoop.hbase.client.Put; + +import java.util.HashMap; +import java.util.Map; + +public class ChainSpecificMinSummary { + + // Key: TraceLevelId + private Map chainNodeSpecificMinSummaryMap; + + public ChainSpecificMinSummary() { + this.chainNodeSpecificMinSummaryMap = new HashMap(); + } + + public void addNodeSummaryResult(ChainNodeSpecificMinSummary chainNodeSpecificMinSummary) { + chainNodeSpecificMinSummaryMap.put(chainNodeSpecificMinSummary.getTraceLevelId(), chainNodeSpecificMinSummary); + } + + public void summary(ChainSpecificTimeSummary timeSummary) { + Map chainNodeSpecificTimeWindowSummaryMap = timeSummary.getSummaryMap(); + + for (Map.Entry entry : chainNodeSpecificTimeWindowSummaryMap.entrySet()) { + chainNodeSpecificMinSummaryMap.get(entry.getKey()).summary(entry.getValue()); + } + } + + public void save(Put put) { + for (Map.Entry entry : chainNodeSpecificMinSummaryMap.entrySet()) { + put.addColumn(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_INCLUDE_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/chain2summary/model/ChainSpecificMonthSummary.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainSpecificMonthSummary.java new file mode 100644 index 000000000..ac5504bcf --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/chain2summary/model/ChainSpecificMonthSummary.java @@ -0,0 +1,34 @@ +package com.ai.cloud.skywalking.analysis.chain2summary.model; + +import com.ai.cloud.skywalking.analysis.categorize2chain.ChainNodeSpecificTimeWindowSummary; +import com.ai.cloud.skywalking.analysis.chain2summary.ChainSpecificTimeSummary; +import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData; +import org.apache.hadoop.hbase.client.Put; + +import java.util.Map; + +public class ChainSpecificMonthSummary { + + // Key: TraceLevelId + private Map chainNodeSpecificMinSummaryMap; + + public void addNodeSummaryResult(ChainNodeSpecificMonthSummary chainNodeSpecificDaySummary) { + chainNodeSpecificMinSummaryMap.put(chainNodeSpecificDaySummary.getTraceLevelId(), chainNodeSpecificDaySummary); + } + + public void summary(ChainSpecificTimeSummary timeSummary) { + Map chainNodeSpecificTimeWindowSummaryMap = timeSummary.getSummaryMap(); + + for (Map.Entry entry : chainNodeSpecificTimeWindowSummaryMap.entrySet()) { + chainNodeSpecificMinSummaryMap.get(entry.getKey()).summary(timeSummary.getSummaryTimestamp(), entry.getValue()); + } + } + + public void save(Put put) { + for (Map.Entry entry : chainNodeSpecificMinSummaryMap.entrySet()) { + put.addColumn(HBaseTableMetaData.TABLE_CHAIN_ONE_MONTH_SUMMARY_INCLUDE_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/config/HBaseTableMetaData.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/HBaseTableMetaData.java index d212e352d..a43cfc634 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 @@ -1,65 +1,104 @@ package com.ai.cloud.skywalking.analysis.config; public class HBaseTableMetaData { - /** - * 调用链明细表,前端收集程序入库数据 - * - * @author wusheng - * - */ - public final static class TABLE_CALL_CHAIN { - public static final String TABLE_NAME = "sw-call-chain"; - } + /** + * 调用链明细表,前端收集程序入库数据 + * + * @author wusheng + */ + public final static class TABLE_CALL_CHAIN { + public static final String TABLE_NAME = "sw-call-chain"; + } - /** - * HBase 表:用于存放CID,TID的映射关系表 - * - * @author wusheng - * - */ - public final static class TABLE_CID_TID_MAPPING { - public static final String TABLE_NAME = "sw-cid-tid-mapping"; + /** + * HBase 表:用于存放CID,TID的映射关系表 + * + * @author wusheng + */ + public final static class TABLE_CID_TID_MAPPING { + public static final String TABLE_NAME = "sw-cid-tid-mapping"; - public static final String COLUMN_FAMILY_NAME = "trace_info"; - - public static final String CID_COLUMN_NAME = "cid"; - } + public static final String COLUMN_FAMILY_NAME = "trace_info"; - /** - * CID明细信息表 - * - * @author wusheng - * - */ - public final static class TABLE_CHAIN_DETAIL { - public static final String TABLE_NAME = "sw-chain-detail"; + public static final String CID_COLUMN_NAME = "cid"; + } - public static final String COLUMN_FAMILY_NAME = "chain_detail"; - } - - /** - * CID间关系表,记录异常CID和正常CID间的归属关系 - * - * @author wusheng - * - */ - public final static class TABLE_CALL_CHAIN_RELATIONSHIP { - 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"; - } - - /** - * 用于存放每个CID在一分钟内的汇总,汇总结果不包含关系汇总 - * - * @author wusheng - * - */ - public final static class TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP{ - public static final String TABLE_NAME = "sw-chain-1min-summary-ex-rela"; - - public static final String COLUMN_FAMILY_NAME = "chain_summary"; - } + /** + * CID明细信息表 + * + * @author wusheng + */ + public final static class TABLE_CHAIN_DETAIL { + public static final String TABLE_NAME = "sw-chain-detail"; + + public static final String COLUMN_FAMILY_NAME = "chain_detail"; + } + + /** + * CID间关系表,记录异常CID和正常CID间的归属关系 + * + * @author wusheng + */ + public final static class TABLE_CALL_CHAIN_RELATIONSHIP { + 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"; + } + + /** + * 用于存放每个CID在一分钟内的汇总,汇总结果不包含关系汇总 + * + * @author wusheng + */ + public final static class TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP { + public static final String TABLE_NAME = "sw-chain-1min-summary-ex-rela"; + + public static final String COLUMN_FAMILY_NAME = "chain_summary"; + } + + /** + * 用于存放每个CID在一分钟内的汇总,汇总结果不包含关系汇总 + * + * @author wusheng + */ + public final static class TABLE_CHAIN_ONE_MINUTE_SUMMARY_INCLUDE_RELATIONSHIP { + public static final String TABLE_NAME = "sw-chain-1min-summary-ic-rela"; + + public static final String COLUMN_FAMILY_NAME = "chain_summary"; + } + + /** + * 用于存放每个CID在一小时内的汇总,汇总结果不包含关系汇总 + * + * @author wusheng + */ + public final static class TABLE_CHAIN_ONE_HOUR_SUMMARY_INCLUDE_RELATIONSHIP { + public static final String TABLE_NAME = "sw-chain-1hour-summary-ic-rela"; + + public static final String COLUMN_FAMILY_NAME = "chain_summary"; + } + + /** + * 用于存放每个CID在一天内的汇总,汇总结果不包含关系汇总 + * + * @author wusheng + */ + public final static class TABLE_CHAIN_ONE_DAY_SUMMARY_INCLUDE_RELATIONSHIP { + public static final String TABLE_NAME = "sw-chain-1day-summary-ic-rela"; + + public static final String COLUMN_FAMILY_NAME = "chain_summary"; + } + + /** + * 用于存放每个CID在一月内的汇总,汇总结果不包含关系汇总 + * + * @author wusheng + */ + public final static class TABLE_CHAIN_ONE_MONTH_SUMMARY_INCLUDE_RELATIONSHIP { + public static final String TABLE_NAME = "sw-chain-1mon-summary-ic-rela"; + + public static final String COLUMN_FAMILY_NAME = "chain_summary"; + } }