diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java index cef44a8f6..bcda4911e 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java @@ -78,4 +78,7 @@ public class Config { public static String CLIENT_PORT; } + public static class StorageChain { + public static long RETRY_STORAGE_WAIT_TIME = 50L; + } } \ No newline at end of file diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BuriedPointEntry.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BuriedPointEntry.java index a35166442..26072d72a 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BuriedPointEntry.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BuriedPointEntry.java @@ -18,7 +18,7 @@ public class BuriedPointEntry { private String processNo; - private BuriedPointEntry(){ + private BuriedPointEntry() { } @@ -87,7 +87,7 @@ public class BuriedPointEntry { result.exceptionStack = fieldValues[7]; result.spanType = fieldValues[8].charAt(0); result.isReceiver = Boolean.getBoolean(fieldValues[9]); - result.businessKey = fieldValues[10]; + result.businessKey = fieldValues[10].replace('^', '-'); result.processNo = fieldValues[11]; return result; } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/ChainException.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/ChainException.java new file mode 100644 index 000000000..7410c97f6 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/ChainException.java @@ -0,0 +1,12 @@ +package com.ai.cloud.skywalking.reciever.storage; + +public class ChainException extends RuntimeException { + + public ChainException(Throwable cause) { + super(cause); + } + + public ChainException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java index 7a13367b3..dd7aa3d02 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java @@ -1,41 +1,48 @@ package com.ai.cloud.skywalking.reciever.storage; +import com.ai.cloud.skywalking.reciever.conf.Config; +import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry; +import com.ai.cloud.skywalking.reciever.storage.chain.SaveToHBaseChain; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + import java.util.ArrayList; import java.util.List; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; - -import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry; -import com.ai.cloud.skywalking.reciever.storage.chain.SaveToHBaseChain; - public class StorageChainController { - private static Logger logger = LogManager - .getLogger(StorageChainController.class); - - private static List chainArray = new ArrayList(); - - static{ - chainArray.add(new SaveToHBaseChain()); - } + private static Logger logger = LogManager + .getLogger(StorageChainController.class); - public static void doStorage(String buriedPointDatas) { - String[] buriedPointData = buriedPointDatas.split(";"); - if (buriedPointData == null || buriedPointData.length == 0) { - return; - } - for (String buriedPoint : buriedPointData) { - try { - if(buriedPoint == null || buriedPoint.trim().length() == 0){ - continue; - } - BuriedPointEntry entry = BuriedPointEntry.convert(buriedPoint); - Chain chain = new Chain(chainArray); - chain.doChain(entry, buriedPoint); - } catch (Throwable e) { - logger.error("ready to save buriedPoint error, choose to ignore. data=" - + buriedPoint, e); - } - } - } + private static List chainArray = new ArrayList(); + + static { + chainArray.add(new SaveToHBaseChain()); + } + + public static void doStorage(String buriedPointDatas) { + String[] buriedPointData = buriedPointDatas.split(";"); + if (buriedPointData == null || buriedPointData.length == 0) { + return; + } + for (String buriedPoint : buriedPointData) { + try { + if (buriedPoint == null || buriedPoint.trim().length() == 0) { + continue; + } + BuriedPointEntry entry = BuriedPointEntry.convert(buriedPoint); + while(true) { + try { + Chain chain = new Chain(chainArray); + chain.doChain(entry, buriedPoint); + break; + } catch (Throwable e) { + Thread.sleep(Config.StorageChain.RETRY_STORAGE_WAIT_TIME); + } + } + } catch (Throwable e) { + logger.error("ready to save buriedPoint error, choose to ignore. data=" + + buriedPoint, e); + } + } + } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/HBaseOperator.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/HBaseOperator.java deleted file mode 100644 index a3695b773..000000000 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/HBaseOperator.java +++ /dev/null @@ -1,79 +0,0 @@ -package com.ai.cloud.skywalking.reciever.storage.chain; - -import com.ai.cloud.skywalking.reciever.conf.Config; -import com.ai.cloud.skywalking.reciever.conf.ConfigInitializer; -import org.apache.hadoop.conf.Configuration; -import org.apache.hadoop.hbase.HBaseConfiguration; -import org.apache.hadoop.hbase.HColumnDescriptor; -import org.apache.hadoop.hbase.HTableDescriptor; -import org.apache.hadoop.hbase.TableName; -import org.apache.hadoop.hbase.client.*; -import org.apache.hadoop.hbase.util.Bytes; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; - -import java.io.IOException; -import java.util.Properties; -import java.util.UUID; - -public class HBaseOperator { - private static Logger logger = LogManager.getLogger(HBaseOperator.class); - private static Configuration configuration = null; - private static Connection connection; - - private static void initHBaseClient() throws IOException { - if (configuration == null) { - configuration = HBaseConfiguration.create(); - if (Config.HBaseConfig.ZK_HOSTNAME == null || "".equals(Config.HBaseConfig.ZK_HOSTNAME)) { - logger.error("Miss HBase ZK quorum Configuration", new IllegalArgumentException("Miss HBase ZK quorum Configuration")); - System.exit(-1); - } - configuration.set("hbase.zookeeper.quorum", Config.HBaseConfig.ZK_HOSTNAME); - configuration.set("hbase.zookeeper.property.clientPort", Config.HBaseConfig.CLIENT_PORT); - connection = ConnectionFactory.createConnection(configuration); - } - } - - private static void createTable(String tableName) { - - try { - initHBaseClient(); - Admin admin = connection.getAdmin(); - if (!admin.isTableAvailable(TableName.valueOf(tableName))) { - HTableDescriptor tableDesc = new HTableDescriptor(TableName.valueOf(tableName)); - tableDesc.addFamily(new HColumnDescriptor(Config.HBaseConfig.FAMILY_COLUMN_NAME)); - admin.createTable(tableDesc); - logger.info("Create table [{}] ok!", tableName); - } - } catch (IOException e) { - logger.error("Create table[{}] failed", tableName, e); - } - } - - public static void insert(String rowKey, String qualifier, String value) { - insert(Config.HBaseConfig.TABLE_NAME, rowKey, qualifier, value); - } - - public static void insert(String tableName, String rowKey, String qualifier, String value) { - try { - createTable(tableName); - Table table = connection.getTable(TableName.valueOf(tableName)); - Put put = new Put(Bytes.toBytes(rowKey)); - put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(qualifier), Bytes - .toBytes(value)); - table.put(put); - if (logger.isDebugEnabled()) { - logger.debug("Insert data[RowKey:{}] success.", rowKey); - } - } catch (IOException e) { - logger.error("Insert the data error.RowKey:[{}],Qualifier[{}],value[{}]", rowKey, qualifier, value, e); - } - } - - public static void main(String[] args) throws IllegalAccessException, IOException { - Properties config = new Properties(); - config.load(HBaseOperator.class.getResourceAsStream("/config.properties")); - ConfigInitializer.initialize(config, Config.class); - HBaseOperator.createTable("test3"); - } -} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToHBaseChain.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToHBaseChain.java index 61599c78a..9801651d1 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToHBaseChain.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToHBaseChain.java @@ -1,22 +1,85 @@ package com.ai.cloud.skywalking.reciever.storage.chain; -import org.apache.commons.lang.StringUtils; - +import com.ai.cloud.skywalking.reciever.conf.Config; import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry; import com.ai.cloud.skywalking.reciever.storage.Chain; +import com.ai.cloud.skywalking.reciever.storage.ChainException; import com.ai.cloud.skywalking.reciever.storage.IStorageChain; +import org.apache.commons.lang.StringUtils; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.hbase.HBaseConfiguration; +import org.apache.hadoop.hbase.HColumnDescriptor; +import org.apache.hadoop.hbase.HTableDescriptor; +import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.client.*; +import org.apache.hadoop.hbase.util.Bytes; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; -public class SaveToHBaseChain implements IStorageChain{ +import java.io.IOException; - @Override - public void doChain(BuriedPointEntry entry, String entryOriginData, Chain chain) { +public class SaveToHBaseChain implements IStorageChain { + private static Logger logger = LogManager.getLogger(SaveToHBaseChain.class); + private static Configuration configuration = null; + private static Connection connection; + + @Override + public void doChain(BuriedPointEntry entry, String entryOriginData, Chain chain) { if (StringUtils.isEmpty(entry.getParentLevel().trim())) { - HBaseOperator.insert(entry.getTraceId(), String.valueOf(entry.getLevelId()), entryOriginData); + insert(entry.getTraceId(), String.valueOf(entry.getLevelId()), entryOriginData); } else { - HBaseOperator.insert(entry.getTraceId(), entry.getParentLevel() + "." + entry.getLevelId(), entryOriginData); + insert(entry.getTraceId(), entry.getParentLevel() + "." + entry.getLevelId(), entryOriginData); } - - - } + + chain.doChain(entry, entryOriginData); + } + + private static void initHBaseClient() throws IOException { + if (configuration == null) { + configuration = HBaseConfiguration.create(); + if (Config.HBaseConfig.ZK_HOSTNAME == null || "".equals(Config.HBaseConfig.ZK_HOSTNAME)) { + logger.error("Miss HBase ZK quorum Configuration", new IllegalArgumentException("Miss HBase ZK quorum Configuration")); + System.exit(-1); + } + configuration.set("hbase.zookeeper.quorum", Config.HBaseConfig.ZK_HOSTNAME); + configuration.set("hbase.zookeeper.property.clientPort", Config.HBaseConfig.CLIENT_PORT); + connection = ConnectionFactory.createConnection(configuration); + } + } + + static { + try { + initHBaseClient(); + Admin admin = connection.getAdmin(); + if (!admin.isTableAvailable(TableName.valueOf(Config.HBaseConfig.TABLE_NAME))) { + HTableDescriptor tableDesc = new HTableDescriptor(TableName.valueOf(Config.HBaseConfig.TABLE_NAME)); + tableDesc.addFamily(new HColumnDescriptor(Config.HBaseConfig.FAMILY_COLUMN_NAME)); + admin.createTable(tableDesc); + logger.info("Create table [{}] ok!", Config.HBaseConfig.TABLE_NAME); + } + } catch (IOException e) { + logger.error("Create table[{}] failed", Config.HBaseConfig.TABLE_NAME, e); + } + } + + public static void insert(String rowKey, String qualifier, String value) { + insert(Config.HBaseConfig.TABLE_NAME, rowKey, qualifier, value); + } + + public static void insert(String tableName, String rowKey, String qualifier, String value) { + try { + Table table = connection.getTable(TableName.valueOf(tableName)); + Put put = new Put(Bytes.toBytes(rowKey)); + put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(qualifier), Bytes + .toBytes(value)); + table.put(put); + if (logger.isDebugEnabled()) { + logger.debug("Insert data[RowKey:{}] success.", rowKey); + } + } catch (IOException e) { + logger.error("Insert the data error.RowKey:[{}],Qualifier[{}],value[{}]", rowKey, qualifier, value, e); + throw new ChainException(e); + } + } }