修复Reduce的bug

This commit is contained in:
ascrutae 2016-01-24 09:45:17 +08:00
parent daf6eab08a
commit ebd75d6ef9
7 changed files with 34 additions and 30 deletions

View File

@ -15,6 +15,7 @@ import org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.mapreduce.lib.output.NullOutputFormat;
import org.apache.hadoop.util.GenericOptionsParser;
import org.apache.hadoop.util.Tool;
import org.apache.hadoop.util.ToolRunner;
@ -55,13 +56,10 @@ public class AnalysisServerDriver extends Configured implements Tool {
TableMapReduceUtil.initTableMapperJob(Config.HBase.TABLE_CALL_CHAIN, scan, Categorize2ChainMapper.class,
Text.class, ChainInfo.class, job);
int regions = MetaTableAccessor.getRegionCount(conf, TableName.valueOf(Config.HBase.TABLE_CHAIN_SUMMARY));
if (regions == 0) {
regions = 1;
}
job.setReducerClass(Categorize2ChainReducer.class);
job.setNumReduceTasks(regions);
FileOutputFormat.setOutputPath(job, new Path("/tmp/mr/mySummaryFile"));
job.setNumReduceTasks(Config.Reducer.REDUCER_NUMBER);
job.setOutputFormatClass(NullOutputFormat.class);
return job.waitForCompletion(true) ? 0 : 1;
}

View File

