diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/AnalysisServerDriver.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/AnalysisServerDriver.java index 0763cec10..02ebee054 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/AnalysisServerDriver.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/AnalysisServerDriver.java @@ -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; } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java index 82f2d9879..fb7b114cf 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainMapper.java @@ -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 { @Override protected void map(ImmutableBytesWritable key, Result value, Context context) throws IOException, InterruptedException { + ConfigInitializer.initialize(); List spanList = new ArrayList(); ChainInfo chainInfo = null; try { @@ -32,7 +34,7 @@ public class Categorize2ChainMapper extends TableMapper { 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) { diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainReducer.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainReducer.java index dd6fc489b..d87a71894 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainReducer.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/categorize2chain/Categorize2ChainReducer.java @@ -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 values, Context context) throws IOException, InterruptedException { + ConfigInitializer.initialize(); int totalCount = reduceAction(key.toString(), values.iterator()); context.write(new Text(key.toString()), new IntWritable(totalCount)); } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Config.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Config.java index 2edad608e..a03c4c554 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Config.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/Config.java @@ -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; } } diff --git a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/ConfigInitializer.java b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/ConfigInitializer.java index 68c862555..f28ea1854 100644 --- a/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/ConfigInitializer.java +++ b/skywalking-analysis/src/main/java/com/ai/cloud/skywalking/analysis/config/ConfigInitializer.java @@ -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."); } } } diff --git a/skywalking-analysis/src/main/resources/analysis.conf b/skywalking-analysis/src/main/resources/analysis.conf index 050bf9b02..e0ec5cc78 100644 --- a/skywalking-analysis/src/main/resources/analysis.conf +++ b/skywalking-analysis/src/main/resources/analysis.conf @@ -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 \ No newline at end of file +chainnodesummary.interval=5 + +reducer.reducer_number=2 \ No newline at end of file diff --git a/skywalking-analysis/src/test/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapperTest.java b/skywalking-analysis/src/test/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapperTest.java index 7d6475150..28b9be157 100644 --- a/skywalking-analysis/src/test/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapperTest.java +++ b/skywalking-analysis/src/test/java/com/ai/cloud/skywalking/analysis/mapper/CallChainMapperTest.java @@ -43,7 +43,7 @@ public class CallChainMapperTest { List chainInfos = new ArrayList(); chainInfos.add(chainInfo); - Categorize2ChainReducer.reduceAction(chainInfo.getUserId() + ":" + chainInfo.getEntranceNodeToken(), chainInfos.iterator()); + Categorize2ChainReducer.reduceAction(chainInfo.getUserId() + ":" + chainInfo.getEntranceNodeToken(), chainInfos.iterator(), context); } public static List selectByTraceId(String traceId) throws IOException {