提交最新的maper-reduce代码,统计部分未完成

This commit is contained in:
ascrutae 2016-03-06 18:46:15 +08:00
parent 235ec75c84
commit f40fd0302d
41 changed files with 446 additions and 1823 deletions

View File

@ -1,117 +0,0 @@
<?xml version="1.0" encoding="UTF-8"?>
<module org.jetbrains.idea.maven.project.MavenProjectsManager.isMavenModule="true" type="JAVA_MODULE" version="4">
<component name="NewModuleRootManager" LANGUAGE_LEVEL="JDK_1_6" inherit-compiler-output="false">
<output url="file://$MODULE_DIR$/target/classes" />
<output-test url="file://$MODULE_DIR$/target/test-classes" />
<content url="file://$MODULE_DIR$">
<sourceFolder url="file://$MODULE_DIR$/src/main/java" isTestSource="false" />
<sourceFolder url="file://$MODULE_DIR$/src/test/java" isTestSource="true" />
<sourceFolder url="file://$MODULE_DIR$/src/main/resources" type="java-resource" />
<excludeFolder url="file://$MODULE_DIR$/target" />
</content>
<orderEntry type="inheritedJdk" />
<orderEntry type="sourceFolder" forTests="false" />
<orderEntry type="library" scope="TEST" name="Maven: junit:junit:4.12" level="project" />
<orderEntry type="library" scope="TEST" name="Maven: org.hamcrest:hamcrest-core:1.3" level="project" />
<orderEntry type="library" name="Maven: org.apache.hbase:hbase-client:1.1.2" level="project" />
<orderEntry type="library" name="Maven: org.apache.hbase:hbase-annotations:1.1.2" level="project" />
<orderEntry type="module-library">
<library name="Maven: jdk.tools:jdk.tools:1.7">
<CLASSES>
<root url="jar://D:/Programs/Java/jdk1.7.0_75/lib/tools.jar!/" />
</CLASSES>
<JAVADOC />
<SOURCES />
</library>
</orderEntry>
<orderEntry type="library" name="Maven: org.apache.hbase:hbase-common:1.1.2" level="project" />
<orderEntry type="library" name="Maven: org.apache.hbase:hbase-protocol:1.1.2" level="project" />
<orderEntry type="library" name="Maven: commons-codec:commons-codec:1.9" level="project" />
<orderEntry type="library" name="Maven: commons-io:commons-io:2.4" level="project" />
<orderEntry type="library" name="Maven: commons-lang:commons-lang:2.6" level="project" />
<orderEntry type="library" name="Maven: commons-logging:commons-logging:1.2" level="project" />
<orderEntry type="library" name="Maven: com.google.guava:guava:12.0.1" level="project" />
<orderEntry type="library" name="Maven: com.google.code.findbugs:jsr305:1.3.9" level="project" />
<orderEntry type="library" name="Maven: com.google.protobuf:protobuf-java:2.5.0" level="project" />
<orderEntry type="library" name="Maven: io.netty:netty-all:4.0.23.Final" level="project" />
<orderEntry type="library" name="Maven: org.apache.zookeeper:zookeeper:3.4.6" level="project" />
<orderEntry type="library" name="Maven: org.slf4j:slf4j-api:1.6.1" level="project" />
<orderEntry type="library" name="Maven: org.slf4j:slf4j-log4j12:1.6.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.htrace:htrace-core:3.1.0-incubating" level="project" />
<orderEntry type="library" name="Maven: org.codehaus.jackson:jackson-mapper-asl:1.9.13" level="project" />
<orderEntry type="library" name="Maven: org.jruby.jcodings:jcodings:1.0.8" level="project" />
<orderEntry type="library" name="Maven: org.jruby.joni:joni:2.1.2" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-auth:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.httpcomponents:httpclient:4.2.5" level="project" />
<orderEntry type="library" name="Maven: org.apache.httpcomponents:httpcore:4.2.4" level="project" />
<orderEntry type="library" name="Maven: org.apache.directory.server:apacheds-kerberos-codec:2.0.0-M15" level="project" />
<orderEntry type="library" name="Maven: org.apache.directory.server:apacheds-i18n:2.0.0-M15" level="project" />
<orderEntry type="library" name="Maven: org.apache.directory.api:api-asn1-api:1.0.0-M20" level="project" />
<orderEntry type="library" name="Maven: org.apache.directory.api:api-util:1.0.0-M20" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-common:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-annotations:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.commons:commons-math3:3.1.1" level="project" />
<orderEntry type="library" name="Maven: xmlenc:xmlenc:0.52" level="project" />
<orderEntry type="library" name="Maven: commons-net:commons-net:3.1" level="project" />
<orderEntry type="library" name="Maven: commons-el:commons-el:1.0" level="project" />
<orderEntry type="library" name="Maven: commons-configuration:commons-configuration:1.6" level="project" />
<orderEntry type="library" name="Maven: commons-digester:commons-digester:1.8" level="project" />
<orderEntry type="library" name="Maven: commons-beanutils:commons-beanutils:1.7.0" level="project" />
<orderEntry type="library" name="Maven: commons-beanutils:commons-beanutils-core:1.8.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.avro:avro:1.7.4" level="project" />
<orderEntry type="library" name="Maven: com.thoughtworks.paranamer:paranamer:2.3" level="project" />
<orderEntry type="library" name="Maven: org.xerial.snappy:snappy-java:1.0.4.1" level="project" />
<orderEntry type="library" name="Maven: com.jcraft:jsch:0.1.42" level="project" />
<orderEntry type="library" name="Maven: org.apache.commons:commons-compress:1.4.1" level="project" />
<orderEntry type="library" name="Maven: org.tukaani:xz:1.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-mapreduce-client-core:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-yarn-common:2.5.1" level="project" />
<orderEntry type="library" name="Maven: javax.xml.bind:jaxb-api:2.2.2" level="project" />
<orderEntry type="library" name="Maven: javax.xml.stream:stax-api:1.0-2" level="project" />
<orderEntry type="library" name="Maven: javax.activation:activation:1.1" level="project" />
<orderEntry type="library" name="Maven: io.netty:netty:3.6.2.Final" level="project" />
<orderEntry type="library" name="Maven: com.github.stephenc.findbugs:findbugs-annotations:1.3.9-1" level="project" />
<orderEntry type="library" name="Maven: org.apache.hbase:hbase-server:1.1.2" level="project" />
<orderEntry type="library" name="Maven: org.apache.hbase:hbase-procedure:1.1.2" level="project" />
<orderEntry type="library" name="Maven: org.apache.hbase:hbase-common:tests:1.1.2" level="project" />
<orderEntry type="library" scope="RUNTIME" name="Maven: org.apache.hbase:hbase-prefix-tree:1.1.2" level="project" />
<orderEntry type="library" name="Maven: commons-httpclient:commons-httpclient:3.1" level="project" />
<orderEntry type="library" name="Maven: commons-collections:commons-collections:3.2.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.hbase:hbase-hadoop-compat:1.1.2" level="project" />
<orderEntry type="library" name="Maven: org.apache.hbase:hbase-hadoop2-compat:1.1.2" level="project" />
<orderEntry type="library" name="Maven: com.yammer.metrics:metrics-core:2.2.0" level="project" />
<orderEntry type="library" name="Maven: com.sun.jersey:jersey-core:1.9" level="project" />
<orderEntry type="library" name="Maven: com.sun.jersey:jersey-server:1.9" level="project" />
<orderEntry type="library" name="Maven: asm:asm:3.1" level="project" />
<orderEntry type="library" name="Maven: commons-cli:commons-cli:1.2" level="project" />
<orderEntry type="library" name="Maven: org.apache.commons:commons-math:2.2" level="project" />
<orderEntry type="library" name="Maven: log4j:log4j:1.2.17" level="project" />
<orderEntry type="library" name="Maven: org.mortbay.jetty:jetty:6.1.26" level="project" />
<orderEntry type="library" name="Maven: org.mortbay.jetty:jetty-util:6.1.26" level="project" />
<orderEntry type="library" name="Maven: org.mortbay.jetty:jetty-sslengine:6.1.26" level="project" />
<orderEntry type="library" name="Maven: org.mortbay.jetty:jsp-2.1:6.1.14" level="project" />
<orderEntry type="library" name="Maven: org.mortbay.jetty:jsp-api-2.1:6.1.14" level="project" />
<orderEntry type="library" name="Maven: org.mortbay.jetty:servlet-api-2.5:6.1.14" level="project" />
<orderEntry type="library" name="Maven: org.codehaus.jackson:jackson-core-asl:1.9.13" level="project" />
<orderEntry type="library" name="Maven: org.codehaus.jackson:jackson-jaxrs:1.9.13" level="project" />
<orderEntry type="library" name="Maven: tomcat:jasper-compiler:5.5.23" level="project" />
<orderEntry type="library" name="Maven: tomcat:jasper-runtime:5.5.23" level="project" />
<orderEntry type="library" name="Maven: org.jamon:jamon-runtime:2.3.1" level="project" />
<orderEntry type="library" name="Maven: com.lmax:disruptor:3.3.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-client:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-mapreduce-client-app:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-mapreduce-client-common:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-yarn-client:2.5.1" level="project" />
<orderEntry type="library" name="Maven: com.sun.jersey:jersey-client:1.9" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-yarn-server-common:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-mapreduce-client-shuffle:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.fusesource.leveldbjni:leveldbjni-all:1.8" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-yarn-api:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-mapreduce-client-jobclient:2.5.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-hdfs:2.5.1" level="project" />
<orderEntry type="library" name="Maven: commons-daemon:commons-daemon:1.0.13" level="project" />
<orderEntry type="library" name="Maven: org.apache.logging.log4j:log4j-core:2.4.1" level="project" />
<orderEntry type="library" name="Maven: org.apache.logging.log4j:log4j-api:2.4.1" level="project" />
<orderEntry type="module" module-name="skywalking-protocol" />
</component>
</module>

