增加Redis辅助MapReduce

This commit is contained in:
ascrutae 2016-05-15 07:08:15 +08:00
parent 585ddb47d7
commit b665322deb
9 changed files with 209 additions and 49 deletions

View File

@ -59,6 +59,11 @@
<artifactId>log4j-core</artifactId>
<version>2.2</version>
</dependency>
<dependency>
<groupId>redis.clients</groupId>
<artifactId>jedis</artifactId>
<version>2.8.1</version>
</dependency>
</dependencies>
<build>

View File

@ -1,18 +1,12 @@
package com.ai.cloud.skywalking.analysis;
import java.io.IOException;
import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.Date;
import com.ai.cloud.skywalking.analysis.chainbuild.ChainBuildMapper;
import com.ai.cloud.skywalking.analysis.chainbuild.ChainBuildReducer;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.config.Config;
import com.ai.cloud.skywalking.analysis.config.ConfigInitializer;
import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.conf.Configured;
import org.apache.hadoop.hbase.HConstants;
import org.apache.hadoop.hbase.client.Scan;
import org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil;
import org.apache.hadoop.io.Text;
@ -24,8 +18,10 @@ import org.apache.hadoop.util.ToolRunner;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.ai.cloud.skywalking.analysis.config.Config;
import com.ai.cloud.skywalking.analysis.config.ConfigInitializer;
import java.io.IOException;
import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.Date;
public class AnalysisServerDriver extends Configured implements Tool {
@ -36,10 +32,10 @@ public class AnalysisServerDriver extends Configured implements Tool {
String analysisMode = System.getenv("skywalking.analysis.mode");
if ("rewrite".equalsIgnoreCase(analysisMode)){
if ("rewrite".equalsIgnoreCase(analysisMode)) {
logger.info("Skywalking analysis mode will switch to [REWRITE] mode");
Config.AnalysisServer.IS_ACCUMULATE_MODE = false;
}else{
} else {
logger.info("Skywalking analysis mode will switch to [ACCUMULATE] mode");
Config.AnalysisServer.IS_ACCUMULATE_MODE = true;
}
@ -57,7 +53,7 @@ public class AnalysisServerDriver extends Configured implements Tool {
conf.set("hbase.zookeeper.quorum", Config.HBase.ZK_QUORUM);
conf.set("hbase.zookeeper.property.clientPort", Config.HBase.ZK_CLIENT_PORT);
//-XX:+UseParallelGC -XX:ParallelGCThreads=4 -XX:GCTimeRatio=10 -XX:YoungGenerationSizeIncrement=20 -XX:TenuredGenerationSizeIncrement=20 -XX:AdaptiveSizeDecrementScaleFactor=2
conf.set("mapred.child.java.opts",Config.MapReduce.JAVA_OPTS);
conf.set("mapred.child.java.opts", Config.MapReduce.JAVA_OPTS);
String[] otherArgs = new GenericOptionsParser(conf, args).getRemainingArgs();
if (otherArgs.length != 2) {
System.err.println("Usage: com.ai.cloud.skywalking.analysis.AnalysisServerDriver yyyy-MM-dd/HH:mm:ss yyyy-MM-dd/HH:mm:ss");
@ -83,6 +79,9 @@ public class AnalysisServerDriver extends Configured implements Tool {
Date startDate = simpleDateFormat.parse(args[0]);
Date endDate = simpleDateFormat.parse(args[1]);
Scan scan = new Scan();
//scan.setMaxVersions();
scan.setBatch(2001);
scan.setMaxVersions();
scan.setTimeRange(startDate.getTime(), endDate.getTime());
return scan;
}

View File

@ -6,10 +6,8 @@ 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.po.SummaryType;
import com.ai.cloud.skywalking.analysis.chainbuild.util.HBaseUtil;
import com.ai.cloud.skywalking.analysis.chainbuild.util.SubLevelSpanCostCounter;
import com.ai.cloud.skywalking.analysis.chainbuild.util.TokenGenerator;
import com.ai.cloud.skywalking.analysis.chainbuild.util.VersionIdentifier;
import com.ai.cloud.skywalking.analysis.chainbuild.util.*;
import com.ai.cloud.skywalking.analysis.config.Config;
import com.ai.cloud.skywalking.analysis.config.ConfigInitializer;
import com.ai.cloud.skywalking.protocol.Span;
import com.ai.cloud.skywalking.util.SpanLevelIdComparators;
@ -48,6 +46,7 @@ public class ChainBuildMapper extends TableMapper<Text, Text> {
return;
}
RedisUtil.autoIncrement(Config.Redis.MAPPER_COUNT_KEY);
List<Span> spanList = new ArrayList<Span>();
ChainInfo chainInfo = null;
try {
@ -61,6 +60,12 @@ public class ChainBuildMapper extends TableMapper<Text, Text> {
+ Bytes.toString(key.get()) + "] has no span data.");
}
if (spanList.size() > 2000) {
throw new Tid2CidECovertException("tid["
+ Bytes.toString(key.get()) + "] node size has over 2000.");
}
chainInfo = spanToChainInfo(Bytes.toString(key.get()), spanList);
logger.debug("convert tid[" + Bytes.toString(key.get())
+ "] to chain with cid[" + chainInfo.getCID() + "].");
@ -116,9 +121,11 @@ public class ChainBuildMapper extends TableMapper<Text, Text> {
+ ":" + chainInfo.getCallEntrance()),
new Text(new Gson().toJson(chainInfo)));
}
RedisUtil.autoIncrement(Config.Redis.SUCCESS_MAPPER_COUNT_KEY);
} catch (Exception e) {
logger.error("Failed to mapper call chain[" + key.toString() + "]",
e);
RedisUtil.autoIncrement(Config.Redis.FAILED_MAPPER_COUNT_KEY);
}
}
@ -155,4 +162,10 @@ public class ChainBuildMapper extends TableMapper<Text, Text> {
}
return spanEntryMap;
}
private void clearData(){
RedisUtil.clearData(Config.Redis.MAPPER_COUNT_KEY);
RedisUtil.clearData(Config.Redis.FAILED_MAPPER_COUNT_KEY);
RedisUtil.clearData(Config.Redis.SUCCESS_MAPPER_COUNT_KEY);
}
}

