diff --git a/skywalking-storage-center/skywalking-storage/pom.xml b/skywalking-storage-center/skywalking-storage/pom.xml
index 00e887303..cfac8066d 100644
--- a/skywalking-storage-center/skywalking-storage/pom.xml
+++ b/skywalking-storage-center/skywalking-storage/pom.xml
@@ -46,6 +46,11 @@
skywalking-registry
${project.version}
+
+ com.zaxxer
+ HikariCP
+ 2.4.3
+
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java
index 7379c89af..9fb51248d 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java
@@ -14,10 +14,13 @@ import com.a.eye.skywalking.storage.config.ConfigInitializer;
import com.a.eye.skywalking.storage.data.IndexDataCapacityMonitor;
import com.a.eye.skywalking.storage.notifier.SearchNotifier;
import com.a.eye.skywalking.storage.notifier.StorageNotifier;
+import com.a.eye.skywalking.storage.util.NetUtils;
import java.io.IOException;
import java.util.Properties;
+import static com.a.eye.skywalking.storage.config.Config.RegistryCenter.REGISTRY_PATH_PREFIX;
+
/**
* Created by xin on 2016/11/12.
*/
@@ -34,13 +37,14 @@ public class Main {
public static void main(String[] args) {
try {
initializeParam();
+ new Thread(new IndexDataCapacityMonitor()).start();
transferService =
TransferServiceBuilder.newBuilder(Config.Server.PORT).startSpanStorageService(new StorageNotifier())
.startTraceSearchService(new SearchNotifier()).build();
transferService.start();
logger.info("transfer service started successfully!");
- new Thread(new IndexDataCapacityMonitor()).start();
+
registryNode();
logger.info("storage service started successfully!");
Thread.currentThread().join();
@@ -52,12 +56,15 @@ public class Main {
}
private static void registryNode() {
- //TODO auth info auth schema
RegistryCenter registryCenter =
RegistryCenterFactory.INSTANCE.getRegistryCenter(CenterType.DEFAULT_CENTER_TYPE);
Properties registerConfig = new Properties();
registerConfig.setProperty(ZookeeperConfig.CONNECT_URL, Config.RegistryCenter.CONNECT_URL);
+ registerConfig.setProperty(ZookeeperConfig.AUTH_SCHEMA, Config.RegistryCenter.AUTH_SCHEMA);
+ registerConfig.setProperty(ZookeeperConfig.AUTH_INFO, Config.RegistryCenter.AUTH_INFO);
registryCenter.start(registerConfig);
+ registryCenter.register(
+ REGISTRY_PATH_PREFIX + NetUtils.getLocalAddress().getHostAddress() + ":" + Config.Server.PORT);
}
private static void initializeParam() throws IllegalAccessException, IOException {
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java
index e4af3922c..aa498e994 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java
@@ -33,10 +33,23 @@ public class Config {
public static String STORAGE_INDEX_FILE_NAME = "dataIndex";
public static long MAX_CAPACITY_PER_INDEX = 1000 * 1000 * 1000 * 1000;
+
}
public static class RegistryCenter {
+
+ public static String AUTH_INFO = "";
+
+ public static String AUTH_SCHEMA = "";
+
public static String CONNECT_URL = "127.0.0.1:2181";
+
+ public static String REGISTRY_PATH_PREFIX = "/storage_list/";
+ }
+
+
+ public static class SpanFinder {
+ public static int MAX_CACHE_SIZE = 10;
}
}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java
index 897d2473f..e417b7b74 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java
@@ -1,22 +1,31 @@
package com.a.eye.skywalking.storage.data;
+import com.a.eye.skywalking.logging.api.ILog;
+import com.a.eye.skywalking.logging.api.LogManager;
import com.a.eye.skywalking.storage.block.index.BlockIndexEngine;
+import com.a.eye.skywalking.storage.config.Config;
import com.a.eye.skywalking.storage.data.file.DataFileReader;
import com.a.eye.skywalking.storage.data.index.*;
import com.a.eye.skywalking.storage.data.spandata.SpanData;
+import com.zaxxer.hikari.HikariConfig;
+import com.zaxxer.hikari.HikariDataSource;
-import java.util.ArrayList;
-import java.util.Iterator;
-import java.util.List;
+import java.sql.SQLException;
+import java.util.*;
public class SpanDataFinder {
+ private static ILog logger = LogManager.getLogger(SpanDataFinder.class);
+ private static IndexDataSourceCache datasourceCache = new IndexDataSourceCache(Config.SpanFinder.MAX_CACHE_SIZE);
+
public static List find(String traceId) {
long blockIndex = BlockIndexEngine.newFinder().find(fetchStartTimeFromTraceId(traceId));
if (blockIndex == 0) {
return new ArrayList();
}
- IndexDBConnector indexDBConnector = new IndexDBConnector(blockIndex);
+
+ IndexDBConnector indexDBConnector = fetchIndexDBConnector(blockIndex);
IndexMetaCollection indexMetaCollection = indexDBConnector.queryByTraceId(traceId);
+ indexDBConnector.close();
Iterator> iterator =
IndexMetaCollections.group(indexMetaCollection, new GroupKeyBuilder() {
@@ -35,8 +44,59 @@ public class SpanDataFinder {
return result;
}
+ private static IndexDBConnector fetchIndexDBConnector(long blockIndex) {
+ HikariDataSource datasource = getOrCreate(blockIndex);
+ IndexDBConnector indexDBConnector = null;
+ try {
+ indexDBConnector = new IndexDBConnector(datasource.getConnection());
+ } catch (SQLException e) {
+ logger.warn("Failed to get connection from datasource,", e);
+ indexDBConnector = new IndexDBConnector(blockIndex);
+ }
+ return indexDBConnector;
+ }
+
+ private static HikariDataSource getOrCreate(long blockIndex) {
+ HikariDataSource datasource = datasourceCache.get(blockIndex);
+ if (datasource == null) {
+ HikariConfig dataSourceConfig = generateDatasourceConfig(blockIndex);
+ datasource = new HikariDataSource(dataSourceConfig);
+ datasourceCache.put(blockIndex, datasource);
+ }
+ return datasource;
+ }
+
+ private static HikariConfig generateDatasourceConfig(long blockIndex) {
+ HikariConfig config = new HikariConfig();
+ config.setJdbcUrl(new ConnectURLGenerator(Config.DataIndex.BASE_PATH, Config.DataIndex.STORAGE_INDEX_FILE_NAME)
+ .generate(blockIndex));
+ config.setDriverClassName("org.hsqldb.jdbc.JDBCDriver");
+ config.setMaximumPoolSize(20);
+ config.setMinimumIdle(5);
+ return config;
+ }
+
private static long fetchStartTimeFromTraceId(String traceId) {
String[] traceIdSegment = traceId.split("\\.");
return Long.parseLong(traceIdSegment[traceIdSegment.length - 5]);
}
+
+ private static class IndexDataSourceCache extends LinkedHashMap {
+
+ private int cacheSize;
+
+ @Override
+ protected boolean removeEldestEntry(Map.Entry eldest) {
+ boolean removed = size() > cacheSize;
+ if (removed) {
+ eldest.getValue().close();
+ }
+ return removed;
+ }
+
+ public IndexDataSourceCache(int cacheSize) {
+ super((int) Math.ceil(cacheSize / 0.75) + 1, 0.75f, true);
+ this.cacheSize = cacheSize;
+ }
+ }
}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/ConnectURLGenerator.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/ConnectURLGenerator.java
new file mode 100644
index 000000000..14ad9cf4a
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/ConnectURLGenerator.java
@@ -0,0 +1,20 @@
+package com.a.eye.skywalking.storage.data.index;
+
+/**
+ * Created by xin on 2016/11/13.
+ */
+public class ConnectURLGenerator {
+
+ private String basePath;
+ private String dbFileName;
+
+ public ConnectURLGenerator(String basePath, String dbFileName) {
+ this.basePath = basePath;
+ this.dbFileName = dbFileName;
+ }
+
+
+ public String generate(long timestamp) {
+ return "jdbc:hsqldb:file:" + basePath + "/" + timestamp + "/" + dbFileName;
+ }
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnector.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnector.java
index 3a32e047c..5c3ee447f 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnector.java
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnector.java
@@ -41,6 +41,11 @@ public class IndexDBConnector {
createTableAndIndexIfNecessary();
}
+ public IndexDBConnector(Connection connection){
+ this.connection = connection;
+ createTableAndIndexIfNecessary();
+ }
+
private void createTableAndIndexIfNecessary() {
try {
if (validateTableIsExists()) {
@@ -146,22 +151,6 @@ public class IndexDBConnector {
}
}
- class ConnectURLGenerator {
-
- private String basePath;
- private String dbFileName;
-
- private ConnectURLGenerator(String basePath, String dbFileName) {
- this.basePath = basePath;
- this.dbFileName = dbFileName;
- }
-
-
- public String generate(long timestamp) {
- return "jdbc:hsqldb:file:" + basePath + "/" + timestamp + "/" + dbFileName;
- }
- }
-
public void close() {
try {
connection.close();
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnectorCache.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnectorCache.java
index f3b5a30ca..fa52d96fd 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnectorCache.java
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnectorCache.java
@@ -41,7 +41,7 @@ public class IndexDBConnectorCache {
}
public LRUCache(int cacheSize) {
- super((int) Math.ceil(MAX_CACHE_SIZE / 0.75) + 1, 0.75f, true);
+ super((int) Math.ceil(cacheSize / 0.75) + 1, 0.75f, true);
}
}
}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AckSpanData.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AckSpanData.java
index c9fe3ce57..de14e3d26 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AckSpanData.java
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AckSpanData.java
@@ -40,4 +40,15 @@ public class AckSpanData extends AbstractSpanData {
return buildLevelId(ackSpan.getParentLevel(), ackSpan.getLevelId());
}
+ public long getCost() {
+ return ackSpan.getCost();
+ }
+
+ public String getExceptionStack() {
+ return ackSpan.getExceptionStack();
+ }
+
+ public int getStatusCode() {
+ return ackSpan.getStatusCode();
+ }
}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/RequestSpanData.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/RequestSpanData.java
index 41dd932b2..96218cae2 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/RequestSpanData.java
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/RequestSpanData.java
@@ -39,4 +39,32 @@ public class RequestSpanData extends AbstractSpanData {
public String getLevelId() {
return buildLevelId(requestSpan.getParentLevel(), requestSpan.getLevelId());
}
+
+ public String getAddress(){
+ return requestSpan.getAddress();
+ }
+
+ public String getApplicationId(){
+ return requestSpan.getApplicationId();
+ }
+
+ public String getProcessNo(){
+ return requestSpan.getProcessNo();
+ }
+
+ public long getStartTime(){
+ return requestSpan.getStartDate();
+ }
+
+ public String getBusinessKey() {
+ return requestSpan.getBussinessKey();
+ }
+
+ public String getCallType() {
+ return requestSpan.getCallType();
+ }
+
+ public int getType() {
+ return requestSpan.getSpanType();
+ }
}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/SearchNotifier.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/SearchNotifier.java
index 63ae56435..61967d79e 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/SearchNotifier.java
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/SearchNotifier.java
@@ -3,50 +3,16 @@ package com.a.eye.skywalking.storage.notifier;
import com.a.eye.skywalking.network.grpc.Span;
import com.a.eye.skywalking.network.listener.TraceSearchNotifier;
import com.a.eye.skywalking.storage.data.SpanDataFinder;
-import com.a.eye.skywalking.storage.data.spandata.AckSpanData;
-import com.a.eye.skywalking.storage.data.spandata.RequestSpanData;
import com.a.eye.skywalking.storage.data.spandata.SpanData;
-import java.util.ArrayList;
-import java.util.HashMap;
import java.util.List;
-import java.util.Map;
public class SearchNotifier implements TraceSearchNotifier {
@Override
public List search(String s) {
List data = SpanDataFinder.find(s);
- return mergeSpanData(data);
+ SpanDataHelper helper = new SpanDataHelper(data);
+ return helper.category().mergeData();
}
-
- private List mergeSpanData(List data) {
- //// TODO: 2016/11/12 需要修改
- Map requestSpen = new HashMap();
- Map ackSpen = new HashMap();
-
- for (SpanData spanData : data) {
- if (spanData instanceof RequestSpanData) {
- requestSpen.put(spanData.getLevelId(), (RequestSpanData) spanData);
- } else {
- ackSpen.put(spanData.getLevelId(), (AckSpanData) spanData);
- }
- }
-
- List mergedSpan = new ArrayList();
- for (Map.Entry entry : requestSpen.entrySet()) {
- AckSpanData ackSpanData = ackSpen.get(entry.getKey());
- if (ackSpanData != null) {
- mergedSpan.add(mergeSpan(entry.getValue(), ackSpanData));
- }
- }
-
- return mergedSpan;
- }
-
- private Span mergeSpan(RequestSpanData value, AckSpanData ackSpanData) {
- return null;
- }
-
-
}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/SpanDataHelper.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/SpanDataHelper.java
new file mode 100644
index 000000000..b685f7920
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/SpanDataHelper.java
@@ -0,0 +1,64 @@
+package com.a.eye.skywalking.storage.notifier;
+
+import com.a.eye.skywalking.network.grpc.Span;
+import com.a.eye.skywalking.storage.data.spandata.AckSpanData;
+import com.a.eye.skywalking.storage.data.spandata.RequestSpanData;
+import com.a.eye.skywalking.storage.data.spandata.SpanData;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Created by xin on 2016/11/12.
+ */
+public class SpanDataHelper {
+ public HashMap levelIdRequestSpanDataMapping = new HashMap();
+ public HashMap levelIdAckSpanDataMapping = new HashMap();
+
+ private List data;
+
+ public SpanDataHelper(List data) {
+ this.data = data;
+ }
+
+ public SpanDataHelper category() {
+ for (SpanData spanData : data) {
+ if (spanData instanceof RequestSpanData) {
+ levelIdRequestSpanDataMapping.put(spanData.getLevelId(), (RequestSpanData) spanData);
+ } else {
+ levelIdAckSpanDataMapping.put(spanData.getLevelId(), (AckSpanData) spanData);
+ }
+ }
+
+ return this;
+ }
+
+ public List mergeData() {
+ List span = new ArrayList();
+ for (Map.Entry entry : levelIdRequestSpanDataMapping.entrySet()) {
+ AckSpanData ackSpanData = levelIdAckSpanDataMapping.get(entry.getKey());
+ if (ackSpanData != null) {
+ span.add(mergeSpan(entry.getValue(), ackSpanData));
+ }
+ }
+
+ return span;
+ }
+
+ private Span mergeSpan(RequestSpanData requestSpanData, AckSpanData ackSpanData) {
+ Span.Builder builder = Span.newBuilder().setAddress(requestSpanData.getAddress())
+ .setApplicationId(requestSpanData.getApplicationId()).setBusinessKey(requestSpanData.getBusinessKey())
+ .setCallType(requestSpanData.getCallType()).setCost(ackSpanData.getCost());
+ if (ackSpanData.getExceptionStack() != null && ackSpanData.getExceptionStack().length() > 0) {
+ builder = builder.setExceptionStack(ackSpanData.getExceptionStack());
+ }
+
+ builder = builder.setLevelId(requestSpanData.getLevelId()).setProcessNo(requestSpanData.getProcessNo())
+ .setSpanType(requestSpanData.getType()).setStarttime(requestSpanData.getStartTime())
+ .setStatusCode(ackSpanData.getStatusCode()).setTraceId(requestSpanData.getTraceId());
+ return builder.build();
+ }
+
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/util/NetUtils.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/util/NetUtils.java
new file mode 100644
index 000000000..63237a2d3
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/util/NetUtils.java
@@ -0,0 +1,74 @@
+package com.a.eye.skywalking.storage.util;
+
+import com.a.eye.skywalking.logging.api.ILog;
+import com.a.eye.skywalking.logging.api.LogManager;
+
+import java.net.InetAddress;
+import java.net.NetworkInterface;
+import java.util.Enumeration;
+import java.util.regex.Pattern;
+
+/**
+ * Created by xin on 2016/11/12.
+ */
+public class NetUtils {
+
+ private static ILog logger = LogManager.getLogger(NetUtils.class);
+ public static final String LOCALHOST = "127.0.0.1";
+ public static final String ANYHOST = "0.0.0.0";
+ private static final Pattern IP_PATTERN = Pattern.compile("\\d{1,3}(\\.\\d{1,3}){3,5}$");
+
+ public static InetAddress getLocalAddress() {
+ InetAddress localAddress = null;
+ try {
+ localAddress = InetAddress.getLocalHost();
+ if (isValidAddress(localAddress)) {
+ return localAddress;
+ }
+ } catch (Throwable e) {
+ logger.warn("Failed to get ip address.", e);
+ }
+
+ try {
+ // 获取所有的网卡
+ Enumeration interfaces = NetworkInterface.getNetworkInterfaces();
+ if (interfaces != null) {
+ while (interfaces.hasMoreElements()) {
+ try {
+ NetworkInterface network = interfaces.nextElement();
+ // 遍历网卡中所有绑定的地址
+ Enumeration addresses = network.getInetAddresses();
+ if (addresses != null) {
+ while (addresses.hasMoreElements()) {
+ try {
+ InetAddress address = addresses.nextElement();
+ // 判断地址是否为合法的IP地址
+ if (isValidAddress(address)) {
+ return address;
+ }
+ } catch (Throwable e) {
+ logger.warn("Failed to get ip address.", e);
+ }
+ }
+ }
+ } catch (Throwable e) {
+ logger.warn("Failed to get ip address.", e);
+ }
+ }
+ }
+ } catch (Throwable e) {
+ logger.warn("Failed to get ip address.", e);
+ }
+
+ return localAddress;
+ }
+
+ private static boolean isValidAddress(InetAddress address) {
+ if (address == null || address.isLoopbackAddress())
+ return false;
+ String name = address.getHostAddress();
+ // 不能是0.0.0.0 也不能是127.0.0.1 并且还得符合IP的正则
+ return (name != null && !ANYHOST.equals(name) && !LOCALHOST.equals(name) && IP_PATTERN.matcher(name).matches());
+ }
+
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/resources/config.properties b/skywalking-storage-center/skywalking-storage/src/main/resources/config.properties
index ff4a07e5c..16a43fd1a 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/resources/config.properties
+++ b/skywalking-storage-center/skywalking-storage/src/main/resources/config.properties
@@ -1,21 +1,35 @@
#
-server.port = 34000
-
+server.port=34000
#
-blockindex.storage_base_path= /tmp/skywalking/block-index
#
-blockindex.data_file_index_file_name= data_file.index
-
#
-datafile.base_path= /tmp/skywalking/data/file
+blockindex.storage_base_path=/tmp/skywalking/block-index
+#
+blockindex.data_file_index_file_name=data_file.index
+#
+#
+#
+datafile.base_path=/tmp/skywalking/data/file
+#
+datafile.max_length=3221225472
+#
#
-datafile.max_length= 3221225472
-
#存放数据文件索引表名
-dataindex.table_name= data_index
+dataindex.table_name=data_index
#数据文件索引存储位置
-dataindex.base_path= /tmp/skywalking/data/index
+dataindex.base_path=/tmp/skywalking/data/index
#
-dataindex.storage_index_file_name= dataIndex
+dataindex.storage_index_file_name=dataIndex
+#
+dataindex.max_capacity_per_index=1000000000
+#
+#
+#
+registrycenter.auth_info=
+#
+registrycenter.auth_schema=
+#
+registrycenter.connect_url=
+#
+registrycenter.registry_path_prefix=
#
-dataindex.max_capacity_per_index= 1000000000