View File

@ -1,112 +0,0 @@
package com.ai.cloud.skywalking.analysis.categorize2chain;
import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Comparator;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
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 com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessChain;
import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter;
import com.ai.cloud.skywalking.analysis.chainbuild.util.VersionIdentifier;
import com.ai.cloud.skywalking.analysis.config.ConfigInitializer;
import com.ai.cloud.skywalking.protocol.Span;
public class Categorize2ChainMapper extends TableMapper<Text, ChainInfo> {
private Logger logger = LoggerFactory
.getLogger(Categorize2ChainMapper.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 {
if(!VersionIdentifier.enableAnaylsis(Bytes.toString(key.get()))){
return;
}
List<Span> spanList = new ArrayList<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(Bytes.toString(key.get()), spanList);
logger.info("Success convert span to chain info...."
+ chainInfo.getCID() + " TraceId : " + Bytes.toString(key.get()));
context.write(
new Text(chainInfo.getUserId() + ":"
+ chainInfo.getEntranceNodeToken()), chainInfo);
} catch (Exception e) {
logger.error("Failed to mapper call chain[" + key.toString() + "]",
e);
}
}
public static ChainInfo spanToChainInfo(String key, List<Span> spanList) {
SubLevelSpanCostCounter costMap = new SubLevelSpanCostCounter();
ChainInfo chainInfo = new ChainInfo();
Collections.sort(spanList, new Comparator<Span>() {
@Override
public int compare(Span span1, Span span2) {
String span1TraceLevel = span1.getParentLevel() + "."
+ span1.getLevelId();
String span2TraceLevel = span2.getParentLevel() + "."
+ span2.getLevelId();
return span1TraceLevel.compareTo(span2TraceLevel);
}
});
Map<String, SpanEntry> spanEntryMap = mergeSpanDataSet(spanList);
for (Map.Entry<String, SpanEntry> entry : spanEntryMap.entrySet()) {
ChainNode chainNode = new ChainNode();
SpanNodeProcessFilter filter = SpanNodeProcessChain
.getProcessChainByCallType(entry.getValue().getSpanType());
filter.doFilter(entry.getValue(), chainNode, costMap);
chainInfo.addNodes(chainNode);
}
chainInfo.generateChainToken();
HBaseUtil.saveCidTidMapping(key, chainInfo);
return chainInfo;
}
private static Map<String, SpanEntry> mergeSpanDataSet(List<Span> spanList) {
Map<String, SpanEntry> spanEntryMap = new LinkedHashMap<String, SpanEntry>();
for (int i = spanList.size() - 1; i >= 0; i--) {
Span span = spanList.get(i);
SpanEntry spanEntry = spanEntryMap.get(span.getParentLevel() + "."
+ span.getLevelId());
if (spanEntry == null) {
spanEntry = new SpanEntry();
spanEntryMap.put(
span.getParentLevel() + "." + span.getLevelId(),
spanEntry);
}
spanEntry.setSpan(span);
}
return spanEntryMap;
}
}

View File

@ -1,57 +0,0 @@
package com.ai.cloud.skywalking.analysis.categorize2chain;
import java.io.IOException;
import java.util.Iterator;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainRelationship;
import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainSummaryWithoutRelationship;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil;
import com.ai.cloud.skywalking.analysis.config.ConfigInitializer;
public class Categorize2ChainReducer extends Reducer<Text, ChainInfo, Text, IntWritable> {
private static Logger logger = LoggerFactory.getLogger(Categorize2ChainReducer.class.getName());
@Override
protected void setup(Context context) throws IOException,
InterruptedException {
ConfigInitializer.initialize();
}
@Override
protected void reduce(Text key, Iterable<ChainInfo> values, Context context) throws IOException, InterruptedException {
int totalCount = reduceAction(key.toString(), values.iterator());
context.write(new Text(key.toString()), new IntWritable(totalCount));
}
public static int reduceAction(String key, Iterator<ChainInfo> chainInfoIterator) throws IOException, InterruptedException {
int totalCount = 0;
try {
ChainRelationship chainRelate = HBaseUtil.loadCallChainRelationship(key.toString());
ChainSummaryWithoutRelationship summary = new ChainSummaryWithoutRelationship();
while (chainInfoIterator.hasNext()) {
ChainInfo chainInfo = chainInfoIterator.next();
try {
chainRelate.categoryChain(chainInfo);
summary.summary(chainInfo);
} catch (Exception e) {
continue;
}
totalCount++;
}
chainRelate.save();
summary.save();
} catch (Exception e) {
logger.error("Failed to reduce key[" + key + "]", e);
}
return totalCount;
}
}

View File

@ -1,77 +0,0 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.entity;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.google.gson.Gson;
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;
public class CategorizedChainInfo {
private String cid;
private String chainFullToken;
private List<String> children_Token;
public CategorizedChainInfo(ChainInfo chainInfo) {
cid = chainInfo.getCID();
StringBuilder stringBuilder = new StringBuilder();
boolean flag = false;
for (ChainNode chainNode : chainInfo.getNodes()) {
if (flag) {
stringBuilder.append(";");
}
stringBuilder.append(chainNode.getTraceLevelId() + "-" + chainNode.getNodeToken());
flag = true;
}
chainFullToken = stringBuilder.toString();
children_Token = new ArrayList<String>();
}
public CategorizedChainInfo(String value) {
JsonObject jsonObject = (JsonObject) new JsonParser().parse(value);
cid = jsonObject.get("chainToken").getAsString();
chainFullToken = jsonObject.get("chainFullToken").getAsString();
children_Token = new Gson().fromJson(jsonObject.get("children_Token"),
new TypeToken<List<String>>() {
}.getType());
}
public String getChainFullToken() {
return chainFullToken;
}
public boolean isContained(UncategorizeChainInfo uncategorizeChainInfo) {
Pattern pattern = Pattern.compile(uncategorizeChainInfo.getNodeRegEx());
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());
}
public void add(UncategorizeChainInfo uncategorizeChainInfo) {
children_Token.add(uncategorizeChainInfo.getCID());
}
@Override
public String toString() {
return new Gson().toJson(this);
}
public List<String> getChildren_Token() {
return children_Token;
}
}

View File

@ -1,71 +0,0 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.entity;
import java.util.HashMap;
import java.util.Map;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.google.gson.Gson;
import com.google.gson.JsonObject;
import com.google.gson.JsonParser;
import com.google.gson.reflect.TypeToken;
public class ChainNodeSpecificTimeWindowSummary {
public static final long INTERVAL = 1L;
private String traceLevelId;
private String nodeToken;
// key : 分钟
private Map<String, ChainNodeSpecificTimeWindowSummaryValue> summerValueMap;
public static ChainNodeSpecificTimeWindowSummary newInstance(String traceLevelId, String nodeToken) {
ChainNodeSpecificTimeWindowSummary cns = new ChainNodeSpecificTimeWindowSummary();
cns.traceLevelId = traceLevelId;
cns.nodeToken = nodeToken;
return cns;
}
private ChainNodeSpecificTimeWindowSummary() {
summerValueMap = new HashMap<String, ChainNodeSpecificTimeWindowSummaryValue>();
}
public ChainNodeSpecificTimeWindowSummary(String value) {
JsonObject jsonObject = (JsonObject) new JsonParser().parse(value);
traceLevelId = jsonObject.get("traceLevelId").getAsString();
summerValueMap = new Gson().fromJson(jsonObject.get("summerValueMap").toString(),
new TypeToken<Map<String, ChainNodeSpecificTimeWindowSummaryValue>>() {
}.getType());
nodeToken = jsonObject.get("nodeToken").getAsString();
}
public String getTraceLevelId() {
return traceLevelId;
}
public void summary(ChainNode node) {
String key = generateKey(node.getStartDate());
ChainNodeSpecificTimeWindowSummaryValue summaryResult = summerValueMap.get(key);
if (summaryResult == null) {
summaryResult = new ChainNodeSpecificTimeWindowSummaryValue();
summerValueMap.put(key, summaryResult);
}
summaryResult.summary(node);
}
private String generateKey(long startTime) {
long minutes = (startTime % (1000 * 60 * 60)) / (1000 * 60);
return String.valueOf(minutes / INTERVAL);
}
@Override
public String toString() {
return new Gson().toJson(this);
}
public String getNodeToken() {
return nodeToken;
}
public Map<String, ChainNodeSpecificTimeWindowSummaryValue> getSummerValueMap() {
return summerValueMap;
}
}

View File

@ -1,51 +0,0 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.entity;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
public class ChainNodeSpecificTimeWindowSummaryValue {
private long totalCall;
private long totalCostTime;
private long correctNumber;
private long humanInterruptionNumber;
public ChainNodeSpecificTimeWindowSummaryValue() {
totalCall = 0;
totalCostTime = 0;
correctNumber = 0;
humanInterruptionNumber = 0;
}
public long getTotalCall() {
return totalCall;
}
public long getTotalCostTime() {
return totalCostTime;
}
public long getCorrectNumber() {
return correctNumber;
}
public long getHumanInterruptionNumber() {
return humanInterruptionNumber;
}
public void summary(ChainNode node) {
totalCall++;
if (node.getStatus() == ChainNode.NodeStatus.NORMAL) {
correctNumber++;
}
if (node.getStatus() == ChainNode.NodeStatus.HUMAN_INTERRUPTION) {
humanInterruptionNumber++;
}
totalCostTime += node.getCost();
}
public void accumulate(ChainNodeSpecificTimeWindowSummaryValue value) {
this.totalCall += value.getTotalCall();
this.correctNumber += value.getCorrectNumber();
this.totalCostTime += value.getTotalCostTime();
this.humanInterruptionNumber += value.getHumanInterruptionNumber();
}
}

