From ab08ff1226d8fb1000df959b73f2c32d33c61031 Mon Sep 17 00:00:00 2001 From: zhangxin10 Date: Tue, 19 Jan 2016 10:16:07 +0800 Subject: [PATCH] =?UTF-8?q?1.=20=E4=BF=AE=E6=94=B9NodeToken=E7=9A=84?= =?UTF-8?q?=E8=A7=84=E5=88=99=EF=BC=8C=E5=A2=9E=E5=8A=A0=E5=89=8D=E7=BC=80?= =?UTF-8?q?CID=202.=20=E6=8F=90=E4=BA=A4=E9=83=A8=E5=88=86Reduce=E7=9A=84?= =?UTF-8?q?=E4=BB=A3=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- skywalking-analysis/skywalking-analysis.iml | 117 ++++++++++++++++++ .../analysis/AnalysisServerDriver.java | 1 + .../skywalking/analysis/config/Config.java | 21 ++-- .../analysis/dao/CallChainInfoDao.java | 38 ++++++ .../skywalking/analysis/model/ChainInfo.java | 56 +++++---- .../skywalking/analysis/model/ChainNode.java | 8 ++ .../analysis/reduce/ChainInfoReduce.java | 9 +- .../skywalking/analysis/util/HBaseUtil.java | 20 ++- .../analysis/util/TokenGenerator.java | 2 +- .../src/main/resources/config.properties | 11 +- 10 files changed, 233 insertions(+), 50 deletions(-) create mode 100644 skywalking-analysis/skywalking-analysis.iml create mode 100644 skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/dao/CallChainInfoDao.java diff --git a/skywalking-analysis/skywalking-analysis.iml b/skywalking-analysis/skywalking-analysis.iml new file mode 100644 index 000000000..aa8205c55 --- /dev/null +++ b/skywalking-analysis/skywalking-analysis.iml @@ -0,0 +1,117 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/AnalysisServerDriver.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/AnalysisServerDriver.java index d32049fa7..ab1a2d99d 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/AnalysisServerDriver.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/AnalysisServerDriver.java @@ -43,6 +43,7 @@ public class AnalysisServerDriver extends Configured implements Tool { TableMapReduceUtil.initTableMapperJob(Config.HBase.CALL_CHAIN_TABLE_NAME, scan, CallChainMapper.class, String.class, ChainInfo.class, job); + //TableMapReduceUtil.initTableReducerJob("sw-call-chain-model", CallChainReducer.class, job); return job.waitForCompletion(true) ? 0 : 1; } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Config.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Config.java index 2c3e81434..8f9d4fef2 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Config.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Config.java @@ -13,24 +13,23 @@ public class Config { public static String TRACE_INFO_TABLE_NAME; + public static String TABLE_CALL_CHAIN_RELATIONSHIP = "sw_chain_relationship"; } public static class TraceInfo { - - public static String PARENT_LEVEL_ID = "parentLevelId"; - - public static String LEVEL_ID = "levelId"; - - public static String BUSINESS_KEY = "businessKey"; - - public static String COST = "cost"; - public static String TRACE_INFO_COLUMN_CID = "cid"; + } - public static String STATUS = "status"; + public static class MySql { - public static String USER_ID = "UId"; + public static String url; + + public static String userName; + + public static String password; + + public static String driverClass; } public static class Filter { diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/dao/CallChainInfoDao.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/dao/CallChainInfoDao.java new file mode 100644 index 000000000..93ebc0759 --- /dev/null +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/dao/CallChainInfoDao.java @@ -0,0 +1,38 @@ +package com.ai.cloud.skywalking.analysis.dao; + +import com.ai.cloud.skywalking.analysis.config.Config; +import com.ai.cloud.skywalking.analysis.model.ChainInfo; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.sql.*; + +public class CallChainInfoDao { + private static Logger logger = LoggerFactory.getLogger(CallChainInfoDao.class.getName()); + + private static Connection connection; + + static { + try { + Class.forName(Config.MySql.driverClass); + connection = DriverManager.getConnection(Config.MySql.url, Config.MySql.userName, Config.MySql.password); + } catch (ClassNotFoundException e) { + logger.error("Failed to find jdbc driver class[" + Config.MySql.driverClass + "]", e); + System.exit(-1); + } catch (SQLException e) { + logger.error("Failed to connection database.", e); + System.exit(-1); + } + } + + public static boolean existCallChainInfo(String cid, String userId) throws SQLException { + final String sql = "SELECT COUNT(cid) as TOTAL_SIZE FROM sw_chain_info WHERE cid = ? AND uid = ?"; + PreparedStatement ps = connection.prepareStatement(sql); + ps.setString(1, cid); + ps.setString(2, userId); + ResultSet resultSet = ps.executeQuery(); + resultSet.next(); + return resultSet.getInt("TOTAL_SIZE") > 0 ? true : false; + } + +} diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/model/ChainInfo.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/model/ChainInfo.java index 18cefdfa0..4a4f1d1bd 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/model/ChainInfo.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/model/ChainInfo.java @@ -1,8 +1,7 @@ package com.ai.cloud.skywalking.analysis.model; -import org.apache.hadoop.io.Writable; - import com.ai.cloud.skywalking.analysis.util.TokenGenerator; +import org.apache.hadoop.io.Writable; import java.io.DataInput; import java.io.DataOutput; @@ -25,6 +24,7 @@ public class ChainInfo implements Writable { public void write(DataOutput out) throws IOException { out.write(chainToken.getBytes()); out.writeChar(chainStatus.getValue()); + out.write(userId.getBytes()); out.writeInt(nodes.size()); for (ChainNode chainNode : nodes) { @@ -34,6 +34,7 @@ public class ChainInfo implements Writable { out.writeLong(chainNode.getCost()); out.write(chainNode.getParentLevelId().getBytes()); out.writeInt(chainNode.getLevelId()); + out.write(chainNode.getBusinessKey().getBytes()); } } @@ -41,6 +42,7 @@ public class ChainInfo implements Writable { public void readFields(DataInput in) throws IOException { chainToken = in.readLine(); chainStatus = ChainStatus.convert(in.readChar()); + userId = in.readLine(); int nodeSize = in.readInt(); this.nodes = new ArrayList(); @@ -52,6 +54,8 @@ public class ChainInfo implements Writable { chainNode.setCost(in.readLong()); chainNode.setParentLevelId(in.readLine()); chainNode.setLevelId(in.readInt()); + chainNode.setBusinessKey(in.readLine()); + nodes.add(chainNode); } } @@ -59,20 +63,20 @@ public class ChainInfo implements Writable { public String getChainToken() { return chainToken; } - - public String getEntranceNodeToken(){ - if(firstChainNode == null){ - return ""; - }else{ - return firstChainNode.getNodeToken(); - } + + public String getEntranceNodeToken() { + if (firstChainNode == null) { + return ""; + } else { + return firstChainNode.getNodeToken(); + } } public void generateChainToken() { - StringBuilder chainTokenDesc = new StringBuilder(); - for(ChainNode node: nodes){ - chainTokenDesc.append(node.getParentLevelId() + "." + node.getLevelId() + "-" + node.getNodeToken() + ";"); - } + StringBuilder chainTokenDesc = new StringBuilder(); + for (ChainNode node : nodes) { + chainTokenDesc.append(node.getParentLevelId() + "." + node.getLevelId() + "-" + node.getNodeToken() + ";"); + } this.chainToken = TokenGenerator.generate(chainTokenDesc.toString()); } @@ -83,19 +87,19 @@ public class ChainInfo implements Writable { public void setChainStatus(ChainStatus chainStatus) { this.chainStatus = chainStatus; } - - public void addNodes(ChainNode chainNode){ - this.nodes.add(0, chainNode); - if (chainNode.getStatus() == ChainNode.NodeStatus.ABNORMAL) { - chainStatus = ChainStatus.ABNORMAL; - } - if(userId == null){ - userId = chainNode.getUserId(); - } - if ((chainNode.getParentLevelId() == null || chainNode.getParentLevelId().length() == 0) + + public void addNodes(ChainNode chainNode) { + this.nodes.add(0, chainNode); + if (chainNode.getStatus() == ChainNode.NodeStatus.ABNORMAL) { + chainStatus = ChainStatus.ABNORMAL; + } + if (userId == null) { + userId = chainNode.getUserId(); + } + if ((chainNode.getParentLevelId() == null || chainNode.getParentLevelId().length() == 0) && chainNode.getLevelId() == 0) { - firstChainNode = chainNode; - } + firstChainNode = chainNode; + } } public List getNodes() { @@ -109,7 +113,7 @@ public class ChainInfo implements Writable { public void setUserId(String userId) { this.userId = userId; } - + public enum ChainStatus { NORMAL('N'), ABNORMAL('A'); private char value; diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/model/ChainNode.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/model/ChainNode.java index e18083e60..febbb721f 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/model/ChainNode.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/model/ChainNode.java @@ -79,6 +79,14 @@ public class ChainNode { this.businessKey = businessKey; } + public String getCallType() { + return callType; + } + + public String getBusinessKey() { + return businessKey; + } + public enum NodeStatus { NORMAL('N'), ABNORMAL('A'); private char value; diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/reduce/ChainInfoReduce.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/reduce/ChainInfoReduce.java index d47aa1283..847966cc2 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/reduce/ChainInfoReduce.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/reduce/ChainInfoReduce.java @@ -7,9 +7,14 @@ import org.apache.hadoop.io.Text; import java.io.IOException; -public class ChainInfoReduce extends TableReducer { +public class ChainInfoReduce extends TableReducer { @Override protected void reduce(Text key, Iterable values, Context context) throws IOException, InterruptedException { - super.reduce(key, values, context); + String[] keyArray = key.toString().split(":"); + String userId = keyArray[0]; + String firstNode = keyArray[1]; + + // + } } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/util/HBaseUtil.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/util/HBaseUtil.java index 21201b3ad..4944d911a 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/util/HBaseUtil.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/util/HBaseUtil.java @@ -2,17 +2,17 @@ package com.ai.cloud.skywalking.analysis.util; import com.ai.cloud.skywalking.analysis.config.Config; import com.ai.cloud.skywalking.analysis.model.ChainInfo; +import com.ai.cloud.skywalking.protocol.Span; import org.apache.hadoop.conf.Configuration; -import org.apache.hadoop.hbase.HBaseConfiguration; -import org.apache.hadoop.hbase.HColumnDescriptor; -import org.apache.hadoop.hbase.HTableDescriptor; -import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.*; import org.apache.hadoop.hbase.client.*; import org.apache.hadoop.hbase.util.Bytes; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.IOException; +import java.util.ArrayList; +import java.util.List; public class HBaseUtil { private static Logger logger = LoggerFactory.getLogger(HBaseUtil.class.getName()); @@ -75,4 +75,16 @@ public class HBaseUtil { } } + + public static void selectById(String id) throws IOException { + List entries = new ArrayList(); + Table table = connection.getTable(TableName.valueOf(Config.HBase.TABLE_CALL_CHAIN_RELATIONSHIP)); + Get g = new Get(Bytes.toBytes(id)); + Result r = table.get(g); + for (Cell cell : r.rawCells()) { + if (cell.getValueArray().length > 0) + entries.add(new Span(Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength()))); + } + } + } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/util/TokenGenerator.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/util/TokenGenerator.java index 73c3200e1..321baa8db 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/util/TokenGenerator.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/util/TokenGenerator.java @@ -31,6 +31,6 @@ public class TokenGenerator { System.exit(-1); } } - return result.toString().toUpperCase(); + return "CID:" + result.toString().toUpperCase(); } } diff --git a/skywalking-analysis/src/main/resources/config.properties b/skywalking-analysis/src/main/resources/config.properties index 57d720597..a24df42c4 100644 --- a/skywalking-analysis/src/main/resources/config.properties +++ b/skywalking-analysis/src/main/resources/config.properties @@ -6,10 +6,9 @@ hbase.zk_client_port=29181 hbase.trace_info_table_name=sw-trace-info hbase.trace_info_column_family=trace_info -traceinfo.parent_level_id=parentLevelId -traceinfo.level_id=levelId -traceinfo.business_key=businessKey -traceinfo.cost=cost traceinfo.trace_info_column_cid=cid -traceinfo.status=status -traceinfo.user_id=UId \ No newline at end of file + +mysql.url= +mysql.username= +mysql.password= +mysql.driverclass= \ No newline at end of file