From 207561b1ee8f537a7b6d12568313ef7a29b9f2e8 Mon Sep 17 00:00:00 2001 From: ascrutae Date: Sat, 12 Nov 2016 18:47:34 +0800 Subject: [PATCH] =?UTF-8?q?=E5=AE=8C=E6=88=90storage=E5=8A=9F=E8=83=BD?= =?UTF-8?q?=E8=A1=94=E6=8E=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- skywalking-alarm/pom.xml | 1 - skywalking-logging/pom.xml | 2 +- .../a/eye/skywalking/logging/api/ILog.java | 6 + .../skywalking/logging/api/NoopLogger.java | 10 ++ .../logging/impl/log4j2/Log4j2Logger.java | 10 ++ .../skywalking/network/TransferService.java | 2 +- skywalking-registry/pom.xml | 7 +- .../registry/api/NotifyListener.java | 2 +- .../zookeeper/ZookeeperRegistryCenter.java | 126 +++++++++++++++--- .../registry/RegistryCenterFactoryTest.java | 73 ++++++++++ .../src/test/resources/log4j.properties | 13 ++ .../src/test/resources/log4j2.xml | 13 ++ skywalking-storage-center/pom.xml | 14 ++ .../skywalking-storage/bin/swserver.sh | 52 -------- .../com/a/eye/skywalking/storage/Main.java | 62 +++++++++ .../storage/block/index/BlockFinder.java | 4 - .../eye/skywalking/storage/config/Config.java | 12 +- .../storage/config/ConfigInitializer.java | 70 ++++++++++ .../skywalking/storage/config/Constants.java | 9 +- .../data/IndexDataCapacityMonitor.java | 3 +- .../eye/skywalking/storage/data/SpanData.java | 8 -- .../storage/data/SpanDataConsumer.java | 4 +- .../storage/data/SpanDataFinder.java | 16 +-- .../storage/data/file/DataFile.java | 4 +- .../storage/data/file/DataFileReader.java | 27 +++- .../storage/data/file/DataFileWriter.java | 2 +- .../storage/data/index/IndexDBConnector.java | 33 ++++- .../storage/data/index/IndexMetaInfo.java | 25 ++-- .../data/spandata/AbstractSpanData.java | 16 +++ .../storage/data/spandata/AckSpanData.java | 43 ++++++ .../data/spandata/RequestSpanData.java | 42 ++++++ .../storage/data/spandata/SpanData.java | 14 ++ .../data/spandata/SpanDataBuilder.java | 44 ++++++ .../storage/data/spandata/SpanType.java | 30 +++++ .../storage/notifier/SearchNotifier.java | 52 ++++++++ .../storage/notifier/StorageNotifier.java | 45 +++++++ ...e.skywalking.reciever.processor.IProcessor | 4 - .../src/main/resources/config.properties | 85 +++--------- .../search/HyperSqlSearchSpeedReporter.java | 92 ------------- .../search/SearchSpeedReporter.java | 37 ----- .../a/eye/skywalking/search/TreeSetTest.java | 40 ------ .../block/index/DataIndexFileOperator.java | 7 - 42 files changed, 786 insertions(+), 375 deletions(-) create mode 100644 skywalking-registry/src/test/java/com/a/eye/skywalking/registry/RegistryCenterFactoryTest.java create mode 100644 skywalking-registry/src/test/resources/log4j.properties create mode 100644 skywalking-registry/src/test/resources/log4j2.xml delete mode 100644 skywalking-storage-center/skywalking-storage/bin/swserver.sh create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/ConfigInitializer.java delete mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanData.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AbstractSpanData.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AckSpanData.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/RequestSpanData.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanData.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanDataBuilder.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanType.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/SearchNotifier.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/StorageNotifier.java delete mode 100644 skywalking-storage-center/skywalking-storage/src/main/resources/META-INF/services/com.a.eye.skywalking.reciever.processor.IProcessor delete mode 100644 skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/HyperSqlSearchSpeedReporter.java delete mode 100644 skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/SearchSpeedReporter.java delete mode 100644 skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/TreeSetTest.java delete mode 100644 skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/block/index/DataIndexFileOperator.java diff --git a/skywalking-alarm/pom.xml b/skywalking-alarm/pom.xml index 31b5ad926..956e0ac02 100644 --- a/skywalking-alarm/pom.xml +++ b/skywalking-alarm/pom.xml @@ -142,7 +142,6 @@ *.properties *.xml - sky-alarm-server ${project.build.directory}/installer/lib diff --git a/skywalking-logging/pom.xml b/skywalking-logging/pom.xml index 97579c9fb..d6ae3605b 100644 --- a/skywalking-logging/pom.xml +++ b/skywalking-logging/pom.xml @@ -14,7 +14,7 @@ http://maven.apache.org skywalking-logging-api - skywalking-logging-log4j-impl + skywalking-logging-impl-log4j2 diff --git a/skywalking-logging/skywalking-logging-api/src/main/java/com/a/eye/skywalking/logging/api/ILog.java b/skywalking-logging/skywalking-logging-api/src/main/java/com/a/eye/skywalking/logging/api/ILog.java index 656d36305..e978b78a1 100644 --- a/skywalking-logging/skywalking-logging-api/src/main/java/com/a/eye/skywalking/logging/api/ILog.java +++ b/skywalking-logging/skywalking-logging-api/src/main/java/com/a/eye/skywalking/logging/api/ILog.java @@ -8,7 +8,13 @@ public interface ILog { void info(String format, Object... arguments); + void warn(String format, Object... arguments); + + void warn(String format, Object arguments, Throwable e); + void error(String format, Throwable e); void error(String format, Object argument, Throwable e); + + } diff --git a/skywalking-logging/skywalking-logging-api/src/main/java/com/a/eye/skywalking/logging/api/NoopLogger.java b/skywalking-logging/skywalking-logging-api/src/main/java/com/a/eye/skywalking/logging/api/NoopLogger.java index f1491875e..83479d813 100644 --- a/skywalking-logging/skywalking-logging-api/src/main/java/com/a/eye/skywalking/logging/api/NoopLogger.java +++ b/skywalking-logging/skywalking-logging-api/src/main/java/com/a/eye/skywalking/logging/api/NoopLogger.java @@ -17,6 +17,16 @@ public class NoopLogger implements ILog{ } + @Override + public void warn(String format, Object... arguments) { + + } + + @Override + public void warn(String format, Object arguments, Throwable e) { + + } + @Override public void error(String format, Throwable e) { diff --git a/skywalking-logging/skywalking-logging-impl-log4j2/src/main/java/com/a/eye/skywalking/logging/impl/log4j2/Log4j2Logger.java b/skywalking-logging/skywalking-logging-impl-log4j2/src/main/java/com/a/eye/skywalking/logging/impl/log4j2/Log4j2Logger.java index 50bb95e22..a4aaee754 100644 --- a/skywalking-logging/skywalking-logging-impl-log4j2/src/main/java/com/a/eye/skywalking/logging/impl/log4j2/Log4j2Logger.java +++ b/skywalking-logging/skywalking-logging-impl-log4j2/src/main/java/com/a/eye/skywalking/logging/impl/log4j2/Log4j2Logger.java @@ -23,6 +23,16 @@ public class Log4j2Logger implements ILog { logger.info(message, arguments); } + @Override + public void warn(String format, Object... arguments) { + logger.warn(format, arguments); + } + + @Override + public void warn(String format, Object arguments, Throwable e) { + logger.warn(format, arguments, e); + } + @Override public void error(String message, Throwable e) { logger.error(message, e); diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/TransferService.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/TransferService.java index 901cdc42e..66ccc7ad6 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/TransferService.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/TransferService.java @@ -30,7 +30,7 @@ public class TransferService { blockUntilShutdown(); } - private void stop() { + public void stop() { if (server != null) { server.shutdown(); } diff --git a/skywalking-registry/pom.xml b/skywalking-registry/pom.xml index a8375e5ea..5d7423cd9 100644 --- a/skywalking-registry/pom.xml +++ b/skywalking-registry/pom.xml @@ -21,12 +21,17 @@ org.apache.zookeeper zookeeper - 3.4.7 + 3.4.8 com.a.eye skywalking-logging-api ${project.version} + + com.a.eye + skywalking-logging-impl-log4j2 + ${project.version} + diff --git a/skywalking-registry/src/main/java/com/a/eye/skywalking/registry/api/NotifyListener.java b/skywalking-registry/src/main/java/com/a/eye/skywalking/registry/api/NotifyListener.java index 7cc9a07e9..0e1ea799c 100644 --- a/skywalking-registry/src/main/java/com/a/eye/skywalking/registry/api/NotifyListener.java +++ b/skywalking-registry/src/main/java/com/a/eye/skywalking/registry/api/NotifyListener.java @@ -1,5 +1,5 @@ package com.a.eye.skywalking.registry.api; public interface NotifyListener { - void notify(EventType type, String urls); + void notify(EventType type, String url); } diff --git a/skywalking-registry/src/main/java/com/a/eye/skywalking/registry/impl/zookeeper/ZookeeperRegistryCenter.java b/skywalking-registry/src/main/java/com/a/eye/skywalking/registry/impl/zookeeper/ZookeeperRegistryCenter.java index 791b0e727..f2d4132f4 100644 --- a/skywalking-registry/src/main/java/com/a/eye/skywalking/registry/impl/zookeeper/ZookeeperRegistryCenter.java +++ b/skywalking-registry/src/main/java/com/a/eye/skywalking/registry/impl/zookeeper/ZookeeperRegistryCenter.java @@ -4,8 +4,10 @@ import com.a.eye.skywalking.logging.api.ILog; import com.a.eye.skywalking.logging.api.LogManager; import com.a.eye.skywalking.registry.api.*; import org.apache.zookeeper.*; +import org.apache.zookeeper.data.Stat; import java.io.IOException; +import java.util.ArrayList; import java.util.List; import java.util.Properties; @@ -23,26 +25,37 @@ public class ZookeeperRegistryCenter implements RegistryCenter { createPath = "/" + createPath; } - mkdirs(createPath); + mkdirs(createPath, true); } - private void mkdirs(String path) { - try { + private void mkdirs(String path, boolean bool) { + try { String[] pathArray = path.split("/"); if (pathArray.length == 0) { return; } + StringBuilder currentCreatePath = new StringBuilder(); for (int i = 0; i < pathArray.length - 1; i++) { String pathSegment = pathArray[i]; if (pathSegment.length() == 0) { continue; } - client.create(path, null, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); - } - client.create(pathArray[pathArray.length - 1], null, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL); + currentCreatePath.append("/").append(pathSegment); + if (client.exists(currentCreatePath.toString(), false) == null) { + client.create(currentCreatePath.toString(), null, ZooDefs.Ids.OPEN_ACL_UNSAFE, + CreateMode.PERSISTENT); + } + } + if (bool) { + client.create(currentCreatePath.append("/").append(pathArray[pathArray.length - 1]).toString(), null, + ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL); + } else { + client.create(currentCreatePath.append("/").append(pathArray[pathArray.length - 1]).toString(), null, + ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + } logger.info("register path[{}] success", path); } catch (Exception e) { logger.error("Failed to create path[{}]", path, e); @@ -52,20 +65,29 @@ public class ZookeeperRegistryCenter implements RegistryCenter { @Override public void subscribe(final String path, final NotifyListener listener) { try { - List childrenPath = client.getChildren(path, new SubscribeWatcher(path, listener)); - for (String child : childrenPath) { - listener.notify(EventType.Add, child); + if (client.exists(path, false) == null) { + logger.warn("{} was not exists. "); + mkdirs(path, false); } + + client.getChildren(path, new SubscribeWatcher(path, listener), new AsyncCallback.Children2Callback() { + @Override + public void processResult(int rc, String path, Object ctx, List children, Stat stat) { + for (String child : children) { + listener.notify(EventType.Add, child); + } + } + }, null); } catch (Exception e) { logger.error("Failed to subscribe the path {} ", path, e); } } @Override - public void start(Properties centerConfig) { - ZookeeperConfig config = new ZookeeperConfig(centerConfig); + public void start(final Properties centerConfig) { + final ZookeeperConfig config = new ZookeeperConfig(centerConfig); try { - client = new ZooKeeper(config.getConnectURL(), 60 * 1000, null); + client = new ZooKeeper(config.getConnectURL(), 60 * 1000, new ConnectWatcher(config)); if (config.hasAuthInfo()) { client.addAuthInfo(config.getAutSchema(), config.getAuth()); } @@ -74,20 +96,89 @@ public class ZookeeperRegistryCenter implements RegistryCenter { } } + private class RetryConnected implements Runnable { + + private ZookeeperConfig config; + + public RetryConnected(ZookeeperConfig config) { + this.config = config; + } + + @Override + public void run() { + while (true) { + try { + client = new ZooKeeper(config.getConnectURL(), 60 * 1000, new ConnectWatcher(config)); + } catch (Exception e) { + logger.error("failed to connect zookeeper", e); + } + + if (client.getState() == ZooKeeper.States.CONNECTED) { + logger.info("connected successfully!"); + break; + } + + try { + Thread.sleep(60 * 1000); + } catch (InterruptedException e) { + logger.error("Failed to sleep.", e); + } + } + + } + } + + + private class ConnectWatcher implements Watcher { + + private ZookeeperConfig config; + + public ConnectWatcher(ZookeeperConfig config) { + this.config = config; + } + + @Override + public void process(WatchedEvent watchedEvent) { + if (watchedEvent.getState() == Event.KeeperState.AuthFailed) { + logger.warn("failed to auth.auth url: {} auth schema:{} auth info:{}", config.getConnectURL(), + config.getAutSchema(), new String(config.getAuth())); + } + + if (watchedEvent.getState() == Event.KeeperState.Disconnected) { + logger.warn("Disconnected from zookeeper. retry connecting..."); + new Thread(new RetryConnected(config)).start(); + } + } + } + private class SubscribeWatcher implements Watcher { private String path; private NotifyListener listener; + private List previousChildPath; + public SubscribeWatcher(String path, NotifyListener listener) { this.path = path; this.listener = listener; + previousChildPath = new ArrayList(); } @Override public void process(WatchedEvent event) { - retryWatch(); + try { + client.getChildren(path, this); + + client.getChildren(path, false, new AsyncCallback.Children2Callback() { + @Override + public void processResult(int rc, String path, Object ctx, List children, Stat stat) { + System.out.println("aaaa"); + } + }, null); + }catch (Exception e){ + + } if (event.getType() == Event.EventType.NodeChildrenChanged) { notifyListener(event); @@ -96,11 +187,10 @@ public class ZookeeperRegistryCenter implements RegistryCenter { private void notifyListener(WatchedEvent event) { try { - List tmpChildrenPath = client.getChildren(path, null); - if (tmpChildrenPath.contains(event.getPath())) { - listener.notify(EventType.Add, event.getPath()); - } else { - listener.notify(EventType.Remove, event.getPath()); + List tmpChildrenPath = client.getChildren(path, false); + tmpChildrenPath.removeAll(previousChildPath); + if (tmpChildrenPath.size() == 0) { + } } catch (Exception e) { logger.error("Failed to fetch path[{}] children.", path, e); diff --git a/skywalking-registry/src/test/java/com/a/eye/skywalking/registry/RegistryCenterFactoryTest.java b/skywalking-registry/src/test/java/com/a/eye/skywalking/registry/RegistryCenterFactoryTest.java new file mode 100644 index 000000000..595e62b65 --- /dev/null +++ b/skywalking-registry/src/test/java/com/a/eye/skywalking/registry/RegistryCenterFactoryTest.java @@ -0,0 +1,73 @@ +package com.a.eye.skywalking.registry; + +import com.a.eye.skywalking.logging.api.LogManager; +import com.a.eye.skywalking.logging.impl.log4j2.Log4j2Resolver; +import com.a.eye.skywalking.registry.api.CenterType; +import com.a.eye.skywalking.registry.api.EventType; +import com.a.eye.skywalking.registry.api.NotifyListener; +import com.a.eye.skywalking.registry.api.RegistryCenter; +import com.a.eye.skywalking.registry.impl.zookeeper.ZookeeperConfig; +import org.apache.zookeeper.KeeperException; +import org.apache.zookeeper.WatchedEvent; +import org.apache.zookeeper.Watcher; +import org.apache.zookeeper.ZooKeeper; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import java.io.IOException; +import java.util.Properties; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; + +/** + * Created by xin on 2016/11/12. + */ +public class RegistryCenterFactoryTest { + + private RegistryCenter registryCenter; + private ZooKeeper zooKeeper; + + @Before + public void setUp() throws IOException { + LogManager.setLogResolver(new Log4j2Resolver()); + registryCenter = RegistryCenterFactory.INSTANCE.getRegistryCenter(CenterType.DEFAULT_CENTER_TYPE); + Properties config = new Properties(); + config.setProperty(ZookeeperConfig.CONNECT_URL, "127.0.0.1:2181"); + registryCenter.start(config); + zooKeeper = new ZooKeeper("127.0.0.1:2181", 60 * 1000, new Watcher(){ + @Override + public void process(WatchedEvent watchedEvent) { + } + }); + } + + @Test + public void testRegistry() throws KeeperException, InterruptedException { + registryCenter.register("/a/b/c"); + assertNotNull(zooKeeper.exists("/a/b/c",false)); + } + + @After + public void clearUp() throws KeeperException, InterruptedException { + //zooKeeper.delete("/a", -1); + } + + @Test + public void testSubscribe(){ + registryCenter.subscribe("/a", new NotifyListener() { + @Override + public void notify(EventType type, String urls) { + assertEquals(type, EventType.Add); + assertEquals(urls,"b"); + } + }); + + registryCenter.register("/a/b"); + + registryCenter.register("/a/d"); + + registryCenter.register("/a/e"); + } +} diff --git a/skywalking-registry/src/test/resources/log4j.properties b/skywalking-registry/src/test/resources/log4j.properties new file mode 100644 index 000000000..14cd87187 --- /dev/null +++ b/skywalking-registry/src/test/resources/log4j.properties @@ -0,0 +1,13 @@ +# log4j properties # + +# logger # +log4j.rootLogger=DEBUG,CONSOLE +log4j.logger.org=ON +#log4j.logger.org.systemgo.devframework=DEBUG + +# Console Appender # +log4j.appender.CONSOLE=org.apache.log4j.ConsoleAppender +log4j.appender.CONSOLE.Target=System.out +log4j.appender.CONSOLE.layout=org.apache.log4j.PatternLayout +log4j.appender.CONSOLE.layout.ConversionPattern=%d %-5p %c{1}:%L - %m%n + diff --git a/skywalking-registry/src/test/resources/log4j2.xml b/skywalking-registry/src/test/resources/log4j2.xml new file mode 100644 index 000000000..91e20d37a --- /dev/null +++ b/skywalking-registry/src/test/resources/log4j2.xml @@ -0,0 +1,13 @@ + + + + + + + + + + + + + \ No newline at end of file diff --git a/skywalking-storage-center/pom.xml b/skywalking-storage-center/pom.xml index 80858de1e..0a5ed77d9 100644 --- a/skywalking-storage-center/pom.xml +++ b/skywalking-storage-center/pom.xml @@ -17,4 +17,18 @@ skywalking-storage skywalking-routing + + + + com.a.eye + skywalking-logging + ${parent.version} + + + + com.a.eye + skywalking-network + ${parent.version} + + diff --git a/skywalking-storage-center/skywalking-storage/bin/swserver.sh b/skywalking-storage-center/skywalking-storage/bin/swserver.sh deleted file mode 100644 index ecd8b5fba..000000000 --- a/skywalking-storage-center/skywalking-storage/bin/swserver.sh +++ /dev/null @@ -1,52 +0,0 @@ -#!/bin/sh - -# OS specific support. $var _must_ be set to either true or false. -cygwin=false -os400=false -darwin=false -case "`uname`" in -CYGWIN*) cygwin=true;; -OS400*) os400=true;; -Darwin*) darwin=true;; -esac - -# resolve links - $0 may be a softlink -SW_SERVER_BIN="$0" - -while [ -h "$SW_SERVER_BIN" ]; do - ls=`ls -ld "$SW_SERVER_BIN"` - link=`expr "$ls" : '.*-> \(.*\)$'` - if expr "$link" : '/.*' > /dev/null; then - SW_SERVER_BIN="$link" - else - SW_SERVER_BIN=`dirname "$SW_SERVER_BIN"`/"$link" - fi -done - -# Get standard environment variables -SW_SERVER_BIN_DIR=`dirname "$SW_SERVER_BIN"` -SW_PREFIX="${SW_SERVER_BIN_DIR}/.." -SW_LOG_DIR="${SW_SERVER_BIN_DIR}/../log" -SW_CFG_DIR="${SW_SERVER_BIN_DIR}/../config" - -#echo $SW_SERVER_BIN_DIR -#set java home -if [ "$JAVA_HOME" != "" ]; then - JAVA="$JAVA_HOME/bin/java" -else - JAVA=java -fi - -CLASSPATH="$SW_CFG_DIR:$CLASSPATH" - -for i in "${SW_SERVER_BIN_DIR}"/../lib/*.jar -do - CLASSPATH="$i:$CLASSPATH" -done - -echo "CLASSPATH=$CLASSPATH" - -JAVA_OPTS="$JAVA_OPTS -Dcom.sun.management.jmxremote.ssl=false -Dcom.sun.management.jmxremote.authenticate=false" - -$JAVA ${JAVA_OPTS} -classpath $CLASSPATH com.a.eye.skywalking.reciever.CollectionServer >> ${SW_SERVER_BIN_DIR}/. -./log/sw-server.log 2>&1 & 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 new file mode 100644 index 000000000..d57372674 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java @@ -0,0 +1,62 @@ +package com.a.eye.skywalking.storage; + +import com.a.eye.skywalking.logging.api.ILog; +import com.a.eye.skywalking.logging.api.LogManager; +import com.a.eye.skywalking.logging.impl.log4j2.Log4j2Resolver; +import com.a.eye.skywalking.network.TransferService; +import com.a.eye.skywalking.network.TransferService.TransferServiceBuilder; +import com.a.eye.skywalking.storage.config.Config; +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 java.io.IOException; +import java.util.Properties; + +/** + * Created by xin on 2016/11/12. + */ +public class Main { + + private static ILog logger = LogManager.getLogger(Main.class); + + static { + LogManager.setLogResolver(new Log4j2Resolver()); + } + + private static TransferService transferService; + + public static void main(String[] args) { + try { + initializeParam(); + + 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(); + logger.info("storage service started successfully!"); + Thread.currentThread().join(); + } catch (Throwable e) { + logger.error("Failed to start service.", e); + } finally { + transferService.stop(); + } + } + + private static void initializeParam() throws IllegalAccessException, IOException { + Properties properties = new Properties(); + try { + properties.load(Main.class.getResourceAsStream("/config.properties")); + ConfigInitializer.initialize(properties, Config.class); + } catch (IllegalAccessException e) { + logger.error("Initialize the collect server configuration failed", e); + throw e; + } catch (IOException e) { + logger.error("Initialize the collect server configuration failed", e); + throw e; + } + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/block/index/BlockFinder.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/block/index/BlockFinder.java index 4dc21537d..440aaf2cc 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/block/index/BlockFinder.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/block/index/BlockFinder.java @@ -19,10 +19,6 @@ public class BlockFinder { index = l2Cache.find(timestamp); } - if (index == null) { - index = System.currentTimeMillis(); - } - return index; } 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 580fc3bfc..b1e984f61 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 @@ -4,16 +4,20 @@ package com.a.eye.skywalking.storage.config; * Created by xin on 2016/11/2. */ public class Config { + public static class Server { + public static int PORT = 34000; + } + public static class BlockIndex { - public static String STORAGE_BASE_PATH = "/tmp/skywalking/index"; + public static String STORAGE_BASE_PATH = "/tmp/skywalking/block_index"; public static String DATA_FILE_INDEX_FILE_NAME = "data_file.index"; } public static class DataFile { - public static String BASE_PATH = ""; + public static String BASE_PATH = "/tmp/skywalking/data/file"; public static long MAX_LENGTH = 3 * 1024 * 1024 * 1024; } @@ -23,9 +27,9 @@ public class Config { public static String TABLE_NAME = "data_index"; - public static String BASE_PATH = ""; + public static String BASE_PATH = "/tmp/skywalking/data/index"; - public static String STORAGE_INDEX_FILE_NAME = ""; + public static String STORAGE_INDEX_FILE_NAME = "dataIndex"; public static long MAX_CAPACITY_PER_INDEX = 1000 * 1000 * 1000 * 1000; } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/ConfigInitializer.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/ConfigInitializer.java new file mode 100644 index 000000000..1a5ea4ad0 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/ConfigInitializer.java @@ -0,0 +1,70 @@ +package com.a.eye.skywalking.storage.config; + +import java.lang.reflect.Field; +import java.lang.reflect.Modifier; +import java.util.LinkedList; +import java.util.Properties; +import java.util.logging.Logger; + +public class ConfigInitializer { + private static Logger logger = Logger.getLogger(ConfigInitializer.class.getName()); + + public static void initialize(Properties properties, Class rootConfigType) throws IllegalAccessException { + initNextLevel(properties, rootConfigType, new ConfigDesc()); + } + + private static void initNextLevel(Properties properties, Class recentConfigType, ConfigDesc parentDesc) throws NumberFormatException, IllegalArgumentException, IllegalAccessException { + for (Field field : recentConfigType.getFields()) { + if (Modifier.isPublic(field.getModifiers()) && Modifier.isStatic(field.getModifiers())) { + String configKey = (parentDesc + "." + + field.getName()).toLowerCase(); + String value = properties.getProperty(configKey); + if (value != null) { + if (field.getType().equals(int.class)) + field.set(null, Integer.valueOf(value)); + if (field.getType().equals(String.class)) + field.set(null, value); + if (field.getType().equals(long.class)) + field.set(null, Long.valueOf(value)); + if (field.getType().equals(boolean.class)) + field.set(null, Boolean.valueOf(value)); + } + } + } + for (Class innerConfiguration : recentConfigType.getClasses()) { + parentDesc.append(innerConfiguration.getSimpleName()); + initNextLevel(properties, innerConfiguration, parentDesc); + parentDesc.removeLastDesc(); + } + } +} + + +class ConfigDesc { + private LinkedList descs = new LinkedList(); + + void append(String currentDesc) { + descs.addLast(currentDesc); + } + + void removeLastDesc() { + descs.removeLast(); + } + + @Override + public String toString() { + if (descs.size() == 0) { + return ""; + } + StringBuilder ret = new StringBuilder(descs.getFirst()); + boolean first = true; + for (String desc : descs) { + if (first) { + first = false; + continue; + } + ret.append(".").append(desc); + } + return ret.toString(); + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Constants.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Constants.java index 322558054..0ddabbd06 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Constants.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Constants.java @@ -8,21 +8,24 @@ public class Constants { public static final String CREATE_TABLE = "CREATE TABLE " + TABLE_NAME + "\n" + "(\n" + " id INT PRIMARY KEY NOT NULL IDENTITY,\n" + " trace_id VARCHAR(32) NOT NULL,\n" - + " parent_level_id VARCHAR(1024) NOT NULL,\n" - + " level_id INT NOT NULL,\n" + + " levelId VARCHAR(1024) NOT NULL,\n" + + " span_type INT NOT NULL, \n" + " file_name VARCHAR(10) NOT NULL,\n" + " offset BIGINT NOT NULL,\n" + " length INT NOT NULL\n" + ");\n"; public static final String CREATE_INDEX = "CREATE INDEX \"index_data_trace_id_index\" ON " + TABLE_NAME + " (trace_id);"; - public static final String INSERT_INDEX = "INSERT INTO " +TABLE_NAME + "(trace_id,parent_level_id,level_id," + public static final String INSERT_INDEX = "INSERT INTO " +TABLE_NAME + "(trace_id,levelId,span_type" + "file_name,offset,length) VALUES(?,?,?,?,?,?)"; public static final String QUERY_TABLES = "SELECT count(1) AS TABLE_COUNT FROM INFORMATION_SCHEMA.TABLES " + "WHERE TABLE_NAME= '" + TABLE_NAME + "';"; public static final String QUERY_INDEX_SIZE = "SELECT count(1) AS INDEX_SIZE FROM " + TABLE_NAME; + + public static final String QUERY_TRACE_ID = "SELECT span_type, file_name, offset, length " + + " FROM "+ TABLE_NAME+ " WHERE trace_id = ?"; } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/IndexDataCapacityMonitor.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/IndexDataCapacityMonitor.java index 066c09619..4554628fe 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/IndexDataCapacityMonitor.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/IndexDataCapacityMonitor.java @@ -13,7 +13,7 @@ import static com.a.eye.skywalking.storage.config.Config.DataIndex.MAX_CAPACITY_ /** * Created by xin on 2016/11/6. */ -public class IndexDataCapacityMonitor extends Thread { +public class IndexDataCapacityMonitor implements Runnable { private static ILog logger = LogManager.getLogger(IndexDataCapacityMonitor.class); private static Detector detector; @@ -76,6 +76,7 @@ public class IndexDataCapacityMonitor extends Thread { logger.error("Failed to to fetch index size from DB:{}", timestamp, e); } detector = new Detector(timestamp, count); + logger.info("Index data capacity monitor started successfully!"); } finally { if (dbConnector != null) { dbConnector.close(); diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanData.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanData.java deleted file mode 100644 index 6828bd10c..000000000 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanData.java +++ /dev/null @@ -1,8 +0,0 @@ -package com.a.eye.skywalking.storage.data; - -public interface SpanData { - - long getStartTime(); - - byte[] toByteArray(); -} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java index 57a555e6e..aab9fc9f6 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java @@ -1,8 +1,10 @@ package com.a.eye.skywalking.storage.data; import com.a.eye.datacarrier.consumer.IConsumer; +import com.a.eye.skywalking.storage.block.index.BlockIndexEngine; import com.a.eye.skywalking.storage.data.file.DataFileWriter; import com.a.eye.skywalking.storage.data.index.*; +import com.a.eye.skywalking.storage.data.spandata.SpanData; import java.util.Iterator; import java.util.List; @@ -24,7 +26,7 @@ public class SpanDataConsumer implements IConsumer { IndexMetaCollections.group(fileWriter.write(data), new GroupKeyBuilder() { @Override public Long buildKey(IndexMetaInfo metaInfo) { - return metaInfo.getStartTime(); + return BlockIndexEngine.newFinder().find(metaInfo.getTraceStartTime()); } }).iterator(); 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 9dfd66d3d..897d2473f 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 @@ -2,20 +2,19 @@ package com.a.eye.skywalking.storage.data; import com.a.eye.skywalking.storage.block.index.BlockIndexEngine; import com.a.eye.skywalking.storage.data.file.DataFileReader; -import com.a.eye.skywalking.storage.data.file.DataFileWriter; import com.a.eye.skywalking.storage.data.index.*; +import com.a.eye.skywalking.storage.data.spandata.SpanData; import java.util.ArrayList; import java.util.Iterator; import java.util.List; -/** - * Created by xin on 2016/11/6. - */ public class SpanDataFinder { - - public static List find(String traceId) { + public static List find(String traceId) { long blockIndex = BlockIndexEngine.newFinder().find(fetchStartTimeFromTraceId(traceId)); + if (blockIndex == 0) { + return new ArrayList(); + } IndexDBConnector indexDBConnector = new IndexDBConnector(blockIndex); IndexMetaCollection indexMetaCollection = indexDBConnector.queryByTraceId(traceId); @@ -27,7 +26,7 @@ public class SpanDataFinder { } }).iterator(); - List result = new ArrayList(); + List result = new ArrayList(); while (iterator.hasNext()) { IndexMetaGroup group = iterator.next(); result.addAll(new DataFileReader(group.getKey()).read(group.getMetaInfo())); @@ -37,6 +36,7 @@ public class SpanDataFinder { } private static long fetchStartTimeFromTraceId(String traceId) { - return -1; + String[] traceIdSegment = traceId.split("\\."); + return Long.parseLong(traceIdSegment[traceIdSegment.length - 5]); } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFile.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFile.java index d057edadc..7546fd6a4 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFile.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFile.java @@ -1,7 +1,7 @@ package com.a.eye.skywalking.storage.data.file; import com.a.eye.skywalking.storage.config.Config; -import com.a.eye.skywalking.storage.data.SpanData; +import com.a.eye.skywalking.storage.data.spandata.SpanData; import com.a.eye.skywalking.storage.data.exception.DataFileOperatorCreateFailedException; import com.a.eye.skywalking.storage.data.exception.SpanDataPersistenceFailedException; import com.a.eye.skywalking.storage.data.exception.SpanDataReadFailedException; @@ -46,7 +46,7 @@ public class DataFile { byte[] bytes = data.toByteArray(); try { operator.getWriter().write(bytes); - IndexMetaInfo metaInfo = new IndexMetaInfo(fileName, currentOffset, bytes.length); + IndexMetaInfo metaInfo = new IndexMetaInfo(data,fileName, currentOffset, bytes.length); currentOffset += bytes.length; return metaInfo; } catch (IOException e) { diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileReader.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileReader.java index c0b4ba40a..e237401c7 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileReader.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileReader.java @@ -1,26 +1,43 @@ package com.a.eye.skywalking.storage.data.file; +import com.a.eye.skywalking.logging.api.ILog; +import com.a.eye.skywalking.logging.api.LogManager; +import com.a.eye.skywalking.network.grpc.AckSpan; +import com.a.eye.skywalking.network.grpc.RequestSpan; import com.a.eye.skywalking.storage.data.index.IndexMetaInfo; +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 com.a.eye.skywalking.storage.data.spandata.SpanType; import java.util.ArrayList; -import java.util.Collection; import java.util.List; /** * Created by xin on 2016/11/6. */ public class DataFileReader { + private static ILog logger = LogManager.getLogger(DataFileReader.class); private DataFile dataFile; public DataFileReader(String fileName) { dataFile = new DataFile(fileName); } - public List read(List metaInfo) { - List metaData = new ArrayList(); + public List read(List metaInfo) { + List metaData = new ArrayList(); - for (IndexMetaInfo indexMetaInfo : metaInfo){ - metaData.add(dataFile.read(indexMetaInfo.getOffset(), indexMetaInfo.getLength())); + for (IndexMetaInfo indexMetaInfo : metaInfo) { + byte[] dataByte = dataFile.read(indexMetaInfo.getOffset(), indexMetaInfo.getLength()); + try { + if (indexMetaInfo.getSpanType() == SpanType.RequestSpan) { + metaData.add(new RequestSpanData(RequestSpan.parseFrom(dataByte))); + } else { + metaData.add(new AckSpanData(AckSpan.parseFrom(dataByte))); + } + } catch (Exception e) { + logger.error("Failed to conver to data", e); + } } return metaData; diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java index a6878bc83..0b2587cf2 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java @@ -1,6 +1,6 @@ package com.a.eye.skywalking.storage.data.file; -import com.a.eye.skywalking.storage.data.SpanData; +import com.a.eye.skywalking.storage.data.spandata.SpanData; import com.a.eye.skywalking.storage.data.index.IndexMetaCollection; import java.util.List; 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 8914ea261..3a32e047c 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 @@ -4,6 +4,10 @@ import com.a.eye.skywalking.logging.api.ILog; import com.a.eye.skywalking.logging.api.LogManager; import com.a.eye.skywalking.storage.config.Config; import com.a.eye.skywalking.storage.data.exception.ConnectorInitializeFailedException; +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 com.a.eye.skywalking.storage.data.spandata.SpanType; import java.sql.*; @@ -91,8 +95,8 @@ public class IndexDBConnector { PreparedStatement ps = connection.prepareStatement(INSERT_INDEX); for (IndexMetaInfo metaInfo : metaGroup.getMetaInfo()) { ps.setString(1, metaInfo.getTraceId()); - ps.setString(2, metaInfo.getParentLevelId()); - ps.setInt(3, metaInfo.getLevelId()); + ps.setString(2, metaInfo.getLevelId()); + ps.setInt(3, metaInfo.getSpanType().getValue()); ps.setString(4, metaInfo.getFileName()); ps.setLong(5, metaInfo.getOffset()); ps.setInt(6, metaInfo.getLength()); @@ -117,8 +121,29 @@ public class IndexDBConnector { return indexSize; } - public IndexMetaCollection queryByTraceId(String traceId) { - return null; + public IndexMetaCollection queryByTraceId(String traceId){ + try { + PreparedStatement ps = connection.prepareStatement(QUERY_TRACE_ID); + ps.setString(1, traceId); + ResultSet rs = ps.executeQuery(); + + IndexMetaCollection collection = new IndexMetaCollection(); + while (rs.next()) { + SpanType spanType = SpanType.convert(rs.getInt("span_type")); + SpanData spanData = null; + + if (SpanType.ACKSpan == spanType) { + spanData = new AckSpanData(); + } else if (SpanType.RequestSpan == spanType) { + spanData = new RequestSpanData(); + } + + collection.add(new IndexMetaInfo(spanData, rs.getString("file_name"), rs.getLong("offset"), rs.getInt("length"))); + } + return collection; + }catch(SQLException e){ + return new IndexMetaCollection(); + } } class ConnectURLGenerator { diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaInfo.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaInfo.java index 70d257362..2eab368e6 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaInfo.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaInfo.java @@ -1,7 +1,11 @@ package com.a.eye.skywalking.storage.data.index; +import com.a.eye.skywalking.storage.data.spandata.SpanData; +import com.a.eye.skywalking.storage.data.spandata.SpanType; + public class IndexMetaInfo { - private String traceId; + + private SpanData spanData; private String fileName; @@ -9,9 +13,8 @@ public class IndexMetaInfo { private int length; - private long startTime; - - public IndexMetaInfo(String fileName, long offset, int length) { + public IndexMetaInfo(SpanData data, String fileName, long offset, int length) { + this.spanData = data; this.fileName = fileName; this.offset = offset; this.length = length; @@ -29,19 +32,19 @@ public class IndexMetaInfo { return length; } - public long getStartTime() { - return startTime; + public long getTraceStartTime() { + return spanData.getTraceStartTime(); } public String getTraceId() { - return null; + return spanData.getTraceId(); } - public String getParentLevelId() { - return null; + public String getLevelId() { + return spanData.getLevelId(); } - public int getLevelId() { - return 0; + public SpanType getSpanType() { + return spanData.getSpanType(); } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AbstractSpanData.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AbstractSpanData.java new file mode 100644 index 000000000..2f703a9ad --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AbstractSpanData.java @@ -0,0 +1,16 @@ +package com.a.eye.skywalking.storage.data.spandata; + +/** + * Created by xin on 2016/11/12. + */ +public abstract class AbstractSpanData implements SpanData{ + + protected String buildLevelId(String parentLevelId, int levelId) { + return (parentLevelId == null || parentLevelId.length() == 0) ? levelId + "" : parentLevelId + "." + levelId; + } + + protected static long buildTraceStartTime(String traceId) { + String[] traceIdSegment = traceId.split("\\."); + return Long.parseLong(traceIdSegment[traceIdSegment.length - 5]); + } +} 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 new file mode 100644 index 000000000..c9fe3ce57 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AckSpanData.java @@ -0,0 +1,43 @@ +package com.a.eye.skywalking.storage.data.spandata; + +import com.a.eye.skywalking.network.grpc.AckSpan; + +/** + * Created by xin on 2016/11/12. + */ +public class AckSpanData extends AbstractSpanData { + private AckSpan ackSpan; + + public AckSpanData(AckSpan ackSpan) { + this.ackSpan = ackSpan; + } + + public AckSpanData() { + } + + @Override + public SpanType getSpanType() { + return SpanType.ACKSpan; + } + + @Override + public long getTraceStartTime() { + return buildTraceStartTime(ackSpan.getTraceId()); + } + + @Override + public byte[] toByteArray() { + return ackSpan.toByteArray(); + } + + @Override + public String getTraceId() { + return ackSpan.getTraceId(); + } + + @Override + public String getLevelId() { + return buildLevelId(ackSpan.getParentLevel(), ackSpan.getLevelId()); + } + +} 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 new file mode 100644 index 000000000..41dd932b2 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/RequestSpanData.java @@ -0,0 +1,42 @@ +package com.a.eye.skywalking.storage.data.spandata; + +import com.a.eye.skywalking.network.grpc.RequestSpan; + +/** + * Created by xin on 2016/11/12. + */ +public class RequestSpanData extends AbstractSpanData { + private RequestSpan requestSpan; + + public RequestSpanData(RequestSpan requestSpan) { + this.requestSpan = requestSpan; + } + + public RequestSpanData() { + } + + @Override + public SpanType getSpanType() { + return SpanType.RequestSpan; + } + + @Override + public long getTraceStartTime() { + return buildTraceStartTime(requestSpan.getTraceId()); + } + + @Override + public byte[] toByteArray() { + return requestSpan.toByteArray(); + } + + @Override + public String getTraceId() { + return requestSpan.getTraceId(); + } + + @Override + public String getLevelId() { + return buildLevelId(requestSpan.getParentLevel(), requestSpan.getLevelId()); + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanData.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanData.java new file mode 100644 index 000000000..8e8490413 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanData.java @@ -0,0 +1,14 @@ +package com.a.eye.skywalking.storage.data.spandata; + +public interface SpanData { + + SpanType getSpanType(); + + long getTraceStartTime(); + + byte[] toByteArray(); + + String getTraceId(); + + String getLevelId(); +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanDataBuilder.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanDataBuilder.java new file mode 100644 index 000000000..0d524167c --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanDataBuilder.java @@ -0,0 +1,44 @@ +package com.a.eye.skywalking.storage.data.spandata; + +import com.a.eye.skywalking.logging.api.ILog; +import com.a.eye.skywalking.logging.api.LogManager; +import com.a.eye.skywalking.network.dependencies.com.google.protobuf.InvalidProtocolBufferException; +import com.a.eye.skywalking.network.grpc.AckSpan; +import com.a.eye.skywalking.network.grpc.RequestSpan; +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; + +/** + * Created by xin on 2016/11/12. + */ +public class SpanDataBuilder { + + private static ILog logger = LogManager.getLogger(SpanDataBuilder.class); + + public static SpanData build(RequestSpan requestSpan) { + return new RequestSpanData(requestSpan); + } + + public static SpanData build(AckSpan ackSpan) { + return new AckSpanData(ackSpan); + } + + public static AckSpan buildAckSpan(byte[] data) { + try { + return AckSpan.parseFrom(data); + } catch (InvalidProtocolBufferException e) { + logger.error("Failed to convert data to ack span.", e); + return null; + } + } + + public static RequestSpan buildRequestSpan(byte[] data) { + try { + return RequestSpan.parseFrom(data); + } catch (InvalidProtocolBufferException e) { + logger.error("Failed to convert data to request span.", e); + return null; + } + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanType.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanType.java new file mode 100644 index 000000000..7aa5f56cd --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/SpanType.java @@ -0,0 +1,30 @@ +package com.a.eye.skywalking.storage.data.spandata; + +/** + * Created by xin on 2016/11/12. + */ +public enum SpanType { + RequestSpan(1), + ACKSpan(2); + + int value; + + SpanType(int value) { + this.value = value; + } + + public int getValue() { + return value; + } + + public static SpanType convert(int value) { + switch (value) { + case 1: + return RequestSpan; + case 2: + return ACKSpan; + default: + throw new IllegalArgumentException("Failed to convert to value" + value); + } + } +} 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 new file mode 100644 index 000000000..63ae56435 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/SearchNotifier.java @@ -0,0 +1,52 @@ +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); + } + + 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/StorageNotifier.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/StorageNotifier.java new file mode 100644 index 000000000..d309449a3 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/notifier/StorageNotifier.java @@ -0,0 +1,45 @@ +package com.a.eye.skywalking.storage.notifier; + +import com.a.eye.datacarrier.DataCarrier; +import com.a.eye.skywalking.logging.api.ILog; +import com.a.eye.skywalking.logging.api.LogManager; +import com.a.eye.skywalking.network.grpc.AckSpan; +import com.a.eye.skywalking.network.grpc.RequestSpan; +import com.a.eye.skywalking.network.listener.SpanStorageNotifier; +import com.a.eye.skywalking.storage.data.spandata.SpanData; +import com.a.eye.skywalking.storage.data.spandata.SpanDataBuilder; +import com.a.eye.skywalking.storage.data.SpanDataConsumer; + +public class StorageNotifier implements SpanStorageNotifier { + + private ILog logger = LogManager.getLogger(StorageNotifier.class); + + private DataCarrier spanDataDataCarrier; + + public StorageNotifier() { + spanDataDataCarrier = new DataCarrier<>(10, 1000); + spanDataDataCarrier.consume(new SpanDataConsumer(), 5, true); + } + + @Override + public boolean storage(RequestSpan requestSpan) { + try { + spanDataDataCarrier.produce(SpanDataBuilder.build(requestSpan)); + return true; + } catch (Exception e) { + logger.error("Failed to storage request span. Span Data:\n {}.", requestSpan.toByteString(), e); + return false; + } + } + + @Override + public boolean storage(AckSpan ackSpan) { + try { + spanDataDataCarrier.produce(SpanDataBuilder.build(ackSpan)); + return true; + } catch (Exception e) { + logger.error("Failed to storage ack span. ack Data:\n {}.", ackSpan.toByteString(), e); + return false; + } + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/resources/META-INF/services/com.a.eye.skywalking.reciever.processor.IProcessor b/skywalking-storage-center/skywalking-storage/src/main/resources/META-INF/services/com.a.eye.skywalking.reciever.processor.IProcessor deleted file mode 100644 index 097774748..000000000 --- a/skywalking-storage-center/skywalking-storage/src/main/resources/META-INF/services/com.a.eye.skywalking.reciever.processor.IProcessor +++ /dev/null @@ -1,4 +0,0 @@ -com.a.eye.skywalking.reciever.processor.RequestSpanProcessor -com.a.eye.skywalking.reciever.processor.AckSpanProcessor -com.a.eye.skywalking.reciever.processor.InputParameterSpanProcessor -com.a.eye.skywalking.reciever.processor.OutputParameterSpanProcessor 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 6d3e561a8..ff4a07e5c 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/resources/config.properties +++ b/skywalking-storage-center/skywalking-storage/src/main/resources/config.properties @@ -1,72 +1,21 @@ -#采集服务器的端口 -server.port=34000 - -server.failed_package_watching_time_windowss=300 -server.max_watching_failed_package_size=200 # -buffer.buffer_deal_thread_number=1 -#每个线程最大缓存数量 -buffer.per_thread_max_buffer_number=1024 -#无数据处理时轮询等待时间(单位:毫秒) -buffer.max_wait_time=5000 -#数据冲突时等待时间(单位:毫秒) -buffer.data_conflict_wait_time=10 -#数据缓存文件目录 -buffer.data_buffer_file_parent_dir=/tmp/skywalking/data/buffer -#缓存数据文件最大长度(单位:byte) -buffer.buffer_file_max_length=104857600 -#每次缓存数据写入失败,最大尝试时间 -buffer.write_data_failure_retry_interval = 10000 +server.port = 34000 -#切换数据文件,等待时间(单位:毫秒) -persistence.switch_file_wait_time=5000 -#追加EOF标志位的线程数量 -persistence.max_append_eof_flags_thread_number=1 -#持久化线程个数 -persistence.max_deal_data_thread_number=1 +# +blockindex.storage_base_path= /tmp/skywalking/block-index +# +blockindex.data_file_index_file_name= data_file.index -#偏移量注册文件的目录 -registerpersistence.register_file_parent_directory=/tmp/skywalking/data/offset -#偏移量注册文件名 -registerpersistence.register_file_name=offset.txt -#偏移量注册备份文件名 -registerpersistence.register_bak_file_name=offset.txt.bak -#偏移量写入文件等待周期(单位:毫秒) -registerpersistence.offset_written_file_wait_cycle=5000 +# +datafile.base_path= /tmp/skywalking/data/file +# +datafile.max_length= 3221225472 - -#trace data hbase表名 -hbaseconfig.tracedatatable.table_name=trace-data -#trace data hbase列簇名字 -hbaseconfig.tracedatatable.family_column_name=call-chain - -#trace data hbase表名 -hbaseconfig.traceparamtable.table_name=trace-param -#trace data hbase列簇名字 -hbaseconfig.traceparamtable.family_column_name=param-data - -#hbase zk quorum -hbaseconfig.zk_hostname=swhbaseenv -#hbase zk port -hbaseconfig.client_port=2181 - -#告警失效时间 -alarm.alarm_expire_seconds=5400 -#是否关闭告警 -alarm.larm_off_flag=false -#告警redis检测器检测周期 -alarm.alarm_redis_inspector_interval=5000 - -#Redis配置 -redis.redis_server=10.1.241.18:16379 -#Redis最大空闲数量 -redis.edis_max_idle=10 -#Redis最小空闲数量 -redis.edis_min_idle=1 -#Redis最大个数 -redis.edis_max_total=20 - -#告警检查器:异常告警检查 -alarm.checker.turn_on_exception_checker=true -#告警检查器:执行时间超时告警检查 -alarm.checker.turn_on_execute_time_checker=true +#存放数据文件索引表名 +dataindex.table_name= data_index +#数据文件索引存储位置 +dataindex.base_path= /tmp/skywalking/data/index +# +dataindex.storage_index_file_name= dataIndex +# +dataindex.max_capacity_per_index= 1000000000 diff --git a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/HyperSqlSearchSpeedReporter.java b/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/HyperSqlSearchSpeedReporter.java deleted file mode 100644 index a55c386f5..000000000 --- a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/HyperSqlSearchSpeedReporter.java +++ /dev/null @@ -1,92 +0,0 @@ -package com.a.eye.skywalking.search; - -import com.zaxxer.hikari.HikariConfig; -import com.zaxxer.hikari.HikariDataSource; - -import java.sql.Connection; -import java.sql.PreparedStatement; -import java.sql.ResultSet; -import java.sql.SQLException; - -/** - * Created by xin on 2016/11/1. - */ -public class HyperSqlSearchSpeedReporter { - - private static final long BASE_TIME_STAMP = 1477983548L; - - private static boolean useSingConnection = true; - private static HikariDataSource hikariDataSource; - private static Connection connection; - - private static String CREATE_TABLE_SQL = - "CREATE TABLE data_index\n" + "(\n" + " id INT IDENTITY PRIMARY KEY NOT NULL,\n" - + " startTime BIGINT NOT NULL\n" + ");\n"; - private static String CREATE_INDEX_SQL = - "CREATE UNIQUE INDEX \"table_name_startTime_uindex\" ON data_index (startTime);"; - - private static String INSERT_DATA_SQL = "INSERT INTO data_index(startTime) VALUES(?);"; - - private static String QUERY_DATA_SQL = - "SELECT startTime FROM data_index WHERE startTime > ? ORDER BY startTime" + " ASC LIMIT 1"; - - public static void initData() throws SQLException { - HikariConfig config = new HikariConfig(); - config.setJdbcUrl("jdbc:hsqldb:mem:test-speed"); - config.setUsername("root"); - config.setPassword("root"); - hikariDataSource = new HikariDataSource(config); - connection = hikariDataSource.getConnection(); - - PreparedStatement ps = connection.prepareStatement(CREATE_TABLE_SQL); - ps.execute(); - ps = connection.prepareStatement(CREATE_INDEX_SQL); - ps.execute(); - - ps = connection.prepareStatement(INSERT_DATA_SQL); - for (int i = 0; i < 3000; i++) { - ps.setLong(1, BASE_TIME_STAMP + i * 1000 * 60 * 60); - //System.out.print(BASE_TIME_STAMP + i * 1000 * 60 * 60); - //System.out.print(","); - ps.execute(); - } - - //System.out.println(); - - ps.close(); - } - - public static long find(long element) throws SQLException { - Connection connection = null; - if (!useSingConnection) { - connection = hikariDataSource.getConnection(); - }else{ - connection = HyperSqlSearchSpeedReporter.connection; - } - PreparedStatement preparedStatement = connection.prepareStatement(QUERY_DATA_SQL); - preparedStatement.setLong(1, element); - ResultSet resultSet = preparedStatement.executeQuery(); - resultSet.next(); - long result = resultSet.getLong("startTime"); - preparedStatement.close(); - - if (!useSingConnection){ - connection.close(); - } - - return result; - } - - public static void main(String[] args) throws SQLException { - initData(); - long startTime = System.nanoTime(); - - for (long i = 0; i < 100000000L; i++) { - find(1478323448L); - } - - long totalTime = System.nanoTime() - startTime; - System.out.println("total time : " + totalTime + " " + (totalTime * 1.0 / 100000000L)); - - } -} diff --git a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/SearchSpeedReporter.java b/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/SearchSpeedReporter.java deleted file mode 100644 index acd881d94..000000000 --- a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/SearchSpeedReporter.java +++ /dev/null @@ -1,37 +0,0 @@ -package com.a.eye.skywalking.search; - -import java.util.Arrays; -import java.util.TreeSet; - -public class SearchSpeedReporter { - - private static final long BASE_TIME_STAMP = 1477983548L; - private static Long[] testedData = new Long[3000]; - - private static TreeSet tree = new TreeSet(); - - public static void initData() { - for (int i = 0; i < 3000; i++) { - testedData[i] = new Long(BASE_TIME_STAMP + i * 1000 * 60 * 60L); - } - tree.addAll(Arrays.asList(testedData)); - } - - public static long find(long toElement) { - return tree.higher(toElement); - } - - - public static void main(String[] args) { - initData(); - - long startTime = System.nanoTime(); - - for (long i = 0; i < 100000000L; i++) { - find(1478323448L); - } - - long totalTime = System.nanoTime() - startTime; - System.out.println("total time : " + totalTime + " " + (totalTime * 1.0 / 100000000L)); - } -} diff --git a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/TreeSetTest.java b/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/TreeSetTest.java deleted file mode 100644 index adac9c616..000000000 --- a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/search/TreeSetTest.java +++ /dev/null @@ -1,40 +0,0 @@ -package com.a.eye.skywalking.search; - -import org.junit.Before; -import org.junit.Test; - -import java.util.TreeSet; - -import static org.junit.Assert.assertEquals; - -/** - * Created by xin on 2016/11/1. - */ -public class TreeSetTest { - - private TreeSet treeSet = new TreeSet(); - - @Before - public void setup(){ - treeSet.add(9L); - treeSet.add(3L); - treeSet.add(13L); - treeSet.add(15L); - treeSet.add(1L); - treeSet.add(11L); - treeSet.add(5L); - treeSet.add(7L); - } - - @Test - public void testGetElement(){ - assertEquals(new Long(5), treeSet.higher(4L)); - } - - @Test - public void testRemoveElement(){ - assertEquals(new Long(1), treeSet.first()); - treeSet.pollFirst(); - assertEquals(new Long(3), treeSet.first()); - } -} diff --git a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/block/index/DataIndexFileOperator.java b/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/block/index/DataIndexFileOperator.java deleted file mode 100644 index 66024e856..000000000 --- a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/block/index/DataIndexFileOperator.java +++ /dev/null @@ -1,7 +0,0 @@ -package com.a.eye.skywalking.storage.block.index; - -/** - * Created by xin on 2016/10/31. - */ -public class DataIndexFileOperator { -}