View File

@ -0,0 +1,69 @@
package com.ai.cloud.skywalking.analysis.chainbuild.util;
import com.ai.cloud.skywalking.analysis.config.Config;
import org.apache.commons.pool2.impl.GenericObjectPoolConfig;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisPool;
/**
* Created by xin on 16-5-13.
*/
public class RedisUtil {
private static Logger logger = LogManager.getLogger(RedisUtil.class);
private static JedisPool jedisPool;
private static boolean turn_on = true;
static {
try {
GenericObjectPoolConfig genericObjectPoolConfig = new GenericObjectPoolConfig();
jedisPool = new JedisPool(genericObjectPoolConfig, Config.Redis.HOST, Config.Redis.PORT);
} catch (Exception e) {
logger.error("Failed to create jedis pool", e);
turn_on = false;
}
}
public static void autoIncrement(String key) {
if (!turn_on) {
return;
}
Jedis jedis = null;
try {
jedis = jedisPool.getResource();
jedis.incrBy(key, 1);
} catch (Exception e) {
logger.error("Failed to auto increment .", e);
} finally {
if (jedis != null) {
jedis.close();
}
}
}
public static void clearData(String key) {
if (!turn_on) {
return;
}
Jedis jedis = null;
try {
jedis = jedisPool.getResource();
jedis.setnx(key, "0");
} catch (Exception e) {
logger.error("Failed to auto increment .", e);
} finally {
if (jedis != null) {
jedis.close();
}
}
}
}

View File

