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