@ -4,6 +4,7 @@ import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessC
import com.ai.cloud.skywalking.analysis.categorize2chain.filter.SpanNodeProcessFilter;
import com.ai.cloud.skywalking.analysis.categorize2chain.model.ChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.model.ChainNode;
import com.ai.cloud.skywalking.analysis.config.ConfigInitializer;
import com.ai.cloud.skywalking.analysis.util.HBaseUtil;
import com.ai.cloud.skywalking.protocol.Span;
import org.apache.hadoop.hbase.Cell;
@ -24,6 +25,7 @@ public class Categorize2ChainMapper extends TableMapper<Text, ChainInfo> {
@Override
protected void map(ImmutableBytesWritable key, Result value, Context context) throws IOException,
InterruptedException {
ConfigInitializer.initialize();
List<Span> spanList = new ArrayList<Span>();
ChainInfo chainInfo = null;
try {
@ -32,7 +34,7 @@ public class Categorize2ChainMapper extends TableMapper<Text, ChainInfo> {
spanList.add(span);
}
chainInfo = spanToChainInfo(key.toString(), spanList);
chainInfo = spanToChainInfo(Bytes.toString(key.get()), spanList);
logger.info("Success convert span to chain info...." + chainInfo.getCID());
context.write(new Text(chainInfo.getUserId() + ":" + chainInfo.getEntranceNodeToken()), chainInfo);
} catch (Exception e) {

View File

@ -3,6 +3,7 @@ package com.ai.cloud.skywalking.analysis.categorize2chain;
import java.io.IOException;
import java.util.Iterator;
import com.ai.cloud.skywalking.analysis.config.ConfigInitializer;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;
@ -17,6 +18,7 @@ public class Categorize2ChainReducer extends Reducer<Text, ChainInfo, Text, IntW
@Override
protected void reduce(Text key, Iterable<ChainInfo> values, Context context) throws IOException, InterruptedException {
ConfigInitializer.initialize();
int totalCount = reduceAction(key.toString(), values.iterator());
context.write(new Text(key.toString()), new IntWritable(totalCount));
}

View File

@ -1,18 +1,20 @@
package com.ai.cloud.skywalking.analysis.config;
public class Config {
public static class Reducer{
public static int REDUCER_NUMBER = 1;
}
public static class HBase {
public static String TRACE_DETAIL_FAMILY_COLUMN = "chain_detail";
public static String CHAIN_SUMMARY_COLUMN_FAMILY = "chain_summary";
public static String TRACE_INFO_COLUMN_FAMILY = "trace_info";
public static String ZK_QUORUM = "10.1.235.197,10.1.235.198,10.1.235.199";
public static String ZK_QUORUM;
public static String ZK_CLIENT_PORT = "29181";
public static String ZK_CLIENT_PORT;
public static String TRACE_INFO_TABLE_NAME = "trace-info";
@ -35,20 +37,20 @@ public class Config {
public static class MySql {
public static String URL = "jdbc:mysql://10.1.228.202:31316/test";
public static String URL;
public static String USERNAME = "devrdbusr21";
public static String USERNAME;
public static String PASSWORD = "devrdbusr21";
public static String PASSWORD;
public static String DRIVER_CLASS = "com.mysql.jdbc.Driver";
}
public static class Filter {
public static String FILTER_PACKAGE_NAME = "com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl";
public static String FILTER_PACKAGE_NAME;
}
public class ChainNodeSummary {
public static final long INTERVAL = 5L;
public static class ChainNodeSummary {
public static long INTERVAL = 5L;
}
}

View File

@ -1,32 +1,31 @@
package com.ai.cloud.skywalking.analysis.config;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.io.InputStream;
import java.lang.reflect.Field;
import java.lang.reflect.Modifier;
import java.util.LinkedList;
import java.util.Properties;
import java.util.logging.Level;
import java.util.logging.Logger;
public class ConfigInitializer {
private static Logger logger = Logger.getLogger(ConfigInitializer.class.getName());
private static Logger logger = LoggerFactory.getLogger(ConfigInitializer.class.getName());
public static void initialize() {
InputStream inputStream = Thread.currentThread()
.getContextClassLoader().getResourceAsStream("/analysis.conf");
//InputStream inputStream = ConfigInitializer.class.getResourceAsStream("/config.properties");
InputStream inputStream = ConfigInitializer.class.getResourceAsStream("/analysis.conf");
if (inputStream == null) {
logger.log(Level.ALL, "No provider sky-walking certification documents, sky-walking api auto shutdown.");
logger.error("No provider sky-walking certification documents, sky-walking api auto shutdown.");
} else {
try {
Properties properties = new Properties();
properties.load(inputStream);
initNextLevel(properties, Config.class, new ConfigDesc());
} catch (IllegalAccessException e) {
logger.log(Level.ALL, "Parsing certification file failed, sky-walking api auto shutdown.");
logger.error("Parsing certification file failed, sky-walking api auto shutdown.");
} catch (IOException e) {
logger.log(Level.ALL, "Failed to read the certification file, sky-walking api auto shutdown.");
logger.error("Failed to read the certification file, sky-walking api auto shutdown.");
}
}
}

View File

@ -22,7 +22,6 @@ hbase.table_chain_detail=sw-chain-detail
hbase.table_call_chain=sw-call-chain
traceinfo.trace_info_column_cid=cid
mysql.url=jdbc:mysql://10.1.228.202:31316/test
@ -30,6 +29,8 @@ mysql.username=devrdbusr21
mysql.password=devrdbusr21
mysql.driver_class=com.mysql.jdbc.Driver
filter.filter_package_name=com.ai.cloud.skywalking.analysis.filter.impl
filter.filter_package_name=com.ai.cloud.skywalking.analysis.categorize2chain.filter.impl
chainnodesummary.interval=5
chainnodesummary.interval=5
reducer.reducer_number=2

View File

@ -43,7 +43,7 @@ public class CallChainMapperTest {
List<ChainInfo> chainInfos = new ArrayList<ChainInfo>();
chainInfos.add(chainInfo);
Categorize2ChainReducer.reduceAction(chainInfo.getUserId() + ":" + chainInfo.getEntranceNodeToken(), chainInfos.iterator());
Categorize2ChainReducer.reduceAction(chainInfo.getUserId() + ":" + chainInfo.getEntranceNodeToken(), chainInfos.iterator(), context);
}
public static List<Span> selectByTraceId(String traceId) throws IOException {