diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Constants.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Constants.java index a8b3f09e5..63304886c 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Constants.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Constants.java @@ -1,5 +1,6 @@ package com.ai.cloud.skywalking.analysis.config; public class Constants { - public static String UNCATEGORIZED_QUALIFIER_NAME = "UNCATEGORIZED_CALL_CHAIN"; + public static final String EXCEPTIONMAPPER = "-1:"; + public static final String UNCATEGORIZED_QUALIFIER_NAME = "UNCATEGORIZED_CALL_CHAIN"; } 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 index 390faf4ea..14d292aad 100644 --- 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 @@ -41,7 +41,9 @@ public class CallChainInfoDao { } int[] result = preparedStatement.executeBatch(); for (int i : result) { - //TODO + if (i != 1) { + logger.error("Failed to save chain detail [" + chainDetail.getChainToken() + "]"); + } } preparedStatement.close(); connection.commit(); @@ -56,7 +58,9 @@ public class CallChainInfoDao { } int[] result = preparedStatement.executeBatch(); for (int i : result) { - //TODO + if (i != 1) { + logger.error("Failed to update chain detail"); + } } preparedStatement.close(); connection.commit(); diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapper.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapper.java index 705248a20..9187fcdfc 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapper.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapper.java @@ -25,11 +25,19 @@ public class CallChainMapper extends TableMapper { protected void map(ImmutableBytesWritable key, Result value, Context context) throws IOException, InterruptedException { List spanList = new ArrayList(); - for (Cell cell : value.rawCells()) { - Span span = new Span(Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength())); - spanList.add(span); + ChainInfo chainInfo = null; + try { + for (Cell cell : value.rawCells()) { + Span span = new Span(Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength())); + spanList.add(span); + } + + chainInfo = spanToChainInfo(key.toString(), spanList); + } catch (Exception e) { + logger.error("Failed to mapper call chain[" + key.toString() + "]", e); + chainInfo = new ChainInfo("-1"); } - ChainInfo chainInfo = spanToChainInfo(key.toString(), spanList); + context.write(new Text(chainInfo.getUserId() + ":" + chainInfo.getEntranceNodeToken()), chainInfo); } 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 26ab6cec3..dc88ed8cf 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 @@ -17,6 +17,11 @@ public class ChainInfo implements Writable { private ChainNode firstChainNode; private long startDate; + public ChainInfo(String userId) { + super(); + this.userId = userId; + } + public ChainInfo() { this.nodes = new ArrayList(); } @@ -78,13 +83,17 @@ public class ChainInfo implements Writable { for (ChainNode node : nodes) { chainTokenDesc.append(node.getParentLevelId() + "." + node.getLevelId() + "-" + node.getNodeToken() + ";"); } - this.chainToken = TokenGenerator.generate(chainTokenDesc.toString()) + "-" + userId; + this.chainToken = TokenGenerator.generate(chainTokenDesc.toString()) + "-"; } public ChainStatus getChainStatus() { return chainStatus; } + public void setChainStatus(ChainStatus chainStatus) { + this.chainStatus = chainStatus; + } + public void addNodes(ChainNode chainNode) { this.nodes.add(0, chainNode); if (chainNode.getStatus() == ChainNode.NodeStatus.ABNORMAL) { @@ -139,6 +148,10 @@ public class ChainInfo implements Writable { public void setChainToken(String chainToken) { this.chainToken = chainToken; } + + public long getStartDate() { + return startDate; + } } 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 10404dffd..871d104b1 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 @@ -1,21 +1,31 @@ package com.ai.cloud.skywalking.analysis.reduce; +import com.ai.cloud.skywalking.analysis.config.Constants; import com.ai.cloud.skywalking.analysis.model.ChainInfo; import com.ai.cloud.skywalking.analysis.util.HBaseUtil; import org.apache.hadoop.hbase.client.Put; import org.apache.hadoop.hbase.mapreduce.TableReducer; import org.apache.hadoop.io.Text; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.io.IOException; import java.util.Iterator; public class ChainInfoReduce extends TableReducer { + private static Logger logger = LoggerFactory.getLogger(ChainInfoReduce.class.getName()); + @Override protected void reduce(Text key, Iterable values, Context context) throws IOException, InterruptedException { reduceAction(key.toString(), values.iterator()); } public static void reduceAction(String key, Iterator chainInfoIterator) throws IOException, InterruptedException { + if (Constants.EXCEPTIONMAPPER.equals(key)) { + logger.info("Skip Exception Mapper....."); + return; + } + try { ChainRelate chainRelate = HBaseUtil.selectCallChainRelationship(key.toString()); Summary summary = new Summary(); @@ -32,7 +42,7 @@ public class ChainInfoReduce extends TableReducer { chainRelate.save(); summary.save(); } catch (Exception e) { - e.printStackTrace(); + logger.error("Failed to reduce key[" + key + "]", e); } } } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/reduce/ChainRelate.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/reduce/ChainRelate.java index 139fe29b1..1dcfe9105 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/reduce/ChainRelate.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/reduce/ChainRelate.java @@ -6,12 +6,16 @@ import com.ai.cloud.skywalking.analysis.model.ChainInfo; import com.ai.cloud.skywalking.analysis.util.HBaseUtil; import com.google.gson.Gson; import org.apache.hadoop.hbase.client.Put; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.io.IOException; import java.sql.SQLException; import java.util.*; public class ChainRelate { + private static Logger logger = LoggerFactory.getLogger(ChainRelate.class.getName()); + private String key; private Map categorizedChainInfoMap = new HashMap(); private List uncategorizeChainInfoList = new ArrayList(); @@ -41,12 +45,17 @@ public class ChainRelate { isContained = true; } else if (entry.getValue().isContained(child)) { entry.getValue().add(child); + chainDetailMap.put(child.getChainToken(), new ChainDetail(child.getChainInfo())); isContained = true; } } if (!isContained) { uncategorizeChainInfoList.add(child); + + if (!uncategorizeChainInfoList.contains(child)) { + chainDetailMap.put(child.getChainToken(), new ChainDetail(child.getChainInfo())); + } } } @@ -71,12 +80,12 @@ public class ChainRelate { } } - public void save() throws SQLException { + public void save() throws SQLException, IOException, InterruptedException { saveChainRelationShip(); saveChainDetail(); } - private void saveChainDetail() throws SQLException { + private void saveChainDetail() throws SQLException, IOException, InterruptedException { List puts = new ArrayList(); for (Map.Entry entry : chainDetailMap.entrySet()) { Put put1 = new Put(entry.getKey().getBytes()); @@ -88,13 +97,15 @@ public class ChainRelate { try { HBaseUtil.saveChainDetails(puts); } catch (IOException e) { - e.printStackTrace(); + logger.error("Faild to save chain detail to hbase.", e); + throw e; } catch (InterruptedException e) { - e.printStackTrace(); + logger.error("Faild to save chain detail to hbase.", e); + throw e; } } - private void saveChainRelationShip() { + private void saveChainRelationShip() throws IOException { Put put = new Put(getKey().getBytes()); put.addColumn(Config.HBase.CHAIN_RELATIONSHIP_COLUMN_FAMILY.getBytes(), Constants.UNCATEGORIZED_QUALIFIER_NAME.getBytes() @@ -108,8 +119,8 @@ public class ChainRelate { try { HBaseUtil.saveChainRelate(put); } catch (IOException e) { - //TODO - e.printStackTrace(); + logger.error("Faild to save chain relationship to hbase.", e); + throw e; } } 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 ffb7be20f..133018ebb 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,7 +2,7 @@ package com.ai.cloud.skywalking.analysis.util; import com.ai.cloud.skywalking.analysis.config.Config; import com.ai.cloud.skywalking.analysis.config.Constants; -import com.ai.cloud.skywalking.analysis.model.*; +import com.ai.cloud.skywalking.analysis.model.ChainInfo; import com.ai.cloud.skywalking.analysis.reduce.*; import com.google.gson.Gson; import com.google.gson.reflect.TypeToken; @@ -150,8 +150,7 @@ public class HBaseUtil { table.batch(puts, resultArrays); for (Object result : resultArrays) { if (result == null) { - //TODO - + logger.error("Failed to save chain specificTimeWindows Summary."); } } } @@ -163,8 +162,7 @@ public class HBaseUtil { table.batch(puts, resultArrays); for (Object result : resultArrays) { if (result == null) { - //TODO - + logger.error("Failed to save chain specificTimeWindows Summary."); } } }