View File

@ -1,152 +0,0 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.entity;
import java.io.IOException;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Set;
import org.apache.hadoop.hbase.client.Put;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil;
import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData;
import com.google.gson.GsonBuilder;
public class ChainRelationship {
private static Logger logger = LoggerFactory.getLogger(ChainRelationship.class.getName());
private String key;
private Map<String, CategorizedChainInfo> categorizedChainInfoMap = new HashMap<String, CategorizedChainInfo>();
private Set<UncategorizeChainInfo> uncategorizeChainInfoSet = new HashSet<UncategorizeChainInfo>();
private Map<String, ChainDetail> chainDetailMap = new HashMap<String, ChainDetail>();
public ChainRelationship(String key) {
this.key = key;
}
private void categoryAllUncategorizedChainInfo(CategorizedChainInfo parentChains) {
if (uncategorizeChainInfoSet != null && uncategorizeChainInfoSet.size() > 0) {
Iterator<UncategorizeChainInfo> uncategorizeChainInfoIterator = uncategorizeChainInfoSet.iterator();
while (uncategorizeChainInfoIterator.hasNext()) {
UncategorizeChainInfo uncategorizeChainInfo = uncategorizeChainInfoIterator.next();
if (parentChains.isContained(uncategorizeChainInfo)) {
parentChains.add(uncategorizeChainInfo);
uncategorizeChainInfoIterator.remove();
}
}
}
}
private void try2CategoryUncategorizedChainInfo(UncategorizeChainInfo child) {
boolean isContained = false;
for (Map.Entry<String, CategorizedChainInfo> entry : categorizedChainInfoMap.entrySet()) {
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;
}
}
if (!isContained) {
if (!uncategorizeChainInfoSet.contains(child)) {
chainDetailMap.put(child.getCID(), new ChainDetail(child.getChainInfo(), false));
uncategorizeChainInfoSet.add(child);
}
}
}
private CategorizedChainInfo addCategorizedChain(ChainInfo chainInfo) {
if (!categorizedChainInfoMap.containsKey(chainInfo.getCID())) {
categorizedChainInfoMap.put(chainInfo.getCID(),
new CategorizedChainInfo(chainInfo));
chainDetailMap.put(chainInfo.getCID(), new ChainDetail(chainInfo, true));
}
return categorizedChainInfoMap.get(chainInfo.getCID());
}
public void categoryChain(ChainInfo chainInfo) {
if (chainInfo.getChainStatus() == ChainInfo.ChainStatus.NORMAL) {
CategorizedChainInfo categorizedChainInfo = addCategorizedChain(chainInfo);
categoryAllUncategorizedChainInfo(categorizedChainInfo);
} else {
UncategorizeChainInfo uncategorizeChainInfo = new UncategorizeChainInfo(chainInfo);
try2CategoryUncategorizedChainInfo(uncategorizeChainInfo);
}
}
public void save() throws SQLException, IOException, InterruptedException {
saveChainRelationship();
saveChainDetail();
}
private void saveChainDetail() throws SQLException, IOException, InterruptedException {
List<Put> puts = new ArrayList<Put>();
for (Map.Entry<String, ChainDetail> entry : chainDetailMap.entrySet()) {
Put put1 = new Put(entry.getKey().getBytes());
entry.getValue().save(put1);
puts.add(put1);
}
try {
HBaseUtil.saveChainDetails(puts);
} catch (IOException e) {
logger.error("Faild to save chain detail to hbase.", e);
throw e;
} catch (InterruptedException e) {
logger.error("Faild to save chain detail to hbase.", e);
throw e;
}
}
private void saveChainRelationship() throws IOException {
Put put = new Put(getKey().getBytes());
put.addColumn(HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.COLUMN_FAMILY_NAME.getBytes(), HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.UNCATEGORIZE_COLUMN_NAME.getBytes()
, new GsonBuilder().excludeFieldsWithoutExposeAnnotation().create().toJson(getUncategorizeChainInfoList()).getBytes());
for (Map.Entry<String, CategorizedChainInfo> entry : getCategorizedChainInfoMap().entrySet()) {
put.addColumn(HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.COLUMN_FAMILY_NAME.getBytes(), entry.getKey().getBytes()
, entry.getValue().toString().getBytes());
}
try {
HBaseUtil.saveChainRelationship(put);
} catch (IOException e) {
logger.error("Faild to save chain relationship to hbase.", e);
throw e;
}
}
public void addCategorizeChain(String qualifierName, CategorizedChainInfo categorizedChainInfo) {
categorizedChainInfoMap.put(qualifierName, categorizedChainInfo);
}
public String getKey() {
return key;
}
public Map<String, CategorizedChainInfo> getCategorizedChainInfoMap() {
return categorizedChainInfoMap;
}
public Set<UncategorizeChainInfo> getUncategorizeChainInfoList() {
return uncategorizeChainInfoSet;
}
public void addUncategorizeChain(List<UncategorizeChainInfo> uncategorizeChainInfos) {
uncategorizeChainInfoSet.addAll(uncategorizeChainInfos);
}
}

View File

@ -1,60 +0,0 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.entity;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil;
import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData;
import org.apache.hadoop.hbase.client.Put;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
public class ChainSpecificTimeWindowSummary {
private static Logger logger = LoggerFactory.getLogger(ChainSpecificTimeWindowSummary.class.getName());
/**
* key : cid + uid + 时间窗口
*/
private Map<String, ChainNodeSpecificTimeWindowSummary> chainNodeSummaryResultMap;
public ChainSpecificTimeWindowSummary() {
chainNodeSummaryResultMap = new HashMap<String, ChainNodeSpecificTimeWindowSummary>();
}
public static ChainSpecificTimeWindowSummary load(String cid_uid_time) {
ChainSpecificTimeWindowSummary result = null;
try {
result = HBaseUtil.selectChainSummaryResult(cid_uid_time);
} catch (IOException e) {
logger.error("Failed to load the key[" + cid_uid_time + "] summary result.", e);
}
if (result == null) {
result = new ChainSpecificTimeWindowSummary();
}
return result;
}
public void addNodeSummaryResult(ChainNodeSpecificTimeWindowSummary chainNodeSummaryResult) {
chainNodeSummaryResultMap.put(chainNodeSummaryResult.getTraceLevelId(), chainNodeSummaryResult);
}
public void summaryNodeValue(ChainNode node) {
String tlid = node.getTraceLevelId();
ChainNodeSpecificTimeWindowSummary chainNodeSummaryResult = chainNodeSummaryResultMap.get(tlid);
if (chainNodeSummaryResult == null) {
chainNodeSummaryResult = ChainNodeSpecificTimeWindowSummary.newInstance(tlid, node.getNodeToken());
chainNodeSummaryResultMap.put(tlid, chainNodeSummaryResult);
}
chainNodeSummaryResult.summary(node);
}
public void save(Put put) {
for (Map.Entry<String, ChainNodeSpecificTimeWindowSummary> entry : chainNodeSummaryResultMap.entrySet()) {
put.addColumn(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP.COLUMN_FAMILY_NAME.getBytes(), entry.getKey().getBytes(), entry.getValue().toString().getBytes());
}
}
}

View File

@ -1,65 +0,0 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.entity;
import com.ai.cloud.skywalking.analysis.categorize2chain.DBCallChainInfoDao;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil;
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.sql.Timestamp;
import java.text.SimpleDateFormat;
import java.util.*;
public class ChainSummaryWithoutRelationship {
private static Logger logger = LoggerFactory.getLogger(ChainSummaryWithoutRelationship.class.getName());
private Map<String, ChainSpecificTimeWindowSummary> loadedChainSpecificTimeWindowSummary;
private Map<String, Timestamp> updateChainInfo;
public ChainSummaryWithoutRelationship() {
loadedChainSpecificTimeWindowSummary = new HashMap<String, ChainSpecificTimeWindowSummary>();
updateChainInfo = new HashMap<String, Timestamp>();
}
public void summary(ChainInfo chainInfo) {
for (ChainNode node : chainInfo.getNodes()) {
String csk = generateChainSummaryKey(chainInfo, node.getStartDate());
if (!loadedChainSpecificTimeWindowSummary.containsKey(csk)) {
loadedChainSpecificTimeWindowSummary.put(csk, ChainSpecificTimeWindowSummary.load(csk));
}
loadedChainSpecificTimeWindowSummary.get(csk).summaryNodeValue(node);
}
updateChainInfo.put(chainInfo.getCID(), new Timestamp(System.currentTimeMillis()));
}
private String generateChainSummaryKey(ChainInfo chainInfo, long startDate) {
return chainInfo.getCID() + "-" + chainInfo.getUserId() + "-" + new SimpleDateFormat("yyyy/MM/dd HH:mm:ss").
format(new Date(startDate / (1000 * 60 * 60) * (1000 * 60 * 60)));
}
public void save() throws IOException, InterruptedException, SQLException {
batchSaveChainSpecificTimeWindowSummary();
updateChainLastActiveTime();
}
private void updateChainLastActiveTime() throws SQLException {
DBCallChainInfoDao.updateChainLastActiveTime(updateChainInfo);
}
private void batchSaveChainSpecificTimeWindowSummary() throws IOException, InterruptedException {
List<Put> puts = new ArrayList<Put>();
logger.info("There are [" + loadedChainSpecificTimeWindowSummary.size() + "] summary data will be storage to HBase");
for (Map.Entry<String, ChainSpecificTimeWindowSummary> entry : loadedChainSpecificTimeWindowSummary.entrySet()) {
Put put = new Put(entry.getKey().getBytes());
entry.getValue().save(put);
puts.add(put);
}
HBaseUtil.batchSaveChainSpecificTimeWindowSummary(puts);
}
}

