1. 增加
This commit is contained in:
parent
f7f8f44cb6
commit
e032ef282f
|
|
@ -78,4 +78,7 @@ public class Config {
|
|||
public static String CLIENT_PORT;
|
||||
}
|
||||
|
||||
public static class StorageChain {
|
||||
public static long RETRY_STORAGE_WAIT_TIME = 50L;
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<IStorageChain> chainArray = new ArrayList<IStorageChain>();
|
||||
|
||||
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<IStorageChain> chainArray = new ArrayList<IStorageChain>();
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
}
|
||||
}
|
||||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue