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 dfbd0275e..915308fdf 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 @@ -2,7 +2,8 @@ package com.a.eye.skywalking.storage.config; public class Constants { - public final static String TABLE_NAME = "data_index"; + public final static String TABLE_NAME = "data_index"; + public static final String DRIVER_CLASS_NAME = "org.hsqldb.jdbc.JDBCDriver"; public static int MAX_BATCH_SIZE = 50; 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 11d69edae..9fc064f15 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 @@ -3,6 +3,8 @@ package com.a.eye.skywalking.storage.data; import com.a.eye.datacarrier.consumer.IConsumer; import com.a.eye.skywalking.health.report.HealthCollector; import com.a.eye.skywalking.health.report.HeathReading; +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.data.file.DataFileWriter; import com.a.eye.skywalking.storage.data.index.*; @@ -13,6 +15,7 @@ import java.util.List; public class SpanDataConsumer implements IConsumer { + private static ILog logger = LogManager.getLogger(SpanDataConsumer.class); private IndexDBConnectorCache cache; private DataFileWriter fileWriter; @@ -23,7 +26,6 @@ public class SpanDataConsumer implements IConsumer { @Override public void consume(List data) { - Iterator> iterator = IndexMetaCollections.group(fileWriter.write(data), new GroupKeyBuilder() { @Override @@ -39,7 +41,6 @@ public class SpanDataConsumer implements IConsumer { HealthCollector.getCurrentHeathReading("SpanDataConsumer") .updateData(HeathReading.INFO, "%s messages were successful consumed .", data.size()); } - } private IndexDBConnector getDBConnector(IndexMetaGroup metaGroup) { @@ -47,7 +48,9 @@ public class SpanDataConsumer implements IConsumer { } @Override - public void onError(List list, Throwable throwable) { - + public void onError(List span, Throwable throwable) { + logger.error("Failed to consumer span data.", throwable); + HealthCollector.getCurrentHeathReading("SpanDataConsumer").updateData(HeathReading.ERROR, + "Failed to consume span data. error message : " + throwable.getMessage()); } } 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 aa2b8bdfc..62110905b 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 @@ -4,6 +4,7 @@ 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.config.Constants; 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; @@ -18,8 +19,8 @@ import static com.a.eye.skywalking.storage.config.Constants.SQL.DEFAULT_USER; import static com.a.eye.skywalking.storage.util.PathResolver.getAbsolutePath; public class SpanDataFinder { - private static ILog logger = LogManager.getLogger(SpanDataFinder.class); - private static IndexDataSourceCache datasourceCache = + private static ILog logger = LogManager.getLogger(SpanDataFinder.class); + private static IndexDataSourceCache datasourceCache = new IndexDataSourceCache(Config.Finder.CACHED_SIZE); public static List find(String traceId) { @@ -75,7 +76,7 @@ public class SpanDataFinder { HikariConfig config = new HikariConfig(); config.setJdbcUrl(new ConnectURLGenerator(getAbsolutePath(Config.DataIndex.PATH), Config.DataIndex.FILE_NAME).generate(blockIndex)); - config.setDriverClassName("org.hsqldb.jdbc.JDBCDriver"); + config.setDriverClassName(Constants.DRIVER_CLASS_NAME); config.setUsername(DEFAULT_USER); config.setPassword(DEFAULT_PASSWORD); config.setMaximumPoolSize(Config.Finder.DataSource.MAX_POOL_SIZE); 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 737902015..86779a97c 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 @@ -24,7 +24,7 @@ public class IndexDBConnector { static { try { - Class.forName("org.hsqldb.jdbc.JDBCDriver"); + Class.forName(Constants.DRIVER_CLASS_NAME); } catch (ClassNotFoundException e) { //never } @@ -32,8 +32,8 @@ public class IndexDBConnector { private long timestamp; private Connection connection; - private ConnectURLGenerator generator = new ConnectURLGenerator(getAbsolutePath(Config.DataIndex.PATH), - Config.DataIndex.FILE_NAME); + private ConnectURLGenerator generator = + new ConnectURLGenerator(getAbsolutePath(Config.DataIndex.PATH), Config.DataIndex.FILE_NAME); public IndexDBConnector(long timestamp) { this.timestamp = timestamp; diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/SearchListener.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/SearchListener.java index 2fdce57ae..dab422fdf 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/SearchListener.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/SearchListener.java @@ -1,19 +1,36 @@ package com.a.eye.skywalking.storage.listener; +import com.a.eye.skywalking.health.report.HealthCollector; +import com.a.eye.skywalking.health.report.HeathReading; +import com.a.eye.skywalking.logging.api.ILog; +import com.a.eye.skywalking.logging.api.LogManager; import com.a.eye.skywalking.network.grpc.Span; import com.a.eye.skywalking.network.listener.TraceSearchListener; import com.a.eye.skywalking.storage.data.SpanDataFinder; import com.a.eye.skywalking.storage.data.spandata.SpanData; import com.a.eye.skywalking.storage.data.spandata.SpanDataHelper; +import java.util.ArrayList; import java.util.List; public class SearchListener implements TraceSearchListener { + private static ILog logger = LogManager.getLogger(SearchListener.class); + @Override - public List search(String s) { - List data = SpanDataFinder.find(s); - SpanDataHelper helper = new SpanDataHelper(data); - return helper.category().mergeData(); + public List search(String traceId) { + try { + List data = SpanDataFinder.find(traceId); + SpanDataHelper helper = new SpanDataHelper(data); + List span = helper.category().mergeData(); + HealthCollector.getCurrentHeathReading("SearchListener") + .updateData(HeathReading.INFO, span.size() + " spans was founded by trace Id [" + traceId + "]."); + return span; + } catch (Exception e) { + logger.error("Failed to search trace Id [{}]", traceId, e); + HealthCollector.getCurrentHeathReading("SearchListener") + .updateData(HeathReading.ERROR, "Failed to search trace Id" + traceId + "."); + return new ArrayList(); + } } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java index 7cef89232..90aad32e2 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java @@ -1,6 +1,8 @@ package com.a.eye.skywalking.storage.listener; import com.a.eye.datacarrier.DataCarrier; +import com.a.eye.skywalking.health.report.HealthCollector; +import com.a.eye.skywalking.health.report.HeathReading; import com.a.eye.skywalking.logging.api.ILog; import com.a.eye.skywalking.logging.api.LogManager; import com.a.eye.skywalking.network.grpc.AckSpan; @@ -26,9 +28,13 @@ public class StorageListener implements SpanStorageListener { public boolean storage(RequestSpan requestSpan) { try { spanDataDataCarrier.produce(SpanDataBuilder.build(requestSpan)); + HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.INFO,"Request span " + + "consume successfully"); return true; } catch (Exception e) { logger.error("Failed to storage request span. Span Data:\n {}.", requestSpan.toByteString(), e); + HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.ERROR,"Request span " + + "consume failed"); return false; } }