View File

@ -1,69 +0,0 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.entity;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.google.gson.GsonBuilder;
import com.google.gson.annotations.Expose;
public class UncategorizeChainInfo {
@Expose
private String cid;
@Expose
private String nodeRegEx;
private ChainInfo chainInfo;
public UncategorizeChainInfo() {
}
public UncategorizeChainInfo(ChainInfo chainInfo) {
this.cid = chainInfo.getCID();
StringBuilder stringBuilder = new StringBuilder();
boolean flag = false;
for (ChainNode node : chainInfo.getNodes()) {
if (flag) {
stringBuilder.append(";*");
}
stringBuilder.append((node.getTraceLevelId() + "-" + node.getNodeToken()));
flag = true;
}
nodeRegEx = stringBuilder.toString();
this.chainInfo = chainInfo;
}
public String getCID() {
return cid;
}
public String getNodeRegEx() {
return nodeRegEx;
}
public ChainInfo getChainInfo() {
return chainInfo;
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (!(o instanceof UncategorizeChainInfo)) return false;
UncategorizeChainInfo that = (UncategorizeChainInfo) o;
return cid != null ? cid.equals(that.cid) : that.cid == null;
}
@Override
public int hashCode() {
return cid != null ? cid.hashCode() : 0;
}
@Override
public String toString() {
GsonBuilder gsonBuilder = new GsonBuilder();
gsonBuilder.excludeFieldsWithoutExposeAnnotation();
return gsonBuilder.create().toJson(this);
}
}

View File

@ -1,16 +0,0 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl;
import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry;
import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter;
public class AppendBusinessKeyFilter extends SpanNodeProcessFilter {
@Override
public void doFilter(SpanEntry spanEntry, ChainNode node, SubLevelSpanCostCounter costMap) {
node.setViewPoint(node.getViewPoint() + spanEntry.getBusinessKey());
this.doNext(spanEntry, node, costMap);
}
}

View File

@ -0,0 +1,88 @@
package com.ai.cloud.skywalking.analysis.chainbuild;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.analysis.chainbuild.po.CallChainTreeNode;
import com.ai.cloud.skywalking.analysis.chainbuild.util.HBaseUtil;
import org.apache.hadoop.hbase.client.Put;
import java.io.IOException;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
public class CallChainTree {
private String callEntrance;
//存放已经合并过的调用链ID
private List<String> hasBeenMergedChainIds;
// 本次Reduce合并过的调用链
private Map<String, ChainInfo> combineChains;
//合并之后的节点
// key : trace level Id
private Map<String, CallChainTreeNode> nodes;
public CallChainTree(String callEntrance) {
hasBeenMergedChainIds = new ArrayList<String>();
combineChains = new HashMap<String, ChainInfo>();
nodes = new HashMap<String, CallChainTreeNode>();
this.callEntrance = callEntrance;
}
public static CallChainTree load(String callEntrance) throws IOException {
CallChainTree chain = HBaseUtil.loadMergedCallChain(callEntrance);
chain.hasBeenMergedChainIds.addAll(HBaseUtil.loadHasBeenMergeChainIds(callEntrance));
if (chain == null) {
chain = new CallChainTree(callEntrance);
}
return chain;
}
public void processMerge(ChainInfo chainInfo) {
if (hasBeenMergedChainIds.contains(chainInfo.getCID())) {
return;
}
for (ChainNode node : chainInfo.getNodes()) {
CallChainTreeNode callChainTreeNode = nodes.get(node.getTraceLevelId());
if (callChainTreeNode != null) {
callChainTreeNode.mergeIfNess(node);
} else {
nodes.put(node.getTraceLevelId(), new CallChainTreeNode(node));
}
}
hasBeenMergedChainIds.add(chainInfo.getChainToken());
combineChains.put(chainInfo.getChainToken(), chainInfo);
}
public void summary(ChainInfo chainInfo) {
for (ChainNode node : chainInfo.getNodes()) {
CallChainTreeNode callChainTreeNode = nodes.get(node.getTraceLevelId());
callChainTreeNode.summary(node);
}
}
public void saveToHbase() {
List<Put> chainInfoPuts = new ArrayList<Put>();
for (Map.Entry<String, ChainInfo> entry : combineChains.entrySet()) {
Put put = new Put(entry.getKey().getBytes());
entry.getValue().saveToHBase(put);
chainInfoPuts.add(put);
}
HBaseUtil.saveMergedCallChain(this);
}
public String getCallEntrance() {
return callEntrance;
}
public void addMergedChainNode(CallChainTreeNode chainNode) {
nodes.put(chainNode.getTraceLevelId(), chainNode);
}
}

View File