@ -26,11 +26,24 @@ public class Config {
public static String FILTER_PACKAGE_NAME;
}
public static class AnalysisServer{
public static class AnalysisServer {
public static boolean IS_ACCUMULATE_MODE = true;
}
public static class MapReduce{
public static class MapReduce {
public static String JAVA_OPTS = "-Xmx200m";
}
public static class Redis {
public static String HOST = "127.0.0.1";
public static int PORT = 6379;
public static String MAPPER_COUNT_KEY = "ANALYSIS_TOTAL_SIZE";
public static String SUCCESS_MAPPER_COUNT_KEY = "ANALYSIS_SUCCESS_TOTAL_SIZE";
public static String FAILED_MAPPER_COUNT_KEY = "ANALYSIS_FAILED_TOTAL_SIZE";
}
}

View File

@ -15,6 +15,12 @@ filter.filter_package_name=com.ai.cloud.skywalking.analysis.chainbuild.filter.im
chainnodesummary.interval=1
reducer.reducer_number=4
reducer.reducer_number=1
mapreduce.java_opts=-Xmx768m
mapreduce.java_opts=-Xmx768m
redis.host=10.1.55.11
redis.port=6379
redis.mapper_count_key=ANALYSIS_TOTAL_SIZE

View File

@ -5,6 +5,7 @@ import com.ai.cloud.skywalking.analysis.chainbuild.ChainBuildMapper;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo;
import com.ai.cloud.skywalking.analysis.config.ConfigInitializer;
import com.ai.cloud.skywalking.analysis.config.HBaseTableMetaData;
import com.ai.cloud.skywalking.analysis.mapper.util.Convert;
import com.ai.cloud.skywalking.analysis.mapper.util.HBaseUtils;
import com.ai.cloud.skywalking.protocol.Span;
import org.apache.hadoop.hbase.Cell;
@ -22,19 +23,27 @@ public class CallChainMapperTest {
private static Connection connection = HBaseUtils.getConnection();
private static SimpleDateFormat simpleDateFormat = new SimpleDateFormat("yyyy-MM-dd/HH:mm:ss");
private static final int EXCEPT_CHAIN_INFO_SIZE = 124985;
public static void main(String[] args) throws Exception {
ConfigInitializer.initialize();
SimpleDateFormat simpleDateFormat = new SimpleDateFormat("yyyy-MM-dd/HH:mm:ss");
Date startDate = simpleDateFormat.parse("2016-04-22/23:57:03");
Date endDate = simpleDateFormat.parse("2016-05-02/23:47:03");
Connection connection = HBaseUtils.getConnection();
Table table = connection.getTable(TableName.valueOf
(HBaseTableMetaData.TABLE_CALL_CHAIN.TABLE_NAME));
Scan scan = new Scan();
//2016-04-13/16:59:24 to 2016-05-13/16:49:24
Date startDate = simpleDateFormat.parse("2016-04-13/16:59:24");
Date endDate = simpleDateFormat.parse("2016-05-13/16:49:24");
scan.setBatch(2001);
scan.setTimeRange(startDate.getTime(), endDate.getTime());
Table table = connection.getTable(TableName.valueOf(HBaseTableMetaData.TABLE_CALL_CHAIN.TABLE_NAME));
ResultScanner result = table.getScanner(scan);
int count = 0;
for (Result result1 : result) {
count++;
List<ChainInfo> chainInfos = Convert.convert(table.getScanner(scan));
if (EXCEPT_CHAIN_INFO_SIZE != chainInfos.size()) {
System.out.println("except size :" + EXCEPT_CHAIN_INFO_SIZE + " accutal size:" + chainInfos.size());
System.exit(-1);
}
System.out.println(count);
}
}

View File

@ -1,47 +1,47 @@
package com.ai.cloud.skywalking.analysis.mapper;
import com.ai.cloud.skywalking.analysis.chainbuild.ChainBuildMapper;
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.ai.cloud.skywalking.analysis.mapper.util.Convert;
import com.ai.cloud.skywalking.analysis.mapper.util.HBaseUtils;
import com.ai.cloud.skywalking.protocol.Span;
import org.apache.hadoop.hbase.Cell;
import org.apache.hadoop.hbase.TableName;
import org.apache.hadoop.hbase.client.*;
import org.apache.hadoop.hbase.client.Connection;
import org.apache.hadoop.hbase.client.Scan;
import org.apache.hadoop.hbase.client.Table;
import org.apache.hadoop.hbase.filter.CompareFilter;
import org.apache.hadoop.hbase.filter.Filter;
import org.apache.hadoop.hbase.filter.SingleColumnValueFilter;
import org.apache.hadoop.hbase.util.Bytes;
import java.io.IOException;
import java.util.ArrayList;
import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
public class ValidateMonthSummaryResult {
public static void main(String[] args) throws IOException {
public class ValidateSpecialTimeSummaryResult {
private static SimpleDateFormat simpleDateFormat = new SimpleDateFormat("yyyy-MM-dd/HH:mm:ss");
public static void main(String[] args) throws IOException, ParseException {
Connection connection = HBaseUtils.getConnection();
Table table = connection.getTable(TableName.valueOf
(HBaseTableMetaData.TABLE_CALL_CHAIN.TABLE_NAME));
Scan scan = new Scan();
Date startDate = simpleDateFormat.parse("2016-04-13/16:59:24");
Date endDate = simpleDateFormat.parse("2016-05-13/16:49:24");
scan.setMaxVersions();
scan.setTimeRange(startDate.getTime(), endDate.getTime());
Filter filter = new SingleColumnValueFilter(HBaseTableMetaData.TABLE_CALL_CHAIN.FAMILY_NAME.getBytes(),
"0-S".getBytes(), CompareFilter.CompareOp.EQUAL,
"http://hire.asiainfo.com/Aisse-Mobile-Web/aisseWorkUser/queryProsonLoding".getBytes());
scan.setFilter(filter);
ResultScanner resultScanner = table.getScanner(scan);
List<ChainInfo> chainInfos = new ArrayList<ChainInfo>();
List<ChainInfo> chainInfos = Convert.convert(table.getScanner(scan));
for (Result result : resultScanner) {
List<Span> spanList = new ArrayList<Span>();
for (Cell cell : result.rawCells()) {
spanList.add(new Span(Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength())));
}
chainInfos.add(ChainBuildMapper.spanToChainInfo(Bytes.toString(result.getRow()), spanList));
}
Map<String, Integer> totalCallSize = new HashMap<String, Integer>();
Map<String, Float> totalCallTimes = new HashMap<>();

View File

@ -0,0 +1,46 @@
package com.ai.cloud.skywalking.analysis.mapper.util;
import com.ai.cloud.skywalking.analysis.chainbuild.ChainBuildMapper;
import com.ai.cloud.skywalking.analysis.chainbuild.po.ChainInfo;
import com.ai.cloud.skywalking.protocol.Span;
import org.apache.hadoop.hbase.Cell;
import org.apache.hadoop.hbase.client.Result;
import org.apache.hadoop.hbase.client.ResultScanner;
import org.apache.hadoop.hbase.util.Bytes;
import java.util.*;
/**
* Created by xin on 16-5-13.
*/
public class Convert {
private static final Map<String, String> traceIds = new HashMap<>();
public static List<ChainInfo> convert(ResultScanner resultScanner) {
List<ChainInfo> chainInfos = new ArrayList<ChainInfo>();
for (Result result : resultScanner) {
try {
List<Span> spanList = new ArrayList<Span>();
for (Cell cell : result.rawCells()) {
spanList.add(new Span(Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength())));
}
if (spanList.size() == 0 || spanList.size() > 2000) {
continue;
}
ChainInfo chainInfo = ChainBuildMapper.spanToChainInfo(Bytes.toString(result.getRow()), spanList);
chainInfos.add(chainInfo);
traceIds.put(spanList.get(0).getTraceId(), chainInfo.getCID());
} catch (Exception e) {
continue;
}
}
System.out.println(traceIds.size());
return chainInfos;
}
}