@ -1,6 +1,10 @@
package com.ai.cloud.skywalking.analysis.chainbuild;
import com.ai.cloud.skywalking.analysis.chainbuild.entity.TraceSpanTree;
import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessChain;
import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter;
import com.ai.cloud.skywalking.analysis.chainbuild.util.VersionIdentifier;
import com.ai.cloud.skywalking.analysis.config.ConfigInitializer;
import com.ai.cloud.skywalking.protocol.Span;
@ -14,12 +18,12 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.*;
public class ChainBuildMapper extends TableMapper<Text, ChainInfo> {
public class ChainBuildMapper extends TableMapper<Text, TraceSpanTree> {
private Logger logger = LoggerFactory
.getLogger(ChainBuildMapper.class);
.getLogger(ChainBuildMapper.class.getName());
@Override
protected void setup(Context context) throws IOException,
@ -27,6 +31,7 @@ public class ChainBuildMapper extends TableMapper<Text, TraceSpanTree> {
ConfigInitializer.initialize();
}
@Override
protected void map(ImmutableBytesWritable key, Result value, Context context)
throws IOException, InterruptedException {
@ -34,21 +39,68 @@ public class ChainBuildMapper extends TableMapper<Text, TraceSpanTree> {
return;
}
List<Span> spanList = new ArrayList<Span>();
ChainInfo chainInfo = null;
try {
List<Span> spanList = new ArrayList<Span>();
for (Cell cell : value.rawCells()) {
Span span = new Span(Bytes.toString(cell.getValueArray(),
cell.getValueOffset(), cell.getValueLength()));
spanList.add(span);
}
TraceSpanTree tree = new TraceSpanTree();
tree.build(spanList);
context.write(new Text(tree.getCid()), tree);
} catch (Throwable e) {
chainInfo = spanToChainInfo(Bytes.toString(key.get()), spanList);
logger.info("Success convert span to chain info...."
+ chainInfo.getCID() + " TraceId : " + Bytes.toString(key.get()));
context.write(
new Text(chainInfo.getUserId() + ":"
+ chainInfo.getEntranceNodeToken()), chainInfo);
} catch (Exception e) {
logger.error("Failed to mapper call chain[" + key.toString() + "]",
e);
}
}
public static ChainInfo spanToChainInfo(String key, List<Span> spanList) {
SubLevelSpanCostCounter costMap = new SubLevelSpanCostCounter();
ChainInfo chainInfo = new ChainInfo();
Collections.sort(spanList, new Comparator<Span>() {
@Override
public int compare(Span span1, Span span2) {
String span1TraceLevel = span1.getParentLevel() + "."
+ span1.getLevelId();
String span2TraceLevel = span2.getParentLevel() + "."
+ span2.getLevelId();
return span1TraceLevel.compareTo(span2TraceLevel);
}
});
Map<String, SpanEntry> spanEntryMap = mergeSpanDataSet(spanList);
for (Map.Entry<String, SpanEntry> entry : spanEntryMap.entrySet()) {
ChainNode chainNode = new ChainNode();
SpanNodeProcessFilter filter = SpanNodeProcessChain
.getProcessChainByCallType(entry.getValue().getSpanType());
filter.doFilter(entry.getValue(), chainNode, costMap);
chainInfo.addNodes(chainNode);
}
//chainInfo.generateChainToken();
//HBaseUtil.saveCidTidMapping(key, chainInfo);
return chainInfo;
}
private static Map<String, SpanEntry> mergeSpanDataSet(List<Span> spanList) {
Map<String, SpanEntry> spanEntryMap = new LinkedHashMap<String, SpanEntry>();
for (int i = spanList.size() - 1; i >= 0; i--) {
Span span = spanList.get(i);
SpanEntry spanEntry = spanEntryMap.get(span.getParentLevel() + "."
+ span.getLevelId());
if (spanEntry == null) {
spanEntry = new SpanEntry();
spanEntryMap.put(
span.getParentLevel() + "." + span.getLevelId(),
spanEntry);
}
spanEntry.setSpan(span);
}
return spanEntryMap;
}
}

View File

@ -0,0 +1,42 @@
package com.ai.cloud.skywalking.analysis.chainbuild;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo;
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 org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.Iterator;
public class ChainBuildReducer extends Reducer<Text, ChainInfo, Text, IntWritable> {
private Logger logger = LoggerFactory
.getLogger(ChainBuildReducer.class.getName());
@Override
protected void setup(Context context) throws IOException,
InterruptedException {
ConfigInitializer.initialize();
}
@Override
protected void reduce(Text key, Iterable<ChainInfo> values, Context context) throws IOException,
InterruptedException {
CallChainTree chainTree = CallChainTree.load(Bytes.toString(key.getBytes()));
Iterator<ChainInfo> chainInfoIterator = values.iterator();
while (chainInfoIterator.hasNext()) {
ChainInfo chainInfo = chainInfoIterator.next();
if (chainInfo.getChainStatus() == ChainInfo.ChainStatus.NORMAL) {
chainTree.processMerge(chainInfo);
}
//合并数据
chainTree.summary(chainInfo);
}
chainTree.saveToHbase();
}
}

View File

@ -1,9 +1,8 @@
package com.ai.cloud.skywalking.analysis.categorize2chain;
package com.ai.cloud.skywalking.analysis.chainbuild;
import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainDetail;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.chainbuild.entity.CallChainDetail;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.analysis.config.Config;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -31,16 +30,16 @@ public class DBCallChainInfoDao {
}
}
public synchronized static void saveChainDetail(ChainDetail chainDetail)
public synchronized static void saveChainDetail(CallChainDetail callChainDetail)
throws SQLException {
PreparedStatement preparedStatement = null;
try {
preparedStatement = connection
.prepareStatement("INSERT INTO sw_chain_detail(cid,uid,traceLevelId,viewpoint,create_time)"
+ " VALUES(?,?,?,?,?)");
for (ChainNode chainNode : chainDetail.getChainNodes()) {
preparedStatement.setString(1, chainDetail.getChainToken());
preparedStatement.setString(2, chainDetail.getUserId());
for (ChainNode chainNode : callChainDetail.getChainNodes()) {
preparedStatement.setString(1, callChainDetail.getChainToken());
preparedStatement.setString(2, callChainDetail.getUserId());
preparedStatement.setString(3, chainNode.getTraceLevelId());
preparedStatement.setString(4, chainNode.getViewPoint() + ":"
+ chainNode.getBusinessKey());
@ -52,7 +51,7 @@ public class DBCallChainInfoDao {
for (int i : result) {
if (i != 1) {
logger.error("Failed to save chain detail ["
+ chainDetail.getChainToken() + "]");
+ callChainDetail.getChainToken() + "]");
}
}
} finally {

View File

@ -1,6 +1,6 @@
package com.ai.cloud.skywalking.analysis.categorize2chain;
package com.ai.cloud.skywalking.analysis.chainbuild;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.protocol.CallType;
import com.ai.cloud.skywalking.protocol.Span;

View File

@ -1,39 +0,0 @@
package com.ai.cloud.skywalking.analysis.chainbuild.entity;
import java.util.List;
public class BranchTraceSpanNode extends TraceSpanNode {
protected BranchTraceSpanNode(TraceSpanNode origin, TraceSpanNode dest,
List<TraceSpanNode> spanContainer) {
setNextBranchNode(dest);
dest.parent = this;
dest.setNextBranchNode(origin);
origin.parent = this;
this.setParent(dest.parent);
this.branchNode = true;
spanContainer.add(this);
}
public boolean hasNextBranch() {
return nextBranchNode != null;
}
public TraceSpanNode nextBranch() {
return nextBranchNode;
}
public void addBranch(TraceSpanNode branch) {
TraceSpanNode lastBranchNode = null;
while (hasNextBranch()) {
lastBranchNode = nextBranchNode;
}
lastBranchNode.nextBranchNode = branch;
branch.parent = this;
}
}

View File

@ -1,12 +1,10 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.entity;
package com.ai.cloud.skywalking.analysis.chainbuild.entity;
import com.ai.cloud.skywalking.analysis.categorize2chain.DBCallChainInfoDao;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.config.Config;
import com.ai.cloud.skywalking.analysis.chainbuild.DBCallChainInfoDao;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData;
import com.google.gson.Gson;
import org.apache.hadoop.hbase.client.Put;
import java.sql.SQLException;
@ -14,13 +12,13 @@ import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
public class ChainDetail {
public class CallChainDetail {
private boolean isNormal = true;
private String chainToken;
private Map<String, ChainNode> chainNodeMap = new HashMap<String, ChainNode>();
private String userId;
public ChainDetail(ChainInfo chainInfo, boolean isNormal) {
public CallChainDetail(ChainInfo chainInfo, boolean isNormal) {
chainToken = chainInfo.getCID();
for (ChainNode chainNode : chainInfo.getNodes()) {
chainNodeMap.put(chainNode.getTraceLevelId(), chainNode);

View File

@ -0,0 +1,18 @@
package com.ai.cloud.skywalking.analysis.chainbuild.entity;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
public class ChainNodeForSummary {
private String traceLevelId;
private String viewPointId;
public ChainNodeForSummary(ChainNode node) {
this.traceLevelId = node.getTraceLevelId();
this.viewPointId = node.getViewPoint();
}
public void summary(ChainNode node) {
}
}

View File

@ -1,323 +0,0 @@
package com.ai.cloud.skywalking.analysis.chainbuild.entity;
import java.util.List;
import com.ai.cloud.skywalking.analysis.chainbuild.exception.TraceSpanTreeNotFountException;
import com.ai.cloud.skywalking.analysis.chainbuild.exception.TraceSpanTreeSerializeException;
import com.ai.cloud.skywalking.analysis.chainbuild.util.StringUtil;
import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator;
import com.ai.cloud.skywalking.protocol.CallType;
import com.ai.cloud.skywalking.protocol.Span;
import com.google.gson.annotations.Expose;
public class TraceSpanNode {
protected TraceSpanNode prev = null;
protected TraceSpanNode next = null;
protected TraceSpanNode parent = null;
protected TraceSpanNode sub = null;
@Expose
protected String prevNodeRefToken = null;
@Expose
protected String nextNodeRefToken = null;
@Expose
protected String parentNodeRefToken = null;
@Expose
protected String subNodeRefToken = null;
@Expose
protected String nodeRefToken = null;
@Expose
protected boolean visualNode = true;
@Expose
protected String parentLevel;
@Expose
protected int levelId;
@Expose
protected String viewPointId = "";
@Expose
protected long cost = 0;
@Expose
protected long callTimes = 0;
/**
* 节点调用的状态<br/>
* 0成功<br/>
* 1异常<br/>
* 异常判断原则代码产生exception并且此exception不在忽略列表中
*/
@Expose
protected byte statusCode = 0;
/**
* 节点调用的错误堆栈<br/>
* 堆栈以JAVA的exception为主要判断依据
*/
@Expose
protected String exceptionStack;
/**
* 节点类型描述<br/>
* 已字符串的形式描述<br/>
* java,dubbo等
*/
@Expose
protected String spanType = "";
/**
* 节点调用过程中的业务字段<br/>
* 业务系统设置的订单号SQL语句等
*/
@Expose
protected String businessKey = "";
/**
* 节点调用所在的系统逻辑名称<br/>
* 由授权文件指定
*/
@Expose
protected String applicationId = "";
/**
* 是否为分支节点
*/
@Expose
protected boolean branchNode;
@Expose
protected TraceSpanNode nextBranchNode;
/**
* Warning: call this constructor ONLY by gson for deserialize
*/
public TraceSpanNode(){
}
public TraceSpanNode(TraceSpanNode parent, TraceSpanNode sub, TraceSpanNode prev, TraceSpanNode next, Span span, List<TraceSpanNode> spanContainer) {
this(parent, sub, prev, next, spanContainer);
this.visualNode = false;
this.parentLevel = span.getParentLevel();
this.levelId = span.getLevelId();
this.viewPointId = span.getViewPointId();
this.cost = span.getCost();
this.callTimes = 1;
this.statusCode = span.getStatusCode();
if (span.isReceiver()) {
this.exceptionStack = "server stack:";
} else {
this.exceptionStack = "client stack:";
}
this.exceptionStack += span.getExceptionStack();
this.spanType = span.getSpanType();
this.businessKey = span.getBusinessKey();
this.applicationId = span.getApplicationId();
//nodeToken : MD5(parentLevelId + levelId + viewpoint)
nodeRefToken = TokenGenerator.generateNodeToken(parentLevel + "-" + levelId + "-" + viewPointId);
}
protected TraceSpanNode(TraceSpanNode parent, TraceSpanNode sub, TraceSpanNode prev, TraceSpanNode next, List<TraceSpanNode> spanContainer) {
this.visualNode = true;
this.setParent(parent);
if (parent != null) {
parent.setSub(this);
}
this.setSub(sub);
if (sub != null) {
sub.setParent(this);
}
this.setPrev(prev);
if (prev != null) {
prev.setNext(this);
}
this.setNext(next);
if (next != null) {
next.setPrev(this);
}
spanContainer.add(this);
}
protected TraceSpanNode(TraceSpanNode parent, TraceSpanNode sub, TraceSpanNode prev, TraceSpanNode next, String parentLevelId, int levelId, List<TraceSpanNode> spanContainer) {
this(parent, sub, prev, next, spanContainer);
this.parentLevel = parentLevelId;
this.levelId = levelId;
this.callTimes = 0;
}
boolean hasNext() {
if (this.next != null) {
return true;
} else {
return false;
}
}
boolean hasSub() {
if (this.sub != null) {
return true;
} else {
return false;
}
}
void mergeSpan(Span span) {
if (CallType.convert(span.getCallType()) == CallType.ASYNC) {
this.cost += span.getCost();
}
if (span.getStatusCode() != 0 && !StringUtil.isBlank(span.getExceptionStack())) {
if (span.isReceiver()) {
this.exceptionStack += "server stack:";
} else {
this.exceptionStack += "client stack:";
}
this.exceptionStack += span.getExceptionStack();
}
}
public TraceSpanNode prev(TraceSpanTree tree) throws TraceSpanTreeNotFountException {
if(prev == null){
if(prevNodeRefToken == null){
throw new TraceSpanTreeNotFountException(getDesc() + " unexpected prev== null and prevNodeRefToken==null");
}else{
prev = tree.findNode(prevNodeRefToken);
}
}
return prev;
}
public TraceSpanNode next(TraceSpanTree tree) throws TraceSpanTreeNotFountException {
if(next == null){
if(nextNodeRefToken == null){
throw new TraceSpanTreeNotFountException(getDesc() + " unexpected next== null and nextNodeRefToken==null");
}else{
next = tree.findNode(nextNodeRefToken);
}
}
return next;
}
public TraceSpanNode parent(TraceSpanTree tree) throws TraceSpanTreeNotFountException {
if(parent == null){
if(parentNodeRefToken == null){
throw new TraceSpanTreeNotFountException(getDesc() + " unexpected parent== null and parentNodeRefToken==null");
}else{
parent = tree.findNode(parentNodeRefToken);
}
}
return parent;
}
public TraceSpanNode sub(TraceSpanTree tree) throws TraceSpanTreeNotFountException {
if(sub == null){
if(subNodeRefToken == null){
throw new TraceSpanTreeNotFountException(getDesc() + " unexpected sub== null and subNodeRefToken==null");
}else{
sub = tree.findNode(subNodeRefToken);
}
}
return sub;
}
public void setPrev(TraceSpanNode prev) {
this.prev = prev;
}
public void setNext(TraceSpanNode next) {
this.next = next;
}
public void setParent(TraceSpanNode parent) {
this.parent = parent;
}
public void setNextBranchNode(TraceSpanNode nextBranchNode){
this.nextBranchNode = nextBranchNode;
}
public void setSub(TraceSpanNode sub) {
this.sub = sub;
}
public boolean isVisualNode() {
return visualNode;
}
public String getParentLevel() {
return parentLevel;
}
public int getLevelId() {
return levelId;
}
public String getViewPointId() {
return viewPointId;
}
public long getCost() {
return cost;
}
public byte getStatusCode() {
return statusCode;
}
public String getExceptionStack() {
return exceptionStack;
}
public String getSpanType() {
return spanType;
}
public String getBusinessKey() {
return businessKey;
}
public String getApplicationId() {
return applicationId;
}
public String getNodeRefToken() throws TraceSpanTreeSerializeException {
if (StringUtil.isBlank(nodeRefToken)) {
throw new TraceSpanTreeSerializeException(getDesc() + " ref token is null.");
}
return nodeRefToken;
}
private String getDesc(){
return "Node[parentLevel=" + parentLevel + ", levelId=" + levelId + ", viewPointId=" + viewPointId + "]";
}
void serializeRef() throws TraceSpanTreeSerializeException {
if (prev != null) {
prevNodeRefToken = prev.getNodeRefToken();
}
if (parent != null) {
parentNodeRefToken = parent.getNodeRefToken();
}
if (next != null) {
nextNodeRefToken = next.getNodeRefToken();
}
if (sub != null) {
subNodeRefToken = sub.getNodeRefToken();
}
}
public boolean isBranchNode() {
return branchNode;
}
}

View File

@ -1,310 +0,0 @@
package com.ai.cloud.skywalking.analysis.chainbuild.entity;
import com.ai.cloud.skywalking.analysis.chainbuild.exception.BuildTraceSpanTreeException;
import com.ai.cloud.skywalking.analysis.chainbuild.exception.TraceSpanTreeNotFountException;
import com.ai.cloud.skywalking.analysis.chainbuild.exception.TraceSpanTreeSerializeException;
import com.ai.cloud.skywalking.analysis.chainbuild.util.StringUtil;
import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator;
import com.ai.cloud.skywalking.protocol.Span;
import com.google.gson.Gson;
import com.google.gson.GsonBuilder;
import com.google.gson.JsonObject;
import com.google.gson.JsonParser;
import com.google.gson.annotations.Expose;
import com.google.gson.reflect.TypeToken;
import org.apache.hadoop.io.Writable;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;
import java.util.*;
public class TraceSpanTree implements Writable {
private Logger logger = LoggerFactory.getLogger(TraceSpanTree.class);
@Expose
private String userId = null;
@Expose
private String cid;
@Expose
private TraceSpanNode treeRoot;
@Expose
private List<TraceSpanNode> spanContainer = new ArrayList<TraceSpanNode>();
private Map<String, TraceSpanNode> traceSpanNodeMap = new HashMap<String, TraceSpanNode>();
public TraceSpanTree() {
}
public String build(List<Span> spanList)
throws BuildTraceSpanTreeException, TraceSpanTreeNotFountException {
if (spanList.size() == 0) {
throw new BuildTraceSpanTreeException("spanList is empty.");
}
Collections.sort(spanList, new Comparator<Span>() {
@Override
public int compare(Span span1, Span span2) {
String span1TraceLevel = span1.getParentLevel() + "."
+ span1.getLevelId();
String span2TraceLevel = span2.getParentLevel() + "."
+ span2.getLevelId();
return span1TraceLevel.compareTo(span2TraceLevel);
}
});
Span span = spanList.get(0);
if (!StringUtil.isBlank(span.getUserId())) {
userId = span.getUserId();
} else {
throw new BuildTraceSpanTreeException(
"spanList[0] 's userId is null");
}
cid = generateCID(spanList.get(0));
treeRoot = new TraceSpanNode(null, null, null, null, spanList.get(0),
spanContainer);
if (spanList.size() > 1) {
for (int i = 1; i < spanList.size(); i++) {
this.build(spanList.get(i));
}
}
return cid;
}
private void build(Span span) throws BuildTraceSpanTreeException,
TraceSpanTreeNotFountException {
if (userId == null && !StringUtil.isBlank(span.getUserId())) {
userId = span.getUserId();
}
TraceSpanNode clientOrServerNode = findNodeAndCreateVisualNodeIfNess(
span.getParentLevel(), span.getLevelId());
if (clientOrServerNode != null) {
clientOrServerNode.mergeSpan(span);
}
if (span.getLevelId() > 0) {
TraceSpanNode foundNode = findNodeAndCreateVisualNodeIfNess(
span.getParentLevel(), span.getLevelId() - 1);
/**
* Create node between foundNode and foundNode.next(maybe
* foundNode.next == null)
*/
new TraceSpanNode(null, null, foundNode, foundNode.next(this),
span, spanContainer);
} else {
/**
* levelId=0 find for parent level if parentLevelId = 0.0.1 then
* find node[parentLevelId=0.0,levelId=1]
*/
String parentLevel = span.getParentLevel();
int idx = parentLevel.lastIndexOf("\\.");
if (idx < 0) {
throw new BuildTraceSpanTreeException("parentLevel="
+ parentLevel + " is unexpected.");
}
TraceSpanNode foundNode = findNodeAndCreateVisualNodeIfNess(
parentLevel.substring(0, idx),
Integer.parseInt(parentLevel.substring(idx + 1)));
/**
* Create sub node of using span data. FoundNode is parent node.
*/
new TraceSpanNode(foundNode, null, null, null, span, spanContainer);
}
}
private TraceSpanNode findNodeAndCreateVisualNodeIfNess(
String parentLevelId, int levelId)
throws TraceSpanTreeNotFountException {
String levelDesc = StringUtil.isBlank(parentLevelId) ? (levelId + "")
: (parentLevelId + "." + levelId);
String[] levelArray = levelDesc.split("\\.");
TraceSpanNode currentNode = treeRoot;
String contextParentLevelId = "";
for (String currentLevel : levelArray) {
int currentLevelInt = Integer.parseInt(currentLevel);
for (int i = 0; i < currentLevelInt; i++) {
if (currentNode.hasNext()) {
currentNode = currentNode.next(this);
} else {
// create visual next node
currentNode = new VisualTraceSpanNode(null, null,
currentNode, null, contextParentLevelId, i,
spanContainer);
}
}
contextParentLevelId = contextParentLevelId == "" ? ("" + currentLevelInt)
: (contextParentLevelId + "." + currentLevelInt);
if (currentNode.hasSub()) {
currentNode = currentNode.sub(this);
} else {
// create visual sub node
currentNode = new VisualTraceSpanNode(currentNode, null, null,
null, contextParentLevelId, 0, spanContainer);
}
}
return currentNode;
}
private String generateCID(Span level0Span)
throws BuildTraceSpanTreeException {
if (StringUtil.isBlank(level0Span.getParentLevel())
&& level0Span.getLevelId() == 0) {
StringBuilder chainTokenDesc = new StringBuilder();
chainTokenDesc.append(userId).append("_");
chainTokenDesc.append(level0Span.getViewPointId());
return getTSBySpanTraceId(level0Span) + "_" + TokenGenerator.generateCID(chainTokenDesc.toString());
} else {
throw new BuildTraceSpanTreeException("tid:"
+ level0Span.getTraceId() + " level0 span data is illegal");
}
}
private static String getTSBySpanTraceId(Span span)
throws BuildTraceSpanTreeException {
try {
Calendar calendar = Calendar.getInstance();
calendar.setTime(new Date(Long.parseLong(span.getTraceId().split(
"\\.")[2])));
return calendar.get(Calendar.YEAR) + "-" + (calendar.get(Calendar.MONTH) + 1);
} catch (Throwable t) {
throw new BuildTraceSpanTreeException("tid:" + span.getTraceId()
+ " is illegal.");
}
}
private void beforeSerialize() throws TraceSpanTreeSerializeException {
for (TraceSpanNode treeNode : spanContainer) {
treeNode.serializeRef();
}
}
public String serialize() throws TraceSpanTreeSerializeException {
beforeSerialize();
return new GsonBuilder().excludeFieldsWithoutExposeAnnotation()
.create().toJson(this);
}
TraceSpanNode findNode(String nodeRefToken)
throws TraceSpanTreeNotFountException {
if (traceSpanNodeMap.containsKey(nodeRefToken)) {
return traceSpanNodeMap.get(nodeRefToken);
} else {
throw new TraceSpanTreeNotFountException("nodeRefToken="
+ nodeRefToken + " not found.");
}
}
@Override
public void write(DataOutput out) throws IOException {
try {
out.write(serialize().getBytes());
} catch (TraceSpanTreeSerializeException e) {
logger.error("Failed to serialize Chain Id[" + cid + "]", e);
}
}
@Override
public void readFields(DataInput in) throws IOException {
String value = in.readLine();
try {
JsonObject jsonObject = (JsonObject) new JsonParser().parse(value);
userId = jsonObject.get("userId").getAsString();
cid = jsonObject.get("cid").getAsString();
treeRoot = new Gson().fromJson(jsonObject.get("treeRoot"),
TraceSpanNode.class);
spanContainer = new Gson().fromJson(
jsonObject.get("spanContainer"),
new TypeToken<List<TraceSpanNode>>() {
}.getType());
for (TraceSpanNode node : spanContainer) {
traceSpanNodeMap.put(node.getNodeRefToken(), node);
}
} catch (Exception e) {
logger.error("Failed to parse the value[" + value
+ "] to TraceSpanTree Object", e);
}
}
public TraceSpanNode getTreeRoot() {
return treeRoot;
}
public String getCid() {
return cid;
}
public void merge(TraceSpanTree spanTree) {
if (spanTree.getTreeRoot().hasNext()) {
SpanTreeMerger.merge(spanTree.getTreeRoot().next, treeRoot.next, spanContainer);
}
if (spanTree.getTreeRoot().hasSub()) {
SpanTreeMerger.merge(spanTree.getTreeRoot().sub, treeRoot.sub, spanContainer);
}
}
private static class SpanTreeMerger {
public static boolean merge(TraceSpanNode origin, TraceSpanNode dest, List<TraceSpanNode> spanContainer) {
boolean flag = false;
if (origin == null || dest == null) {
if (origin != null && dest == null) {
dest.parent.sub = origin;
origin.parent = dest.parent;
return true;
}
return true;
}
if (dest.isBranchNode()) {
BranchTraceSpanNode branchTraceSpanNode = (BranchTraceSpanNode) dest;
boolean branchFlag = false;
while (branchTraceSpanNode.hasNextBranch()) {
branchFlag = merge(origin, branchTraceSpanNode.nextBranch(), spanContainer);
if (branchFlag) {
break;
}
}
if (branchFlag) {
return true;
} else {
branchTraceSpanNode.addBranch(origin);
return false;
}
}
if (origin.isVisualNode() || dest.isVisualNode()) {
boolean nextFlag = merge(origin.next, dest.next, spanContainer);
boolean subFlag = merge(origin.sub, dest.sub, spanContainer);
if (subFlag && nextFlag) {
// 合并子树数据
} else {
new BranchTraceSpanNode(origin, dest, spanContainer);
}
flag = nextFlag && nextFlag;
} else {
if (origin.nodeRefToken.equals(dest.nodeRefToken)) {
// 合并子树数据
flag = true;
} else {
new BranchTraceSpanNode(origin, dest, spanContainer);
flag = false;
}
flag = flag && merge(origin.next, dest.next, spanContainer);
flag = flag && merge(origin.sub, dest.sub, spanContainer);
}
return flag;
}
}
}

View File

@ -1,24 +0,0 @@
package com.ai.cloud.skywalking.analysis.chainbuild.entity;
import java.util.List;
import com.ai.cloud.skywalking.analysis.chainbuild.util.StringUtil;
public class VisualTraceSpanNode extends TraceSpanNode {
protected VisualTraceSpanNode(TraceSpanNode parent, TraceSpanNode sub,
TraceSpanNode prev, TraceSpanNode next, String parentLevelId,
int levelId, List<TraceSpanNode> spanContainer) {
super(parent, sub, prev, next, parentLevelId, levelId, spanContainer);
/**set visual node token.<br/>
* for example: <br/>
* VisualNode[0.0]<br/>
* VisualNode[0.0.1]<br/>
* etc.<br/>
*/
nodeRefToken = "VisualNode[" + (StringUtil.isBlank(parentLevelId) ? "": nodeRefToken + ".") + levelId + "]";
}
}

View File

@ -1,13 +0,0 @@
package com.ai.cloud.skywalking.analysis.chainbuild.exception;
public class BuildTraceSpanTreeException extends Exception {
private static final long serialVersionUID = 5816399370389190974L;
public BuildTraceSpanTreeException(String msg){
super(msg);
}
public BuildTraceSpanTreeException(String msg, Exception cause){
super(msg, cause);
}
}

View File

@ -1,13 +0,0 @@
package com.ai.cloud.skywalking.analysis.chainbuild.exception;
public class TraceSpanTreeNotFountException extends Exception {
private static final long serialVersionUID = 5559441397011866237L;
public TraceSpanTreeNotFountException(String msg){
super(msg);
}
public TraceSpanTreeNotFountException(String msg, Exception cause){
super(msg, cause);
}
}

View File

@ -1,13 +0,0 @@
package com.ai.cloud.skywalking.analysis.chainbuild.exception;
public class TraceSpanTreeSerializeException extends Exception {
private static final long serialVersionUID = 7857716041262993579L;
public TraceSpanTreeSerializeException(String msg){
super(msg);
}
public TraceSpanTreeSerializeException(String msg, Exception cause){
super(msg, cause);
}
}

View File

@ -1,4 +1,4 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.filter;
package com.ai.cloud.skywalking.analysis.chainbuild.filter;
import java.io.IOException;
import java.util.HashMap;

View File

@ -1,8 +1,8 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.filter;
package com.ai.cloud.skywalking.analysis.chainbuild.filter;
import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter;
import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter;
public abstract class SpanNodeProcessFilter {

View File

@ -0,0 +1,16 @@
package com.ai.cloud.skywalking.analysis.chainbuild.filter.impl;
import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry;
import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter;
public class AppendBusinessKeyFilter extends SpanNodeProcessFilter {
@Override
public void doFilter(SpanEntry spanEntry, ChainNode node, SubLevelSpanCostCounter costMap) {
node.setViewPoint(node.getViewPoint() + spanEntry.getBusinessKey());
this.doNext(spanEntry, node, costMap);
}
}

View File

@ -1,9 +1,9 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl;
package com.ai.cloud.skywalking.analysis.chainbuild.filter.impl;
import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry;
import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter;
import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry;
import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter;
public class CopyAttrFilter extends SpanNodeProcessFilter {

View File

@ -1,9 +1,9 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl;
package com.ai.cloud.skywalking.analysis.chainbuild.filter.impl;
import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry;
import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter;
import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry;
import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter;
public class ProcessCostTimeFilter extends SpanNodeProcessFilter {
@Override

View File

@ -1,9 +1,9 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl;
package com.ai.cloud.skywalking.analysis.chainbuild.filter.impl;
import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry;
import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter;
import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry;
import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter;
public class ReplaceAddressFilter extends SpanNodeProcessFilter {

View File

@ -1,9 +1,9 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl;
package com.ai.cloud.skywalking.analysis.chainbuild.filter.impl;
import com.ai.cloud.skywalking.analysis.categorize2chain.SpanEntry;
import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainNode;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.SubLevelSpanCostCounter;
import com.ai.cloud.skywalking.analysis.chainbuild.SpanEntry;
import com.ai.cloud.skywalking.analysis.chainbuild.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainNode;
import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter;
import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator;
public class TokenGenerateFilter extends SpanNodeProcessFilter {

View File

@ -0,0 +1,46 @@
package com.ai.cloud.skywalking.analysis.chainbuild.po;
import com.ai.cloud.skywalking.analysis.chainbuild.entity.ChainNodeForSummary;
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 CallChainTreeNode {
private String traceLevelId;
// key: nodeToken
private Map<String, ChainNodeForSummary> chainNodeContainer;
public CallChainTreeNode(ChainNode node) {
this.traceLevelId = node.getTraceLevelId();
chainNodeContainer.put(node.getNodeToken(), new ChainNodeForSummary(node));
}
public CallChainTreeNode(String originData) {
JsonObject jsonObject = (JsonObject) new JsonParser().parse(originData);
traceLevelId = jsonObject.get("traceLevelId").getAsString();
chainNodeContainer = new Gson().fromJson(jsonObject.get("chainNodeContainer").getAsString(),
new TypeToken<Map<String, ChainNodeForSummary>>() {
}.getType());
}
public void mergeIfNess(ChainNode node) {
if (!chainNodeContainer.containsKey(node.getNodeToken())) {
chainNodeContainer.put(node.getNodeToken(), new ChainNodeForSummary(node));
}
}
public void summary(ChainNode node) {
ChainNodeForSummary chainNode = chainNodeContainer.get(node.getNodeToken());
chainNode.summary(node);
}
public String getTraceLevelId() {
return traceLevelId;
}
}

View File

@ -1,11 +1,11 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.po;
package com.ai.cloud.skywalking.analysis.chainbuild.po;
import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator;
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.hbase.client.Put;
import org.apache.hadoop.io.Writable;
import java.io.DataInput;
@ -21,6 +21,7 @@ public class ChainInfo implements Writable {
private String userId = null;
private ChainNode firstChainNode;
private long startDate;
private String chainToken;
public ChainInfo(String userId) {
super();
@ -61,15 +62,6 @@ public class ChainInfo implements Writable {
}
}
public void generateChainToken() {
StringBuilder chainTokenDesc = new StringBuilder();
for (ChainNode node : nodes) {
chainTokenDesc.append(node.getParentLevelId() + "."
+ node.getLevelId() + "-" + node.getNodeToken() + ";");
}
this.cid = TokenGenerator.generateCID(chainTokenDesc.toString());
}
public ChainStatus getChainStatus() {
return chainStatus;
}
@ -92,6 +84,7 @@ public class ChainInfo implements Writable {
&& chainNode.getLevelId() == 0) {
firstChainNode = chainNode;
startDate = chainNode.getStartDate();
cid = firstChainNode.getViewPoint();
}
}
@ -107,6 +100,14 @@ public class ChainInfo implements Writable {
this.userId = userId;
}
public String getChainToken() {
return chainToken;
}
public void saveToHBase(Put put) {
}
public enum ChainStatus {
NORMAL('N'), ABNORMAL('A');
private char value;

View File

@ -1,4 +1,4 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.po;
package com.ai.cloud.skywalking.analysis.chainbuild.po;
import com.google.gson.GsonBuilder;
import com.google.gson.annotations.Expose;

View File

@ -1,19 +1,13 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.util;
package com.ai.cloud.skywalking.analysis.chainbuild.util;
import com.ai.cloud.skywalking.analysis.categorize2chain.*;
import com.ai.cloud.skywalking.analysis.categorize2chain.entity.CategorizedChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainNodeSpecificTimeWindowSummary;
import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainRelationship;
import com.ai.cloud.skywalking.analysis.categorize2chain.entity.ChainSpecificTimeWindowSummary;
import com.ai.cloud.skywalking.analysis.categorize2chain.entity.UncategorizeChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.chain2summary.ChainRelationship4Search;
import com.ai.cloud.skywalking.analysis.chain2summary.entity.*;
import com.ai.cloud.skywalking.analysis.chainbuild.CallChainTree;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.chainbuild.po.CallChainTreeNode;
import com.ai.cloud.skywalking.analysis.config.Config;
import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData;
import com.google.gson.Gson;
import com.google.gson.reflect.TypeToken;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hbase.*;
import org.apache.hadoop.hbase.client.*;
@ -36,30 +30,35 @@ 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_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_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);
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_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);
createTableIfNeed(HBaseTableMetaData.TABLE_MERGED_CHAIN_DETAIL.TABLE_NAME,
HBaseTableMetaData.TABLE_MERGED_CHAIN_DETAIL.COLUMN_FAMILY_NAME);
createTableIfNeed(HBaseTableMetaData.TABLE_CALL_CHAIN_TREE_ID_AND_CID_MAPPING.TABLE_NAME,
HBaseTableMetaData.TABLE_CALL_CHAIN_TREE_ID_AND_CID_MAPPING.COLUMN_FAMILY_NAME);
} catch (IOException e) {
logger.error("Create tables failed", e);
}
@ -115,52 +114,6 @@ public class HBaseUtil {
return true;
}
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));
Result r = table.get(g);
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<UncategorizeChainInfo> uncategorizeChainInfoList = new Gson().fromJson(Bytes.toString(cell.getValueArray(),
cell.getValueOffset(), cell.getValueLength()),
new TypeToken<List<UncategorizeChainInfo>>() {
}.getType());
chainRelate.addUncategorizeChain(uncategorizeChainInfoList);
} else {
chainRelate.addCategorizeChain(qualifierName, new CategorizedChainInfo(
Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength())
));
}
}
}
return chainRelate;
}
public static ChainSpecificTimeWindowSummary selectChainSummaryResult(String key) throws IOException {
ChainSpecificTimeWindowSummary result = null;
Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CHAIN_ONE_MINUTE_SUMMARY_EXCLUDE_RELATIONSHIP.TABLE_NAME));
Get g = new Get(Bytes.toBytes(key));
Result r = table.get(g);
if (r.rawCells().length == 0) {
return null;
}
result = new ChainSpecificTimeWindowSummary();
for (Cell cell : r.rawCells()) {
if (cell.getValueArray().length > 0)
result.addNodeSummaryResult(new ChainNodeSpecificTimeWindowSummary(Bytes.toString(cell.getValueArray(),
cell.getValueOffset(), cell.getValueLength())));
}
return result;
}
public static void saveChainRelationship(Put put) throws IOException {
Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CALL_CHAIN_RELATIONSHIP.TABLE_NAME));
@ -194,44 +147,6 @@ 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<UncategorizeChainInfo> uncategorizeChainInfoList = new Gson().fromJson(Bytes.toString(cell.getValueArray(),
cell.getValueOffset(), cell.getValueLength()),
new TypeToken<List<UncategorizeChainInfo>>() {
}.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));
@ -358,4 +273,46 @@ public class HBaseUtil {
}
}
}
public static CallChainTree loadMergedCallChain(String callEntrance) throws IOException {
CallChainTree result = null;
Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_MERGED_CHAIN_DETAIL.TABLE_NAME));
Get g = new Get(Bytes.toBytes(callEntrance));
Result r = table.get(g);
if (r.rawCells().length == 0) {
return null;
}
result = new CallChainTree(callEntrance);
for (Cell cell : r.rawCells()) {
if (cell.getValueArray().length > 0)
result.addMergedChainNode(new CallChainTreeNode(Bytes.toString(cell.getValueArray(),
cell.getValueOffset(), cell.getValueLength())));
}
return result;
}
public static void saveMergedCallChain(CallChainTree callChainTree) {
// save
// save relationship
}
public static List<String> loadHasBeenMergeChainIds(String topoId) throws IOException {
List<String> result = new ArrayList<String>();
Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CALL_CHAIN_TREE_ID_AND_CID_MAPPING.TABLE_NAME));
Get g = new Get(Bytes.toBytes(topoId));
Result r = table.get(g);
if (r.rawCells().length == 0) {
return null;
}
for (Cell cell : r.rawCells()) {
if (cell.getValueArray().length > 0) {
List<String> hasBeenMergedCIds = new Gson().fromJson("",
new TypeToken<List<String>>() {
}.getType());
result.addAll(hasBeenMergedCIds);
}
}
return result;
}
}

View File

@ -1,22 +0,0 @@
package com.ai.cloud.skywalking.analysis.chainbuild.util;
public class StringUtil {
public static boolean isBlank(String str){
if(str == null || str == "" || str.trim() == ""){
return true;
}else{
return false;
}
}
public static boolean equal(String str1, String str2){
if(str1 == null){
str1 = "";
}
if(str2 == null){
str2 = "";
}
return str1.trim().equals(str2.trim());
}
}

View File

@ -1,4 +1,4 @@
package com.ai.cloud.skywalking.analysis.categorize2chain.util;
package com.ai.cloud.skywalking.analysis.chainbuild.util;
import java.util.HashMap;
import java.util.Map;

View File

@ -12,13 +12,13 @@ public class TokenGenerator {
private TokenGenerator() {
//Non
}
public static String generateCID(String originData) {
return "CID_" + generate(originData);
return "CID_" + generate(originData);
}
public static String generateNodeToken(String originData){
return "C_NID_" + generate(originData);
return "C_NID_" + generate(originData);
}
private static String generate(String originData) {
@ -41,4 +41,4 @@ public class TokenGenerator {
}
return result.toString().toUpperCase();
}
}
}

View File

@ -1,31 +1,25 @@
package com.ai.cloud.skywalking.analysis.chainbuild.util;
/**
* 版本识别器
*
* @author wusheng
*
*/
public class VersionIdentifier {
/**
* 根据tid识别数据是否可分析<br/>
* 目前允许分析所有1.x的版本号
*
* @param tid
* @return
*/
public static boolean enableAnaylsis(String tid){
if(tid != null){
String[] tidSections = tid.split("\\.");
if(tidSections.length == 7){
String version = tidSections[0];
String subVersion = tidSections[1];
if("1".equals(version) && subVersion.length() > 0){
return true;
}
}
}
return false;
}
}
/**
* 根据tid识别数据是否可分析<br/>
* 目前允许分析所有1.x的版本号
*
* @param tid
* @return
*/
public static boolean enableAnaylsis(String tid) {
if (tid != null) {
String[] tidSections = tid.split("\\.");
if (tidSections.length == 7) {
String version = tidSections[0];
String subVersion = tidSections[1];
if ("1".equals(version) && subVersion.length() > 0) {
return true;
}
}
}
return false;
}
}

View File

@ -43,8 +43,6 @@ public class HBaseTableMetaData {
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";
}
/**
@ -101,4 +99,22 @@ public class HBaseTableMetaData {
public static final String COLUMN_FAMILY_NAME = "chain_summary";
}
/**
* 用于存放已经合并的调用链的信息
*
* @author zhangxin
*/
public final static class TABLE_MERGED_CHAIN_DETAIL {
public static final String TABLE_NAME = "sw-merged-chain-detail";
public static final String COLUMN_FAMILY_NAME = "chain_detail";
}
public final static class TABLE_CALL_CHAIN_TREE_ID_AND_CID_MAPPING {
public static final String TABLE_NAME = "sw-topologyId-cid-mapping";
public static final String COLUMN_FAMILY_NAME = "sw-topologyId-cid-mapping";
}
}