From 0baf691ae48f9b6971a0503d8b210e50e50dcb1a Mon Sep 17 00:00:00 2001 From: ascrutae Date: Fri, 1 Sep 2017 12:12:21 +0800 Subject: [PATCH 1/4] fix jdbc plugin issue --- .../apm/plugin/jdbc/define/JDBCDriverInterceptor.java | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/apm-sniffer/apm-sdk-plugin/jdbc-plugin/src/main/java/org/skywalking/apm/plugin/jdbc/define/JDBCDriverInterceptor.java b/apm-sniffer/apm-sdk-plugin/jdbc-plugin/src/main/java/org/skywalking/apm/plugin/jdbc/define/JDBCDriverInterceptor.java index 535b04e88..244d24e17 100644 --- a/apm-sniffer/apm-sdk-plugin/jdbc-plugin/src/main/java/org/skywalking/apm/plugin/jdbc/define/JDBCDriverInterceptor.java +++ b/apm-sniffer/apm-sdk-plugin/jdbc-plugin/src/main/java/org/skywalking/apm/plugin/jdbc/define/JDBCDriverInterceptor.java @@ -23,8 +23,12 @@ public class JDBCDriverInterceptor implements InstanceMethodsAroundInterceptor { @Override public Object afterMethod(EnhancedInstance objInst, Method method, Object[] allArguments, Class[] argumentsTypes, Object ret) throws Throwable { - return new SWConnection((String)allArguments[0], - (Properties)allArguments[1], (Connection)ret); + if (ret != null) { + return new SWConnection((String)allArguments[0], + (Properties)allArguments[1], (Connection)ret); + } + + return ret; } @Override public void handleMethodException(EnhancedInstance objInst, Method method, Object[] allArguments, From bdf29e5a92fc40c24c6a36a6d55c1decc77b0a61 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Sat, 2 Sep 2017 18:11:53 +0800 Subject: [PATCH 2/4] Jvm metric test success. #365 --- .../JVMMetricsServiceHandlerTestCase.java | 36 +++- .../cost/define/SegmentCostEsTableDefine.java | 4 +- .../agentstream/mock/SegmentPost.java | 31 ++- .../json/segment/normal/dubbox-provider.json | 4 +- .../define/ElasticSearchStorageInstaller.java | 7 +- .../apm/collector/ui/dao/GCMetricEsDAO.java | 13 +- .../collector/ui/dao/IInstPerformanceDAO.java | 8 +- .../apm/collector/ui/dao/ISegmentCostDAO.java | 6 +- .../ui/dao/InstPerformanceEsDAO.java | 44 +++- .../ui/dao/InstPerformanceH2DAO.java | 12 +- .../collector/ui/dao/MemoryMetricEsDAO.java | 3 +- .../ui/dao/MemoryPoolMetricEsDAO.java | 3 +- .../collector/ui/dao/SegmentCostEsDAO.java | 10 +- .../collector/ui/dao/SegmentCostH2DAO.java | 2 +- .../collector/ui/dao/ServiceEntryEsDAO.java | 2 +- .../ui/jetty/UIJettyModuleDefine.java | 6 + .../jetty/handler/SegmentTopGetHandler.java | 11 +- .../InstanceHealthGetHandler.java | 12 +- .../time/OneInstanceLastTimeGetHandler.java | 12 +- .../ui/service/InstanceHealthService.java | 5 +- .../ui/service/InstanceJVMService.java | 201 +++++++++--------- .../ui/service/SegmentTopService.java | 4 +- .../resources/META-INF/defines/es_dao.define | 3 + 23 files changed, 283 insertions(+), 156 deletions(-) diff --git a/apm-collector/apm-collector-agentjvm/src/test/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandlerTestCase.java b/apm-collector/apm-collector-agentjvm/src/test/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandlerTestCase.java index 6183381fb..8c36e2c83 100644 --- a/apm-collector/apm-collector-agentjvm/src/test/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandlerTestCase.java +++ b/apm-collector/apm-collector-agentjvm/src/test/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandlerTestCase.java @@ -78,14 +78,34 @@ public class JVMMetricsServiceHandlerTestCase { } private static void buildMemoryPoolMetric(JVMMetric.Builder jvmMetric) { - MemoryPool.Builder builder_1 = MemoryPool.newBuilder(); - builder_1.setType(PoolType.NEWGEN_USAGE); - builder_1.setIsHeap(true); - builder_1.setInit(20); - builder_1.setMax(100); - builder_1.setUsed(50); - builder_1.setCommited(30); - jvmMetric.addMemoryPool(builder_1.build()); + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.NEWGEN_USAGE, true).build()); + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.NEWGEN_USAGE, false).build()); + + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.OLDGEN_USAGE, true).build()); + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.OLDGEN_USAGE, false).build()); + + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.METASPACE_USAGE, true).build()); + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.METASPACE_USAGE, false).build()); + + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.PERMGEN_USAGE, true).build()); + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.PERMGEN_USAGE, false).build()); + + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.SURVIVOR_USAGE, true).build()); + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.SURVIVOR_USAGE, false).build()); + + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.CODE_CACHE_USAGE, true).build()); + jvmMetric.addMemoryPool(buildMemoryPoolMetric(PoolType.CODE_CACHE_USAGE, false).build()); + } + + private static MemoryPool.Builder buildMemoryPoolMetric(PoolType poolType, boolean isHeap) { + MemoryPool.Builder builder = MemoryPool.newBuilder(); + builder.setType(poolType); + builder.setIsHeap(isHeap); + builder.setInit(20); + builder.setMax(100); + builder.setUsed(50); + builder.setCommited(30); + return builder; } private static void buildGcMetric(JVMMetric.Builder jvmMetric) { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java index 3e047134c..a25bf0fba 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.segment.cost.define; +import org.skywalking.apm.collector.storage.define.segment.SegmentCostTable; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; -import org.skywalking.apm.collector.storage.define.segment.SegmentCostTable; /** * @author pengys5 @@ -27,7 +27,7 @@ public class SegmentCostEsTableDefine extends ElasticSearchTableDefine { @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_SEGMENT_ID, ElasticSearchColumnDefine.Type.Keyword.name())); - addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_SERVICE_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_SERVICE_NAME, ElasticSearchColumnDefine.Type.Text.name())); addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_COST, ElasticSearchColumnDefine.Type.Long.name())); addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_START_TIME, ElasticSearchColumnDefine.Type.Long.name())); addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_END_TIME, ElasticSearchColumnDefine.Type.Long.name())); diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java index 26dd7ce04..23293258b 100644 --- a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java @@ -28,9 +28,9 @@ public class SegmentPost { InstanceEsDAO instanceEsDAO = new InstanceEsDAO(); instanceEsDAO.setClient(client); - InstanceDataDefine.Instance consumerInstance = new InstanceDataDefine.Instance("2", 2, "dubbox-consumer", now, 2, now, ""); + InstanceDataDefine.Instance consumerInstance = new InstanceDataDefine.Instance("2", 2, "dubbox-consumer", now, 2, now, osInfo("consumer").toString()); instanceEsDAO.save(consumerInstance); - InstanceDataDefine.Instance providerInstance = new InstanceDataDefine.Instance("3", 3, "dubbox-provider", now, 3, now, ""); + InstanceDataDefine.Instance providerInstance = new InstanceDataDefine.Instance("3", 3, "dubbox-provider", now, 3, now, osInfo("provider").toString()); instanceEsDAO.save(providerInstance); ApplicationEsDAO applicationEsDAO = new ApplicationEsDAO(); @@ -64,10 +64,13 @@ public class SegmentPost { modifyTime(provider); HttpClientTools.INSTANCE.post("http://localhost:12800/segments", provider.toString()); + diff = 0; Thread.sleep(1000); } } + private static long diff = 0; + private static void modifyTime(JsonElement jsonElement) { JsonArray segmentArray = jsonElement.getAsJsonArray(); for (JsonElement element : segmentArray) { @@ -76,10 +79,28 @@ public class SegmentPost { for (JsonElement span : spans) { long startTime = span.getAsJsonObject().get("st").getAsLong(); long endTime = span.getAsJsonObject().get("et").getAsLong(); - long currentTime = System.currentTimeMillis(); - span.getAsJsonObject().addProperty("st", currentTime); - span.getAsJsonObject().addProperty("et", currentTime + (endTime - startTime)); + + if (diff == 0) { + diff = System.currentTimeMillis() - startTime; + } + + span.getAsJsonObject().addProperty("st", startTime + diff); + span.getAsJsonObject().addProperty("et", endTime + diff); } } } + + private static JsonObject osInfo(String hostName) { + JsonObject osInfoJson = new JsonObject(); + osInfoJson.addProperty("osName", "Linux"); + osInfoJson.addProperty("hostName", hostName); + osInfoJson.addProperty("processId", 1); + + JsonArray ipv4Array = new JsonArray(); + ipv4Array.add("123.123.123.123"); + ipv4Array.add("124.124.124.124"); + osInfoJson.add("ipv4s", ipv4Array); + + return osInfoJson; + } } diff --git a/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-provider.json b/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-provider.json index 6da99d6f7..28a61aee0 100644 --- a/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-provider.json +++ b/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-provider.json @@ -40,8 +40,8 @@ "tv": 0, "lv": 2, "ps": -1, - "st": 1501858094883, - "et": 1501858096950, + "st": 1501858094726, + "et": 1501858096804, "ci": 3, "cn": "", "oi": 0, diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java index 040349742..05e3f2738 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java @@ -53,7 +53,12 @@ public class ElasticSearchStorageInstaller extends StorageInstaller { return Settings.builder() .put("index.number_of_shards", tableDefine.numberOfShards()) .put("index.number_of_replicas", tableDefine.numberOfReplicas()) - .put("index.refresh_interval", String.valueOf(tableDefine.refreshInterval()) + "s").build(); + .put("index.refresh_interval", String.valueOf(tableDefine.refreshInterval()) + "s") + + .put("analysis.analyzer.collector_analyzer.tokenizer", "collector_tokenizer") + .put("analysis.tokenizer.collector_tokenizer.type", "standard") + .put("analysis.tokenizer.collector_tokenizer.max_token_length", 5) + .build(); } private XContentBuilder createMappingBuilder(ElasticSearchTableDefine tableDefine) throws IOException { diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricEsDAO.java index 0e6ce7a45..16b441e5e 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricEsDAO.java @@ -16,7 +16,6 @@ import org.elasticsearch.search.aggregations.AggregationBuilders; import org.elasticsearch.search.aggregations.bucket.terms.Terms; import org.elasticsearch.search.aggregations.metrics.sum.Sum; import org.skywalking.apm.collector.core.util.Const; -import org.skywalking.apm.collector.storage.define.jvm.CpuMetricTable; import org.skywalking.apm.collector.storage.define.jvm.GCMetricTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.skywalking.apm.network.proto.GCPhrase; @@ -92,8 +91,8 @@ public class GCMetricEsDAO extends EsDAO implements IGCMetricDAO { MultiGetRequestBuilder youngPrepareMultiGet = getClient().prepareMultiGet(); int i = 0; do { - String youngId = (startTimeBucket + i) + Const.ID_SPLIT + GCPhrase.NEW_VALUE + instanceId; - youngPrepareMultiGet.add(CpuMetricTable.TABLE, CpuMetricTable.TABLE_TYPE, youngId); + String youngId = (startTimeBucket + i) + Const.ID_SPLIT + instanceId + Const.ID_SPLIT + GCPhrase.NEW_VALUE; + youngPrepareMultiGet.add(GCMetricTable.TABLE, GCMetricTable.TABLE_TYPE, youngId); i++; } while (startTimeBucket + i <= endTimeBucket); @@ -102,7 +101,7 @@ public class GCMetricEsDAO extends EsDAO implements IGCMetricDAO { MultiGetResponse multiGetResponse = youngPrepareMultiGet.get(); for (MultiGetItemResponse itemResponse : multiGetResponse.getResponses()) { if (itemResponse.getResponse().isExists()) { - youngArray.add(((Number)itemResponse.getResponse().getSource().get(CpuMetricTable.COLUMN_USAGE_PERCENT)).intValue()); + youngArray.add(((Number)itemResponse.getResponse().getSource().get(GCMetricTable.COLUMN_COUNT)).intValue()); } else { youngArray.add(0); } @@ -112,8 +111,8 @@ public class GCMetricEsDAO extends EsDAO implements IGCMetricDAO { MultiGetRequestBuilder oldPrepareMultiGet = getClient().prepareMultiGet(); i = 0; do { - String oldId = (startTimeBucket + i) + Const.ID_SPLIT + GCPhrase.OLD_VALUE + instanceId; - oldPrepareMultiGet.add(CpuMetricTable.TABLE, CpuMetricTable.TABLE_TYPE, oldId); + String oldId = (startTimeBucket + i) + Const.ID_SPLIT + instanceId + Const.ID_SPLIT + GCPhrase.OLD_VALUE; + oldPrepareMultiGet.add(GCMetricTable.TABLE, GCMetricTable.TABLE_TYPE, oldId); i++; } while (startTimeBucket + i <= endTimeBucket); @@ -123,7 +122,7 @@ public class GCMetricEsDAO extends EsDAO implements IGCMetricDAO { multiGetResponse = oldPrepareMultiGet.get(); for (MultiGetItemResponse itemResponse : multiGetResponse.getResponses()) { if (itemResponse.getResponse().isExists()) { - oldArray.add(((Number)itemResponse.getResponse().getSource().get(CpuMetricTable.COLUMN_USAGE_PERCENT)).intValue()); + oldArray.add(((Number)itemResponse.getResponse().getSource().get(GCMetricTable.COLUMN_COUNT)).intValue()); } else { oldArray.add(0); } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IInstPerformanceDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IInstPerformanceDAO.java index 488df87b0..850872e38 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IInstPerformanceDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IInstPerformanceDAO.java @@ -9,9 +9,13 @@ import java.util.List; public interface IInstPerformanceDAO { List getMultiple(long timeBucket, int applicationId); - int getMetric(int instanceId, long timeBucket); + int getTpsMetric(int instanceId, long timeBucket); - JsonArray getMetric(int instanceId, long startTimeBucket, long endTimeBucket); + JsonArray getTpsMetric(int instanceId, long startTimeBucket, long endTimeBucket); + + int getRespTimeMetric(int instanceId, long timeBucket); + + JsonArray getRespTimeMetric(int instanceId, long startTimeBucket, long endTimeBucket); class InstPerformance { private final int instanceId; diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ISegmentCostDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ISegmentCostDAO.java index 43466b991..99c09d34e 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ISegmentCostDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ISegmentCostDAO.java @@ -7,5 +7,9 @@ import com.google.gson.JsonObject; */ public interface ISegmentCostDAO { JsonObject loadTop(long startTime, long endTime, long minCost, long maxCost, String operationName, - String globalTraceId, int limit, int from); + String globalTraceId, int limit, int from, Sort sort); + + public enum Sort { + Cost, Time + } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceEsDAO.java index f307512a3..e7dcd7dbd 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceEsDAO.java @@ -18,7 +18,6 @@ import org.elasticsearch.search.aggregations.bucket.terms.Terms; import org.elasticsearch.search.aggregations.metrics.sum.Sum; import org.skywalking.apm.collector.core.util.Const; import org.skywalking.apm.collector.storage.define.instance.InstPerformanceTable; -import org.skywalking.apm.collector.storage.define.jvm.CpuMetricTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; /** @@ -68,7 +67,7 @@ public class InstPerformanceEsDAO extends EsDAO implements IInstPerformanceDAO { return instPerformances; } - @Override public int getMetric(int instanceId, long timeBucket) { + @Override public int getTpsMetric(int instanceId, long timeBucket) { String id = timeBucket + Const.ID_SPLIT + instanceId; GetResponse getResponse = getClient().prepareGet(InstPerformanceTable.TABLE, id).get(); @@ -78,13 +77,13 @@ public class InstPerformanceEsDAO extends EsDAO implements IInstPerformanceDAO { return 0; } - @Override public JsonArray getMetric(int instanceId, long startTimeBucket, long endTimeBucket) { + @Override public JsonArray getTpsMetric(int instanceId, long startTimeBucket, long endTimeBucket) { MultiGetRequestBuilder prepareMultiGet = getClient().prepareMultiGet(); int i = 0; do { String id = (startTimeBucket + i) + Const.ID_SPLIT + instanceId; - prepareMultiGet.add(CpuMetricTable.TABLE, InstPerformanceTable.TABLE_TYPE, id); + prepareMultiGet.add(InstPerformanceTable.TABLE, InstPerformanceTable.TABLE_TYPE, id); i++; } while (startTimeBucket + i <= endTimeBucket); @@ -100,4 +99,41 @@ public class InstPerformanceEsDAO extends EsDAO implements IInstPerformanceDAO { } return metrics; } + + @Override public int getRespTimeMetric(int instanceId, long timeBucket) { + String id = timeBucket + Const.ID_SPLIT + instanceId; + GetResponse getResponse = getClient().prepareGet(InstPerformanceTable.TABLE, id).get(); + + if (getResponse.isExists()) { + int callTimes = ((Number)getResponse.getSource().get(InstPerformanceTable.COLUMN_CALL_TIMES)).intValue(); + int costTotal = ((Number)getResponse.getSource().get(InstPerformanceTable.COLUMN_COST_TOTAL)).intValue(); + return costTotal / callTimes; + } + return 0; + } + + @Override public JsonArray getRespTimeMetric(int instanceId, long startTimeBucket, long endTimeBucket) { + MultiGetRequestBuilder prepareMultiGet = getClient().prepareMultiGet(); + + int i = 0; + do { + String id = (startTimeBucket + i) + Const.ID_SPLIT + instanceId; + prepareMultiGet.add(InstPerformanceTable.TABLE, InstPerformanceTable.TABLE_TYPE, id); + i++; + } + while (startTimeBucket + i <= endTimeBucket); + + JsonArray metrics = new JsonArray(); + MultiGetResponse multiGetResponse = prepareMultiGet.get(); + for (MultiGetItemResponse response : multiGetResponse.getResponses()) { + if (response.getResponse().isExists()) { + int callTimes = ((Number)response.getResponse().getSource().get(InstPerformanceTable.COLUMN_CALL_TIMES)).intValue(); + int costTotal = ((Number)response.getResponse().getSource().get(InstPerformanceTable.COLUMN_COST_TOTAL)).intValue(); + metrics.add(costTotal / callTimes); + } else { + metrics.add(0); + } + } + return metrics; + } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceH2DAO.java index a480c7ce3..fcb20253b 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceH2DAO.java @@ -13,11 +13,19 @@ public class InstPerformanceH2DAO extends H2DAO implements IInstPerformanceDAO { return null; } - @Override public int getMetric(int instanceId, long timeBucket) { + @Override public int getTpsMetric(int instanceId, long timeBucket) { return 0; } - @Override public JsonArray getMetric(int instanceId, long startTimeBucket, long endTimeBucket) { + @Override public JsonArray getTpsMetric(int instanceId, long startTimeBucket, long endTimeBucket) { + return null; + } + + @Override public int getRespTimeMetric(int instanceId, long timeBucket) { + return 0; + } + + @Override public JsonArray getRespTimeMetric(int instanceId, long startTimeBucket, long endTimeBucket) { return null; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryMetricEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryMetricEsDAO.java index 3e481a587..e441b0e1c 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryMetricEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryMetricEsDAO.java @@ -44,9 +44,7 @@ public class MemoryMetricEsDAO extends EsDAO implements IMemoryMetricDAO { while (startTimeBucket + i <= endTimeBucket); JsonObject metric = new JsonObject(); - JsonArray usedMetric = new JsonArray(); - MultiGetResponse multiGetResponse = prepareMultiGet.get(); for (MultiGetItemResponse response : multiGetResponse.getResponses()) { if (response.getResponse().isExists()) { @@ -57,6 +55,7 @@ public class MemoryMetricEsDAO extends EsDAO implements IMemoryMetricDAO { usedMetric.add(0); } } + metric.add("used", usedMetric); return metric; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryPoolMetricEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryPoolMetricEsDAO.java index a2cbbc10d..01b78c3ee 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryPoolMetricEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryPoolMetricEsDAO.java @@ -45,9 +45,7 @@ public class MemoryPoolMetricEsDAO extends EsDAO implements IMemoryPoolMetricDAO while (startTimeBucket + i <= endTimeBucket); JsonObject metric = new JsonObject(); - JsonArray usedMetric = new JsonArray(); - MultiGetResponse multiGetResponse = prepareMultiGet.get(); for (MultiGetItemResponse response : multiGetResponse.getResponses()) { if (response.getResponse().isExists()) { @@ -58,6 +56,7 @@ public class MemoryPoolMetricEsDAO extends EsDAO implements IMemoryPoolMetricDAO usedMetric.add(0); } } + metric.add("used", usedMetric); return metric; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostEsDAO.java index 32e32bc75..f724a7451 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostEsDAO.java @@ -15,9 +15,9 @@ import org.elasticsearch.search.sort.SortOrder; import org.skywalking.apm.collector.core.util.CollectionUtils; import org.skywalking.apm.collector.core.util.StringUtils; import org.skywalking.apm.collector.storage.dao.DAOContainer; -import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.skywalking.apm.collector.storage.define.global.GlobalTraceTable; import org.skywalking.apm.collector.storage.define.segment.SegmentCostTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; /** * @author pengys5 @@ -25,7 +25,7 @@ import org.skywalking.apm.collector.storage.define.segment.SegmentCostTable; public class SegmentCostEsDAO extends EsDAO implements ISegmentCostDAO { @Override public JsonObject loadTop(long startTime, long endTime, long minCost, long maxCost, String operationName, - String globalTraceId, int limit, int from) { + String globalTraceId, int limit, int from, Sort sort) { SearchRequestBuilder searchRequestBuilder = getClient().prepareSearch(SegmentCostTable.TABLE); searchRequestBuilder.setTypes(SegmentCostTable.TABLE_TYPE); searchRequestBuilder.setSearchType(SearchType.DFS_QUERY_THEN_FETCH); @@ -48,7 +48,11 @@ public class SegmentCostEsDAO extends EsDAO implements ISegmentCostDAO { mustQueryList.add(QueryBuilders.matchQuery(SegmentCostTable.COLUMN_SERVICE_NAME, operationName)); } - searchRequestBuilder.addSort(SegmentCostTable.COLUMN_COST, SortOrder.DESC); + if (Sort.Cost.equals(sort)) { + searchRequestBuilder.addSort(SegmentCostTable.COLUMN_COST, SortOrder.DESC); + } else if (Sort.Time.equals(sort)) { + searchRequestBuilder.addSort(SegmentCostTable.COLUMN_START_TIME, SortOrder.DESC); + } searchRequestBuilder.setSize(limit); searchRequestBuilder.setFrom(from); diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostH2DAO.java index fcafd0f60..644eab487 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostH2DAO.java @@ -8,7 +8,7 @@ import org.skywalking.apm.collector.storage.h2.dao.H2DAO; */ public class SegmentCostH2DAO extends H2DAO implements ISegmentCostDAO { @Override public JsonObject loadTop(long startTime, long endTime, long minCost, long maxCost, String operationName, - String globalTraceId, int limit, int from) { + String globalTraceId, int limit, int from, Sort sort) { return null; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ServiceEntryEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ServiceEntryEsDAO.java index 299c4ae46..8d0145a59 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ServiceEntryEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ServiceEntryEsDAO.java @@ -37,7 +37,7 @@ public class ServiceEntryEsDAO extends EsDAO implements IServiceEntryDAO { boolQueryBuilder.must().add(QueryBuilders.matchQuery(ServiceEntryTable.COLUMN_APPLICATION_ID, applicationId)); } if (StringUtils.isNotEmpty(entryServiceName)) { - boolQueryBuilder.must().add(QueryBuilders.termQuery(ServiceEntryTable.COLUMN_ENTRY_SERVICE_NAME, entryServiceName)); + boolQueryBuilder.must().add(QueryBuilders.matchQuery(ServiceEntryTable.COLUMN_ENTRY_SERVICE_NAME, entryServiceName)); } searchRequestBuilder.setQuery(boolQueryBuilder); diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyModuleDefine.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyModuleDefine.java index b78deae2e..1b367ec8a 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyModuleDefine.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyModuleDefine.java @@ -17,6 +17,9 @@ import org.skywalking.apm.collector.ui.jetty.handler.TraceStackGetHandler; import org.skywalking.apm.collector.ui.jetty.handler.UIJettyServerHandler; import org.skywalking.apm.collector.ui.jetty.handler.application.ApplicationsGetHandler; import org.skywalking.apm.collector.ui.jetty.handler.instancehealth.InstanceHealthGetHandler; +import org.skywalking.apm.collector.ui.jetty.handler.instancemetric.InstanceMetricGetOneTimeBucketHandler; +import org.skywalking.apm.collector.ui.jetty.handler.instancemetric.InstanceMetricGetRangeTimeBucketHandler; +import org.skywalking.apm.collector.ui.jetty.handler.instancemetric.InstanceOsInfoGetHandler; import org.skywalking.apm.collector.ui.jetty.handler.servicetree.EntryServiceGetHandler; import org.skywalking.apm.collector.ui.jetty.handler.servicetree.ServiceTreeGetHandler; import org.skywalking.apm.collector.ui.jetty.handler.time.AllInstanceLastTimeGetHandler; @@ -64,6 +67,9 @@ public class UIJettyModuleDefine extends UIModuleDefine { handlers.add(new AllInstanceLastTimeGetHandler()); handlers.add(new InstanceHealthGetHandler()); handlers.add(new ApplicationsGetHandler()); + handlers.add(new InstanceOsInfoGetHandler()); + handlers.add(new InstanceMetricGetOneTimeBucketHandler()); + handlers.add(new InstanceMetricGetRangeTimeBucketHandler()); handlers.add(new EntryServiceGetHandler()); handlers.add(new ServiceTreeGetHandler()); return handlers; diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SegmentTopGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SegmentTopGetHandler.java index 003d37e55..5de03a313 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SegmentTopGetHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SegmentTopGetHandler.java @@ -4,6 +4,7 @@ import com.google.gson.JsonElement; import javax.servlet.http.HttpServletRequest; import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; import org.skywalking.apm.collector.server.jetty.JettyHandler; +import org.skywalking.apm.collector.ui.dao.ISegmentCostDAO; import org.skywalking.apm.collector.ui.service.SegmentTopService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -77,7 +78,15 @@ public class SegmentTopGetHandler extends JettyHandler { operationName = req.getParameter("operationName"); } - return service.loadTop(startTime, endTime, minCost, maxCost, operationName, globalTraceId, limit, from); + ISegmentCostDAO.Sort sort = ISegmentCostDAO.Sort.Cost; + if (req.getParameterMap().containsKey("sort")) { + String sortStr = req.getParameter("sort"); + if (sortStr.toLowerCase().equals(ISegmentCostDAO.Sort.Time.name().toLowerCase())) { + sort = ISegmentCostDAO.Sort.Time; + } + } + + return service.loadTop(startTime, endTime, minCost, maxCost, operationName, globalTraceId, limit, from, sort); } @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancehealth/InstanceHealthGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancehealth/InstanceHealthGetHandler.java index 57832ba20..92ac0b206 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancehealth/InstanceHealthGetHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancehealth/InstanceHealthGetHandler.java @@ -24,13 +24,13 @@ public class InstanceHealthGetHandler extends JettyHandler { private InstanceHealthService service = new InstanceHealthService(); @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { - String timestampStr = req.getParameter("timestamp"); + String timeBucketStr = req.getParameter("timeBucket"); String[] applicationIdsStr = req.getParameterValues("applicationIds"); - logger.debug("instance health get timestamp: {}, applicationIdsStr: {}", timestampStr, applicationIdsStr); + logger.debug("instance health get timeBucket: {}, applicationIdsStr: {}", timeBucketStr, applicationIdsStr); - long timestamp; + long timeBucket; try { - timestamp = Long.parseLong(timestampStr); + timeBucket = Long.parseLong(timeBucketStr); } catch (NumberFormatException e) { throw new ArgumentsParseException("timestamp must be long"); } @@ -45,12 +45,12 @@ public class InstanceHealthGetHandler extends JettyHandler { } JsonObject response = new JsonObject(); - response.addProperty("timestamp", timestamp); + response.addProperty("timeBucket", timeBucket); JsonArray appInstances = new JsonArray(); response.add("appInstances", appInstances); for (int applicationId : applicationIds) { - appInstances.add(service.getInstances(timestamp, applicationId)); + appInstances.add(service.getInstances(timeBucket, applicationId)); } return response; } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/time/OneInstanceLastTimeGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/time/OneInstanceLastTimeGetHandler.java index 3e29f4c08..e9891b9dd 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/time/OneInstanceLastTimeGetHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/time/OneInstanceLastTimeGetHandler.java @@ -23,18 +23,18 @@ public class OneInstanceLastTimeGetHandler extends JettyHandler { private TimeSynchronousService service = new TimeSynchronousService(); @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { - String applicationInstanceIdStr = req.getParameter("applicationInstanceId"); - logger.debug("applicationInstanceId: {}", applicationInstanceIdStr); + String instanceIdStr = req.getParameter("instanceId"); + logger.debug("instanceId: {}", instanceIdStr); - int applicationInstanceId; + int instanceId; try { - applicationInstanceId = Integer.parseInt(applicationInstanceIdStr); + instanceId = Integer.parseInt(instanceIdStr); } catch (NumberFormatException e) { throw new ArgumentsParseException("application instance id must be integer"); } - Long time = service.instanceLastTime(applicationInstanceId); - logger.debug("application instance id: {}, instance last time: {}", applicationInstanceId, time); + Long time = service.instanceLastTime(instanceId); + logger.debug("application instance id: {}, instance last time: {}", instanceId, time); JsonObject timeJson = new JsonObject(); timeJson.addProperty("timeBucket", time); return timeJson; diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java index 2b380df53..fff48a03f 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java @@ -18,11 +18,10 @@ public class InstanceHealthService { private final Logger logger = LoggerFactory.getLogger(InstanceHealthService.class); - public JsonObject getInstances(long timestamp, int applicationId) { + public JsonObject getInstances(long timeBucket, int applicationId) { JsonObject response = new JsonObject(); - long secondTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(timestamp); - long s5TimeBucket = TimeBucketUtils.INSTANCE.getFiveSecondTimeBucket(secondTimeBucket); + long s5TimeBucket = TimeBucketUtils.INSTANCE.getFiveSecondTimeBucket(timeBucket); IInstPerformanceDAO instPerformanceDAO = (IInstPerformanceDAO)DAOContainer.INSTANCE.get(IInstPerformanceDAO.class.getName()); List performances = instPerformanceDAO.getMultiple(s5TimeBucket, applicationId); diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceJVMService.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceJVMService.java index d84aee2dc..8b8fb00a5 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceJVMService.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceJVMService.java @@ -39,53 +39,58 @@ public class InstanceJVMService { public JsonObject getInstanceJvmMetric(int instanceId, Set metricTypes, long timeBucket) { JsonObject metrics = new JsonObject(); - if (metricTypes.contains(MetricType.cpu.name())) { - ICpuMetricDAO cpuMetricDAO = (ICpuMetricDAO)DAOContainer.INSTANCE.get(ICpuMetricDAO.class.getName()); - metrics.addProperty(MetricType.cpu.name(), cpuMetricDAO.getMetric(instanceId, timeBucket)); - } else if (metricTypes.contains(MetricType.gc.name())) { - IGCMetricDAO gcMetricDAO = (IGCMetricDAO)DAOContainer.INSTANCE.get(IGCMetricDAO.class.getName()); - metrics.add(MetricType.gc.name(), gcMetricDAO.getMetric(instanceId, timeBucket)); - } else if (metricTypes.contains(MetricType.tps.name())) { - IInstPerformanceDAO instPerformanceDAO = (IInstPerformanceDAO)DAOContainer.INSTANCE.get(IInstPerformanceDAO.class.getName()); - metrics.addProperty(MetricType.tps.name(), instPerformanceDAO.getMetric(instanceId, timeBucket)); - } else if (metricTypes.contains(MetricType.heapmemory.name())) { - IMemoryMetricDAO memoryMetricDAO = (IMemoryMetricDAO)DAOContainer.INSTANCE.get(IMemoryMetricDAO.class.getName()); - metrics.add(MetricType.heapmemory.name(), memoryMetricDAO.getMetric(instanceId, timeBucket, true)); - } else if (metricTypes.contains(MetricType.nonheapmemory.name())) { - IMemoryMetricDAO memoryMetricDAO = (IMemoryMetricDAO)DAOContainer.INSTANCE.get(IMemoryMetricDAO.class.getName()); - metrics.add(MetricType.heapmemory.name(), memoryMetricDAO.getMetric(instanceId, timeBucket, false)); - } else if (metricTypes.contains(MetricType.heappermgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.heappermgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, true, PoolType.PERMGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.heapmetaspace.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.heapmetaspace.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, true, PoolType.METASPACE_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.heapnewgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.heapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, true, PoolType.NEWGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.heapoldgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.heapoldgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, true, PoolType.OLDGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.heapsurvivor.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.heapsurvivor.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, true, PoolType.SURVIVOR_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.nonheappermgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.nonheappermgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, false, PoolType.PERMGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.nonheapmetaspace.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.nonheapmetaspace.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, false, PoolType.METASPACE_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.nonheapnewgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.nonheapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, false, PoolType.NEWGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.nonheapoldgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.nonheapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, false, PoolType.OLDGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.nonheapsurvivor.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.nonheapsurvivor.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, false, PoolType.OLDGEN_USAGE_VALUE)); - } else { - throw new UnexpectedException("unexpected metric type"); + for (String metricType : metricTypes) { + if (metricType.toLowerCase().equals(MetricType.cpu.name())) { + ICpuMetricDAO cpuMetricDAO = (ICpuMetricDAO)DAOContainer.INSTANCE.get(ICpuMetricDAO.class.getName()); + metrics.addProperty(MetricType.cpu.name(), cpuMetricDAO.getMetric(instanceId, timeBucket)); + } else if (metricType.toLowerCase().equals(MetricType.gc.name())) { + IGCMetricDAO gcMetricDAO = (IGCMetricDAO)DAOContainer.INSTANCE.get(IGCMetricDAO.class.getName()); + metrics.add(MetricType.gc.name(), gcMetricDAO.getMetric(instanceId, timeBucket)); + } else if (metricType.toLowerCase().equals(MetricType.tps.name())) { + IInstPerformanceDAO instPerformanceDAO = (IInstPerformanceDAO)DAOContainer.INSTANCE.get(IInstPerformanceDAO.class.getName()); + metrics.addProperty(MetricType.tps.name(), instPerformanceDAO.getTpsMetric(instanceId, timeBucket)); + } else if (metricType.toLowerCase().equals(MetricType.resptime.name())) { + IInstPerformanceDAO instPerformanceDAO = (IInstPerformanceDAO)DAOContainer.INSTANCE.get(IInstPerformanceDAO.class.getName()); + metrics.addProperty(MetricType.resptime.name(), instPerformanceDAO.getRespTimeMetric(instanceId, timeBucket)); + } else if (metricType.toLowerCase().equals(MetricType.heapmemory.name())) { + IMemoryMetricDAO memoryMetricDAO = (IMemoryMetricDAO)DAOContainer.INSTANCE.get(IMemoryMetricDAO.class.getName()); + metrics.add(MetricType.heapmemory.name(), memoryMetricDAO.getMetric(instanceId, timeBucket, true)); + } else if (metricType.toLowerCase().equals(MetricType.nonheapmemory.name())) { + IMemoryMetricDAO memoryMetricDAO = (IMemoryMetricDAO)DAOContainer.INSTANCE.get(IMemoryMetricDAO.class.getName()); + metrics.add(MetricType.nonheapmemory.name(), memoryMetricDAO.getMetric(instanceId, timeBucket, false)); + } else if (metricType.toLowerCase().equals(MetricType.heappermgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.heappermgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, true, PoolType.PERMGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.heapmetaspace.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.heapmetaspace.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, true, PoolType.METASPACE_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.heapnewgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.heapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, true, PoolType.NEWGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.heapoldgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.heapoldgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, true, PoolType.OLDGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.heapsurvivor.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.heapsurvivor.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, true, PoolType.SURVIVOR_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.nonheappermgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.nonheappermgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, false, PoolType.PERMGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.nonheapmetaspace.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.nonheapmetaspace.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, false, PoolType.METASPACE_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.nonheapnewgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.nonheapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, false, PoolType.NEWGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.nonheapoldgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.nonheapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, false, PoolType.OLDGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.nonheapsurvivor.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.nonheapsurvivor.name(), memoryPoolMetricDAO.getMetric(instanceId, timeBucket, false, PoolType.OLDGEN_USAGE_VALUE)); + } else { + throw new UnexpectedException("unexpected metric type"); + } } return metrics; } @@ -93,59 +98,65 @@ public class InstanceJVMService { public JsonObject getInstanceJvmMetrics(int instanceId, Set metricTypes, long startTimeBucket, long endTimeBucket) { JsonObject metrics = new JsonObject(); - if (metricTypes.contains(MetricType.cpu.name())) { - ICpuMetricDAO cpuMetricDAO = (ICpuMetricDAO)DAOContainer.INSTANCE.get(ICpuMetricDAO.class.getName()); - metrics.add(MetricType.cpu.name(), cpuMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket)); - } else if (metricTypes.contains(MetricType.gc.name())) { - IGCMetricDAO gcMetricDAO = (IGCMetricDAO)DAOContainer.INSTANCE.get(IGCMetricDAO.class.getName()); - metrics.add(MetricType.gc.name(), gcMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket)); - } else if (metricTypes.contains(MetricType.tps.name())) { - IInstPerformanceDAO instPerformanceDAO = (IInstPerformanceDAO)DAOContainer.INSTANCE.get(IInstPerformanceDAO.class.getName()); - metrics.add(MetricType.tps.name(), instPerformanceDAO.getMetric(instanceId, startTimeBucket, endTimeBucket)); - } else if (metricTypes.contains(MetricType.heapmemory.name())) { - IMemoryMetricDAO memoryMetricDAO = (IMemoryMetricDAO)DAOContainer.INSTANCE.get(IMemoryMetricDAO.class.getName()); - metrics.add(MetricType.heapmemory.name(), memoryMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true)); - } else if (metricTypes.contains(MetricType.nonheapmemory.name())) { - IMemoryMetricDAO memoryMetricDAO = (IMemoryMetricDAO)DAOContainer.INSTANCE.get(IMemoryMetricDAO.class.getName()); - metrics.add(MetricType.heapmemory.name(), memoryMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false)); - } else if (metricTypes.contains(MetricType.heappermgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.heappermgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true, PoolType.PERMGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.heapmetaspace.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.heapmetaspace.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true, PoolType.METASPACE_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.heapnewgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.heapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true, PoolType.NEWGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.heapoldgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.heapoldgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true, PoolType.OLDGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.heapsurvivor.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.heapsurvivor.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true, PoolType.SURVIVOR_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.nonheappermgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.nonheappermgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false, PoolType.PERMGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.nonheapmetaspace.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.nonheapmetaspace.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false, PoolType.METASPACE_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.nonheapnewgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.nonheapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false, PoolType.NEWGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.nonheapoldgen.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.nonheapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false, PoolType.OLDGEN_USAGE_VALUE)); - } else if (metricTypes.contains(MetricType.nonheapsurvivor.name())) { - IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); - metrics.add(MetricType.nonheapsurvivor.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false, PoolType.OLDGEN_USAGE_VALUE)); - } else { - throw new UnexpectedException("unexpected metric type"); + for (String metricType : metricTypes) { + if (metricType.toLowerCase().equals(MetricType.cpu.name())) { + ICpuMetricDAO cpuMetricDAO = (ICpuMetricDAO)DAOContainer.INSTANCE.get(ICpuMetricDAO.class.getName()); + metrics.add(MetricType.cpu.name(), cpuMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket)); + } else if (metricType.toLowerCase().equals(MetricType.gc.name())) { + IGCMetricDAO gcMetricDAO = (IGCMetricDAO)DAOContainer.INSTANCE.get(IGCMetricDAO.class.getName()); + metrics.add(MetricType.gc.name(), gcMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket)); + } else if (metricType.toLowerCase().equals(MetricType.tps.name())) { + IInstPerformanceDAO instPerformanceDAO = (IInstPerformanceDAO)DAOContainer.INSTANCE.get(IInstPerformanceDAO.class.getName()); + metrics.add(MetricType.tps.name(), instPerformanceDAO.getTpsMetric(instanceId, startTimeBucket, endTimeBucket)); + } else if (metricType.toLowerCase().equals(MetricType.resptime.name())) { + IInstPerformanceDAO instPerformanceDAO = (IInstPerformanceDAO)DAOContainer.INSTANCE.get(IInstPerformanceDAO.class.getName()); + metrics.add(MetricType.resptime.name(), instPerformanceDAO.getRespTimeMetric(instanceId, startTimeBucket, endTimeBucket)); + } else if (metricType.toLowerCase().equals(MetricType.heapmemory.name())) { + IMemoryMetricDAO memoryMetricDAO = (IMemoryMetricDAO)DAOContainer.INSTANCE.get(IMemoryMetricDAO.class.getName()); + metrics.add(MetricType.heapmemory.name(), memoryMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true)); + } else if (metricType.toLowerCase().equals(MetricType.nonheapmemory.name())) { + IMemoryMetricDAO memoryMetricDAO = (IMemoryMetricDAO)DAOContainer.INSTANCE.get(IMemoryMetricDAO.class.getName()); + metrics.add(MetricType.nonheapmemory.name(), memoryMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false)); + } else if (metricType.toLowerCase().equals(MetricType.heappermgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.heappermgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true, PoolType.PERMGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.heapmetaspace.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.heapmetaspace.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true, PoolType.METASPACE_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.heapnewgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.heapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true, PoolType.NEWGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.heapoldgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.heapoldgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true, PoolType.OLDGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.heapsurvivor.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.heapsurvivor.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, true, PoolType.SURVIVOR_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.nonheappermgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.nonheappermgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false, PoolType.PERMGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.nonheapmetaspace.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.nonheapmetaspace.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false, PoolType.METASPACE_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.nonheapnewgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.nonheapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false, PoolType.NEWGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.nonheapoldgen.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.nonheapnewgen.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false, PoolType.OLDGEN_USAGE_VALUE)); + } else if (metricType.toLowerCase().equals(MetricType.nonheapsurvivor.name())) { + IMemoryPoolMetricDAO memoryPoolMetricDAO = (IMemoryPoolMetricDAO)DAOContainer.INSTANCE.get(IMemoryPoolMetricDAO.class.getName()); + metrics.add(MetricType.nonheapsurvivor.name(), memoryPoolMetricDAO.getMetric(instanceId, startTimeBucket, endTimeBucket, false, PoolType.OLDGEN_USAGE_VALUE)); + } else { + throw new UnexpectedException("unexpected metric type"); + } } + return metrics; } public enum MetricType { - cpu, gc, tps, heapmemory, heappermgen, heapmetaspace, heapnewgen, + cpu, gc, tps, resptime, heapmemory, heappermgen, heapmetaspace, heapnewgen, heapoldgen, heapsurvivor, nonheapmemory, nonheappermgen, nonheapmetaspace, nonheapnewgen, nonheapoldgen, nonheapsurvivor } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/SegmentTopService.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/SegmentTopService.java index 97958b48b..48ed414af 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/SegmentTopService.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/SegmentTopService.java @@ -14,9 +14,9 @@ public class SegmentTopService { private final Logger logger = LoggerFactory.getLogger(SegmentTopService.class); public JsonObject loadTop(long startTime, long endTime, long minCost, long maxCost, String operationName, - String globalTraceId, int limit, int from) { + String globalTraceId, int limit, int from, ISegmentCostDAO.Sort sort) { logger.debug("startTime: {}, endTime: {}, minCost: {}, maxCost: {}, operationName: {}, globalTraceId: {}, limit: {}, from: {}", startTime, endTime, minCost, maxCost, operationName, globalTraceId, limit, from); ISegmentCostDAO segmentCostDAO = (ISegmentCostDAO)DAOContainer.INSTANCE.get(ISegmentCostDAO.class.getName()); - return segmentCostDAO.loadTop(startTime, endTime, minCost, maxCost, operationName, globalTraceId, limit, from); + return segmentCostDAO.loadTop(startTime, endTime, minCost, maxCost, operationName, globalTraceId, limit, from, sort); } } diff --git a/apm-collector/apm-collector-ui/src/main/resources/META-INF/defines/es_dao.define b/apm-collector/apm-collector-ui/src/main/resources/META-INF/defines/es_dao.define index 6d5f760ce..554fbdddb 100644 --- a/apm-collector/apm-collector-ui/src/main/resources/META-INF/defines/es_dao.define +++ b/apm-collector/apm-collector-ui/src/main/resources/META-INF/defines/es_dao.define @@ -8,6 +8,9 @@ org.skywalking.apm.collector.ui.dao.ApplicationEsDAO org.skywalking.apm.collector.ui.dao.ServiceNameEsDAO org.skywalking.apm.collector.ui.dao.InstanceEsDAO org.skywalking.apm.collector.ui.dao.InstPerformanceEsDAO +org.skywalking.apm.collector.ui.dao.CpuMetricEsDAO org.skywalking.apm.collector.ui.dao.GCMetricEsDAO +org.skywalking.apm.collector.ui.dao.MemoryMetricEsDAO +org.skywalking.apm.collector.ui.dao.MemoryPoolMetricEsDAO org.skywalking.apm.collector.ui.dao.ServiceEntryEsDAO org.skywalking.apm.collector.ui.dao.ServiceReferenceEsDAO \ No newline at end of file From a9ff88fe958900e555b75591f06b9dc5af4f6cc1 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Sun, 3 Sep 2017 16:19:44 +0800 Subject: [PATCH 3/4] Performance evaluation test result: Because of laptop just have 100M network card and sata disk, so 100 million segment process by collector used 255 second, but cpu just used 28% --- .../instance/InstanceIDService.java | 4 +- .../worker/cache/ApplicationCache.java | 16 +- .../worker/cache/InstanceCache.java | 2 +- .../worker/cache/ServiceCache.java | 2 +- .../define/GlobalTraceEsTableDefine.java | 2 +- .../ServiceNameRegisterSerialWorker.java | 4 +- .../servicename/dao/ServiceNameEsDAO.java | 2 +- .../cost/define/SegmentCostEsTableDefine.java | 2 +- .../mock/grpc/GrpcSegmentPost.java | 473 ++++++++++++++++++ .../src/test/resources/logback.xml | 16 + .../src/main/resources/application.yml | 2 +- .../collector/core/server/ServerHolder.java | 9 +- .../src/main/proto/RemoteCommonService.proto | 2 +- .../handler/RemoteCommonServiceHandler.java | 34 +- .../stream/worker/RemoteWorkerRef.java | 74 ++- .../stream/worker/impl/AggregationWorker.java | 8 +- .../stream/worker/impl/PersistenceWorker.java | 51 +- .../stream/worker/impl/data/DataCache.java | 10 +- .../worker/impl/data/DataCollection.java | 34 +- .../stream/worker/impl/data/Window.java | 22 +- 20 files changed, 696 insertions(+), 73 deletions(-) create mode 100644 apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/grpc/GrpcSegmentPost.java create mode 100644 apm-collector/apm-collector-agentstream/src/test/resources/logback.xml diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java index fa5f5bb22..40743ab21 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java @@ -1,6 +1,6 @@ package org.skywalking.apm.collector.agentregister.instance; -import org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterRemoteWorker; +import org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceRegisterRemoteWorker; import org.skywalking.apm.collector.agentstream.worker.register.instance.dao.IInstanceDAO; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.storage.dao.DAOContainer; @@ -28,7 +28,7 @@ public class InstanceIDService { StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); InstanceDataDefine.Instance instance = new InstanceDataDefine.Instance("0", applicationId, agentUUID, registerTime, 0, registerTime, osInfo); try { - context.getClusterWorkerContext().lookup(ApplicationRegisterRemoteWorker.WorkerRole.INSTANCE).tell(instance); + context.getClusterWorkerContext().lookup(InstanceRegisterRemoteWorker.WorkerRole.INSTANCE).tell(instance); } catch (WorkerNotFoundException | WorkerInvokeException e) { logger.error(e.getMessage(), e); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ApplicationCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ApplicationCache.java index ff5574029..a38cb3976 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ApplicationCache.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ApplicationCache.java @@ -10,16 +10,26 @@ import org.skywalking.apm.collector.storage.dao.DAOContainer; */ public class ApplicationCache { - private static Cache CACHE = CacheBuilder.newBuilder().maximumSize(1000).build(); + private static Cache CACHE = CacheBuilder.newBuilder().initialCapacity(100).maximumSize(1000).build(); public static int get(String applicationCode) { + int applicationId = 0; try { - return CACHE.get(applicationCode, () -> { + applicationId = CACHE.get(applicationCode, () -> { IApplicationDAO dao = (IApplicationDAO)DAOContainer.INSTANCE.get(IApplicationDAO.class.getName()); return dao.getApplicationId(applicationCode); }); } catch (Throwable e) { - return 0; + return applicationId; } + + if (applicationId == 0) { + IApplicationDAO dao = (IApplicationDAO)DAOContainer.INSTANCE.get(IApplicationDAO.class.getName()); + applicationId = dao.getApplicationId(applicationCode); + if (applicationId != 0) { + CACHE.put(applicationCode, applicationId); + } + } + return applicationId; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java index 54ac23020..e2808ef6f 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java @@ -10,7 +10,7 @@ import org.skywalking.apm.collector.storage.dao.DAOContainer; */ public class InstanceCache { - private static Cache CACHE = CacheBuilder.newBuilder().maximumSize(1000).build(); + private static Cache CACHE = CacheBuilder.newBuilder().initialCapacity(100).maximumSize(5000).build(); public static int get(int applicationInstanceId) { try { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java index 6d0294bd7..1dcc065b0 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java @@ -11,7 +11,7 @@ import org.skywalking.apm.collector.storage.dao.DAOContainer; */ public class ServiceCache { - private static Cache CACHE = CacheBuilder.newBuilder().maximumSize(10000).build(); + private static Cache CACHE = CacheBuilder.newBuilder().initialCapacity(1000).maximumSize(20000).build(); public static String getServiceName(int serviceId) { try { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java index 657606ca8..8a2b244a2 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java @@ -14,7 +14,7 @@ public class GlobalTraceEsTableDefine extends ElasticSearchTableDefine { } @Override public int refreshInterval() { - return 2; + return 5; } @Override public int numberOfShards() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java index 7b98a3556..5d6fd9efd 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java @@ -3,6 +3,7 @@ package org.skywalking.apm.collector.agentstream.worker.register.servicename; import org.skywalking.apm.collector.agentstream.worker.register.IdAutoIncrement; import org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.IServiceNameDAO; import org.skywalking.apm.collector.storage.dao.DAOContainer; +import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.register.ServiceNameDataDefine; import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorker; import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider; @@ -10,7 +11,6 @@ import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; import org.skywalking.apm.collector.stream.worker.Role; import org.skywalking.apm.collector.stream.worker.WorkerException; -import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.stream.worker.selector.ForeverFirstSelector; import org.skywalking.apm.collector.stream.worker.selector.WorkerSelector; import org.slf4j.Logger; @@ -47,8 +47,8 @@ public class ServiceNameRegisterSerialWorker extends AbstractLocalAsyncWorker { } else { int max = dao.getMaxServiceId(); serviceId = IdAutoIncrement.INSTANCE.increment(min, max); - serviceName.setApplicationId(serviceId); serviceName.setId(String.valueOf(serviceId)); + serviceName.setServiceId(serviceId); } dao.save(serviceName); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java index 150a0b601..a2e3c63d8 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java @@ -30,7 +30,7 @@ public class ServiceNameEsDAO extends EsDAO implements IServiceNameDAO { ElasticSearchClient client = getClient(); SearchRequestBuilder searchRequestBuilder = client.prepareSearch(ServiceNameTable.TABLE); - searchRequestBuilder.setTypes("type"); + searchRequestBuilder.setTypes(ServiceNameTable.TABLE_TYPE); searchRequestBuilder.setSearchType(SearchType.QUERY_THEN_FETCH); BoolQueryBuilder builder = QueryBuilders.boolQuery(); builder.must().add(QueryBuilders.termQuery(ServiceNameTable.COLUMN_APPLICATION_ID, applicationId)); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java index a25bf0fba..82f05c550 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java @@ -14,7 +14,7 @@ public class SegmentCostEsTableDefine extends ElasticSearchTableDefine { } @Override public int refreshInterval() { - return 2; + return 5; } @Override public int numberOfShards() { diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/grpc/GrpcSegmentPost.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/grpc/GrpcSegmentPost.java new file mode 100644 index 000000000..3672c633b --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/grpc/GrpcSegmentPost.java @@ -0,0 +1,473 @@ +package org.skywalking.apm.collector.agentstream.mock.grpc; + +import io.grpc.ManagedChannel; +import io.grpc.ManagedChannelBuilder; +import io.grpc.stub.StreamObserver; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import org.junit.Test; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; +import org.skywalking.apm.network.proto.Application; +import org.skywalking.apm.network.proto.ApplicationInstance; +import org.skywalking.apm.network.proto.ApplicationInstanceMapping; +import org.skywalking.apm.network.proto.ApplicationMapping; +import org.skywalking.apm.network.proto.ApplicationRegisterServiceGrpc; +import org.skywalking.apm.network.proto.Downstream; +import org.skywalking.apm.network.proto.InstanceDiscoveryServiceGrpc; +import org.skywalking.apm.network.proto.KeyWithStringValue; +import org.skywalking.apm.network.proto.LogMessage; +import org.skywalking.apm.network.proto.OSInfo; +import org.skywalking.apm.network.proto.RefType; +import org.skywalking.apm.network.proto.ServiceNameCollection; +import org.skywalking.apm.network.proto.ServiceNameDiscoveryServiceGrpc; +import org.skywalking.apm.network.proto.ServiceNameElement; +import org.skywalking.apm.network.proto.ServiceNameMappingCollection; +import org.skywalking.apm.network.proto.SpanLayer; +import org.skywalking.apm.network.proto.SpanObject; +import org.skywalking.apm.network.proto.SpanType; +import org.skywalking.apm.network.proto.TraceSegmentObject; +import org.skywalking.apm.network.proto.TraceSegmentReference; +import org.skywalking.apm.network.proto.TraceSegmentServiceGrpc; +import org.skywalking.apm.network.proto.UniqueId; +import org.skywalking.apm.network.proto.UpstreamSegment; +import org.skywalking.apm.network.trace.component.ComponentsDefine; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class GrpcSegmentPost { + + private final Logger logger = LoggerFactory.getLogger(GrpcSegmentPost.class); + + private AtomicLong sequence = new AtomicLong(1); + + @Test + public void init() { + ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 11800).maxInboundMessageSize(1024 * 1024 * 50).usePlaintext(true).build(); + + int consumerApplicationId = 0; + int providerApplicationId = 0; + int consumerInstanceId = 0; + int providerInstanceId = 0; + int consumerEntryServiceId = 0; + int consumerExitServiceId = 0; + int consumerExitApplicationId = 0; + int providerEntryServiceId = 0; + + while (consumerApplicationId == 0) { + consumerApplicationId = registerApplication(channel, "consumer"); + } + while (consumerExitApplicationId == 0) { + consumerExitApplicationId = registerApplication(channel, "172.25.0.4:20880"); + } + while (providerApplicationId == 0) { + providerApplicationId = registerApplication(channel, "provider"); + } + while (consumerInstanceId == 0) { + consumerInstanceId = registerInstanceId(channel, "ConsumerUUID", consumerApplicationId, "consumer_host_name", 1); + } + while (providerInstanceId == 0) { + providerInstanceId = registerInstanceId(channel, "ProviderUUID", providerApplicationId, "provider_host_name", 2); + } + while (consumerEntryServiceId == 0) { + consumerEntryServiceId = registerServiceId(channel, consumerApplicationId, "/dubbox-case/case/dubbox-rest"); + } + while (consumerExitServiceId == 0) { + consumerExitServiceId = registerServiceId(channel, consumerApplicationId, "org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()"); + } + while (providerEntryServiceId == 0) { + providerEntryServiceId = registerServiceId(channel, providerApplicationId, "org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()"); + } + + Ids ids = new Ids(); + ids.setConsumerApplicationId(consumerApplicationId); + ids.setProviderApplicationId(providerApplicationId); + ids.setConsumerInstanceId(consumerInstanceId); + ids.setProviderInstanceId(providerInstanceId); + ids.setConsumerEntryServiceId(consumerEntryServiceId); + ids.setConsumerExitServiceId(consumerExitServiceId); + ids.setConsumerExitApplicationId(consumerExitApplicationId); + + long startTime = TimeBucketUtils.INSTANCE.getSecondTimeBucket(System.currentTimeMillis()); + logger.info("start time: {}", startTime); + + int count = 10; + ThreadCount threadCount = new ThreadCount(count); + for (int i = 0; i < count; i++) { + Status status = new Status(); + BuildNewSegment buildNewSegment = new BuildNewSegment(channel, ids, threadCount, i, status); + Executors.newSingleThreadExecutor().execute(buildNewSegment); + } + + while (threadCount.getCount() != 0) { + try { + Thread.sleep(100); + } catch (InterruptedException e) { + } + } + long endTime = TimeBucketUtils.INSTANCE.getSecondTimeBucket(System.currentTimeMillis()); + logger.info("end time: {}", endTime); + + channel.shutdownNow(); + while (!channel.isTerminated()) { + try { + channel.awaitTermination(100, TimeUnit.SECONDS); + } catch (InterruptedException e) { + logger.error(e.getMessage(), e); + } + } + } + + private int registerApplication(ManagedChannel channel, String applicationCode) { + ApplicationRegisterServiceGrpc.ApplicationRegisterServiceBlockingStub stub = ApplicationRegisterServiceGrpc.newBlockingStub(channel); + Application application = Application.newBuilder().addApplicationCode(applicationCode).build(); + ApplicationMapping mapping = stub.register(application); + int applicationId = mapping.getApplication(0).getValue(); + + try { + Thread.sleep(10); + } catch (InterruptedException e) { + } + return applicationId; + } + + private int registerInstanceId(ManagedChannel channel, String agentUUId, Integer applicationId, + String hostName, int processNo) { + InstanceDiscoveryServiceGrpc.InstanceDiscoveryServiceBlockingStub stub = InstanceDiscoveryServiceGrpc.newBlockingStub(channel); + ApplicationInstance.Builder instance = ApplicationInstance.newBuilder(); + instance.setApplicationId(applicationId); + instance.setRegisterTime(System.currentTimeMillis()); + instance.setAgentUUID(agentUUId); + + OSInfo.Builder osInfo = OSInfo.newBuilder(); + osInfo.setHostname(hostName); + osInfo.setOsName("Linux"); + osInfo.setProcessNo(processNo); + osInfo.addIpv4S("10.0.0.1"); + osInfo.addIpv4S("10.0.0.2"); + instance.setOsinfo(osInfo.build()); + + ApplicationInstanceMapping mapping = stub.register(instance.build()); + int instanceId = mapping.getApplicationInstanceId(); + + try { + Thread.sleep(10); + } catch (InterruptedException e) { + } + return instanceId; + } + + private int registerServiceId(ManagedChannel channel, int applicationId, String serviceName) { + ServiceNameDiscoveryServiceGrpc.ServiceNameDiscoveryServiceBlockingStub stub = ServiceNameDiscoveryServiceGrpc.newBlockingStub(channel); + ServiceNameCollection.Builder collection = ServiceNameCollection.newBuilder(); + + ServiceNameElement.Builder element = ServiceNameElement.newBuilder(); + element.setApplicationId(applicationId); + element.setServiceName(serviceName); + collection.addElements(element); + + ServiceNameMappingCollection mappingCollection = stub.discovery(collection.build()); + int serviceId = mappingCollection.getElements(0).getServiceId(); + + try { + Thread.sleep(10); + } catch (InterruptedException e) { + } + return serviceId; + } + + class BuildNewSegment implements Runnable { + private final ManagedChannel segmentChannel; + private final Ids ids; + private final ThreadCount threadCount; + private final int procNo; + private final Status status; + private StreamObserver streamObserver; + + public BuildNewSegment(ManagedChannel segmentChannel, + Ids ids, ThreadCount threadCount, int procNo, + Status status) { + this.segmentChannel = segmentChannel; + this.ids = ids; + this.threadCount = threadCount; + this.procNo = procNo; + this.status = status; + } + + @Override public void run() { + statusChange(); + int i = 0; + while (i < 50000) { + send(streamObserver, ids); + + i++; + if (i % 10000 == 0) { + logger.info("process no: {}, send segment count: {}", procNo, i); + streamObserver.onCompleted(); + while (!status.isFinish) { + try { + Thread.sleep(10); + } catch (InterruptedException e) { + } + } + status.setFinish(false); + statusChange(); + } + } + this.threadCount.finishOne(); + } + + private void statusChange() { + TraceSegmentServiceGrpc.TraceSegmentServiceStub stub = TraceSegmentServiceGrpc.newStub(segmentChannel); + streamObserver = stub.collect(new StreamObserver() { + @Override public void onNext(Downstream downstream) { + } + + @Override public void onError(Throwable throwable) { + logger.error(throwable.getMessage(), throwable); + } + + @Override public void onCompleted() { + status.setFinish(true); + logger.info("process no: {}, server completed", procNo); + } + }); + } + } + + public void send(StreamObserver streamObserver, Ids ids) { + long now = System.currentTimeMillis(); + UniqueId consumerSegmentId = createSegmentId(); + UniqueId providerSegmentId = createSegmentId(); + + streamObserver.onNext(createConsumerSegment(consumerSegmentId, ids, now)); + streamObserver.onNext(createProviderSegment(consumerSegmentId, providerSegmentId, ids, now)); + } + + private UpstreamSegment createConsumerSegment(UniqueId segmentId, Ids ids, long timestamp) { + UpstreamSegment.Builder upstream = UpstreamSegment.newBuilder(); + upstream.addGlobalTraceIds(segmentId); + + TraceSegmentObject.Builder segmentBuilder = TraceSegmentObject.newBuilder(); + segmentBuilder.setApplicationId(ids.consumerApplicationId); + segmentBuilder.setApplicationInstanceId(ids.consumerInstanceId); + segmentBuilder.setTraceSegmentId(segmentId); + + SpanObject.Builder entrySpan = SpanObject.newBuilder(); + entrySpan.setSpanId(0); + entrySpan.setSpanType(SpanType.Entry); + entrySpan.setSpanLayer(SpanLayer.Http); + entrySpan.setParentSpanId(-1); + entrySpan.setStartTime(timestamp); + entrySpan.setEndTime(timestamp + 3000); + entrySpan.setComponentId(ComponentsDefine.TOMCAT.getId()); + entrySpan.setOperationNameId(ids.getConsumerEntryServiceId()); + entrySpan.setIsError(false); + + LogMessage.Builder entryLogMessage = LogMessage.newBuilder(); + entryLogMessage.setTime(timestamp); + + KeyWithStringValue.Builder data_1 = KeyWithStringValue.newBuilder(); + data_1.setKey("url"); + data_1.setValue("http://localhost:18080/dubbox-case/case/dubbox-rest"); + entryLogMessage.addData(data_1); + + KeyWithStringValue.Builder data_2 = KeyWithStringValue.newBuilder(); + data_2.setKey("http.method"); + data_2.setValue("GET"); + entryLogMessage.addData(data_2); + entrySpan.addLogs(entryLogMessage); + segmentBuilder.addSpans(entrySpan); + + SpanObject.Builder exitSpan = SpanObject.newBuilder(); + exitSpan.setSpanId(1); + exitSpan.setSpanType(SpanType.Exit); + exitSpan.setSpanLayer(SpanLayer.RPCFramework); + exitSpan.setParentSpanId(0); + exitSpan.setStartTime(timestamp + 500); + exitSpan.setEndTime(timestamp + 2500); + exitSpan.setComponentId(ComponentsDefine.TOMCAT.getId()); + exitSpan.setOperationNameId(ids.getConsumerExitServiceId()); + exitSpan.setPeerId(ids.consumerExitApplicationId); + exitSpan.setIsError(false); + + LogMessage.Builder exitLogMessage = LogMessage.newBuilder(); + exitLogMessage.setTime(timestamp); + + KeyWithStringValue.Builder data = KeyWithStringValue.newBuilder(); + data.setKey("url"); + data.setValue("rest://172.25.0.4:20880/org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()"); + exitLogMessage.addData(data); + exitSpan.addLogs(exitLogMessage); + segmentBuilder.addSpans(exitSpan); + + upstream.setSegment(segmentBuilder.build().toByteString()); + return upstream.build(); + } + + private UpstreamSegment createProviderSegment(UniqueId consumerSegmentId, UniqueId providerSegmentId, Ids ids, + long timestamp) { + UpstreamSegment.Builder upstream = UpstreamSegment.newBuilder(); + upstream.addGlobalTraceIds(consumerSegmentId); + + TraceSegmentObject.Builder segmentBuilder = TraceSegmentObject.newBuilder(); + segmentBuilder.setApplicationId(ids.providerApplicationId); + segmentBuilder.setApplicationInstanceId(ids.providerInstanceId); + segmentBuilder.setTraceSegmentId(providerSegmentId); + + TraceSegmentReference.Builder referenceBuilder = TraceSegmentReference.newBuilder(); + referenceBuilder.setParentTraceSegmentId(consumerSegmentId); + referenceBuilder.setParentApplicationInstanceId(ids.getConsumerInstanceId()); + referenceBuilder.setParentSpanId(1); + referenceBuilder.setParentServiceId(ids.getConsumerExitServiceId()); + referenceBuilder.setEntryApplicationInstanceId(ids.getConsumerInstanceId()); + referenceBuilder.setEntryServiceId(ids.getConsumerEntryServiceId()); + referenceBuilder.setNetworkAddressId(ids.consumerExitApplicationId); + referenceBuilder.setRefType(RefType.CrossProcess); + segmentBuilder.addRefs(referenceBuilder); + + SpanObject.Builder entrySpan = SpanObject.newBuilder(); + entrySpan.setSpanId(0); + entrySpan.setSpanType(SpanType.Entry); + entrySpan.setSpanLayer(SpanLayer.RPCFramework); + entrySpan.setParentSpanId(-1); + entrySpan.setStartTime(timestamp + 1000); + entrySpan.setEndTime(timestamp + 2000); + entrySpan.setComponentId(ComponentsDefine.TOMCAT.getId()); + entrySpan.setOperationNameId(ids.getProviderEntryServiceId()); + entrySpan.setIsError(false); + + LogMessage.Builder entryLogMessage = LogMessage.newBuilder(); + entryLogMessage.setTime(timestamp); + + KeyWithStringValue.Builder data_1 = KeyWithStringValue.newBuilder(); + data_1.setKey("url"); + data_1.setValue("rest://172.25.0.4:20880/org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()"); + entryLogMessage.addData(data_1); + + KeyWithStringValue.Builder data_2 = KeyWithStringValue.newBuilder(); + data_2.setKey("http.method"); + data_2.setValue("GET"); + entryLogMessage.addData(data_2); + entrySpan.addLogs(entryLogMessage); + segmentBuilder.addSpans(entrySpan); + + upstream.setSegment(segmentBuilder.build().toByteString()); + return upstream.build(); + } + + private UniqueId createSegmentId() { + long id = sequence.getAndIncrement(); + UniqueId.Builder builder = UniqueId.newBuilder(); + builder.addIdParts(id); + builder.addIdParts(id); + builder.addIdParts(id); + return builder.build(); + } + + class Ids { + private int consumerApplicationId = 0; + private int providerApplicationId = 0; + private int consumerInstanceId = 0; + private int providerInstanceId = 0; + private int consumerEntryServiceId = 0; + private int consumerExitServiceId = 0; + private int consumerExitApplicationId = 0; + private int providerEntryServiceId = 0; + + public int getConsumerApplicationId() { + return consumerApplicationId; + } + + public void setConsumerApplicationId(int consumerApplicationId) { + this.consumerApplicationId = consumerApplicationId; + } + + public int getProviderApplicationId() { + return providerApplicationId; + } + + public void setProviderApplicationId(int providerApplicationId) { + this.providerApplicationId = providerApplicationId; + } + + public int getConsumerInstanceId() { + return consumerInstanceId; + } + + public void setConsumerInstanceId(int consumerInstanceId) { + this.consumerInstanceId = consumerInstanceId; + } + + public int getProviderInstanceId() { + return providerInstanceId; + } + + public void setProviderInstanceId(int providerInstanceId) { + this.providerInstanceId = providerInstanceId; + } + + public int getConsumerEntryServiceId() { + return consumerEntryServiceId; + } + + public void setConsumerEntryServiceId(int consumerEntryServiceId) { + this.consumerEntryServiceId = consumerEntryServiceId; + } + + public int getConsumerExitServiceId() { + return consumerExitServiceId; + } + + public void setConsumerExitServiceId(int consumerExitServiceId) { + this.consumerExitServiceId = consumerExitServiceId; + } + + public int getConsumerExitApplicationId() { + return consumerExitApplicationId; + } + + public void setConsumerExitApplicationId(int consumerExitApplicationId) { + this.consumerExitApplicationId = consumerExitApplicationId; + } + + public int getProviderEntryServiceId() { + return providerEntryServiceId; + } + + public void setProviderEntryServiceId(int providerEntryServiceId) { + this.providerEntryServiceId = providerEntryServiceId; + } + } + + class ThreadCount { + private int count; + + public ThreadCount(int count) { + this.count = count; + } + + public void finishOne() { + count--; + } + + public int getCount() { + return count; + } + } + + class Status { + private boolean isFinish = false; + + public boolean isFinish() { + return isFinish; + } + + public void setFinish(boolean finish) { + isFinish = finish; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/test/resources/logback.xml b/apm-collector/apm-collector-agentstream/src/test/resources/logback.xml new file mode 100644 index 000000000..46eba2b93 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/resources/logback.xml @@ -0,0 +1,16 @@ + + + + + %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + + + + + + + + \ No newline at end of file diff --git a/apm-collector/apm-collector-boot/src/main/resources/application.yml b/apm-collector/apm-collector-boot/src/main/resources/application.yml index 1a84f5740..b1d440ec5 100644 --- a/apm-collector/apm-collector-boot/src/main/resources/application.yml +++ b/apm-collector/apm-collector-boot/src/main/resources/application.yml @@ -25,6 +25,6 @@ storage: elasticsearch: cluster_name: CollectorDBCluster cluster_transport_sniffer: true - cluster_nodes: 127.0.0.1:9300 + cluster_nodes: 10.0.0.19:9300,10.0.0.6:9300 diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java index 96edf8f02..2653ab530 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java @@ -4,12 +4,16 @@ import java.util.LinkedList; import java.util.List; import org.skywalking.apm.collector.core.framework.Handler; import org.skywalking.apm.collector.core.util.CollectionUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class ServerHolder { + private final Logger logger = LoggerFactory.getLogger(ServerHolder.class); + private List servers; public ServerHolder() { @@ -33,7 +37,10 @@ public class ServerHolder { private void addHandler(List handlers, Server server) { if (CollectionUtils.isNotEmpty(handlers)) { - handlers.forEach(handler -> server.addHandler(handler)); + handlers.forEach(handler -> { + server.addHandler(handler); + logger.debug("add handler into server: {}, handler name: {}", server.hostPort(), handler.getClass().getName()); + }); } } diff --git a/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto b/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto index 159bfe9bc..b2100424f 100644 --- a/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto +++ b/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto @@ -4,7 +4,7 @@ option java_multiple_files = true; option java_package = "org.skywalking.apm.collector.remote.grpc.proto"; service RemoteCommonService { - rpc call (RemoteMessage) returns (Empty) { + rpc call (stream RemoteMessage) returns (Empty) { } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java index 86b5dc2b0..d1d92ebd3 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java @@ -22,17 +22,29 @@ public class RemoteCommonServiceHandler extends RemoteCommonServiceGrpc.RemoteCo private final Logger logger = LoggerFactory.getLogger(RemoteCommonServiceHandler.class); - @Override public void call(RemoteMessage request, StreamObserver responseObserver) { - String roleName = request.getWorkerRole(); - RemoteData remoteData = request.getRemoteData(); + @Override public StreamObserver call(StreamObserver responseObserver) { + return new StreamObserver() { + @Override public void onNext(RemoteMessage message) { + String roleName = message.getWorkerRole(); + RemoteData remoteData = message.getRemoteData(); - StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); - Role role = context.getClusterWorkerContext().getRole(roleName); - Object object = role.dataDefine().deserialize(remoteData); - try { - context.getClusterWorkerContext().lookupInSide(roleName).tell(object); - } catch (WorkerNotFoundException | WorkerInvokeException e) { - logger.error(e.getMessage(), e); - } + StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); + Role role = context.getClusterWorkerContext().getRole(roleName); + Object object = role.dataDefine().deserialize(remoteData); + try { + context.getClusterWorkerContext().lookupInSide(roleName).tell(object); + } catch (WorkerNotFoundException | WorkerInvokeException e) { + logger.error(e.getMessage(), e); + } + } + + @Override public void onError(Throwable throwable) { + logger.error(throwable.getMessage(), throwable); + } + + @Override public void onCompleted() { + responseObserver.onCompleted(); + } + }; } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/RemoteWorkerRef.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/RemoteWorkerRef.java index 5e99ab101..707230f6c 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/RemoteWorkerRef.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/RemoteWorkerRef.java @@ -1,17 +1,24 @@ package org.skywalking.apm.collector.stream.worker; +import io.grpc.stub.StreamObserver; import org.skywalking.apm.collector.client.grpc.GRPCClient; +import org.skywalking.apm.collector.remote.grpc.proto.Empty; import org.skywalking.apm.collector.remote.grpc.proto.RemoteCommonServiceGrpc; import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; import org.skywalking.apm.collector.remote.grpc.proto.RemoteMessage; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class RemoteWorkerRef extends WorkerRef { + private final Logger logger = LoggerFactory.getLogger(RemoteWorkerRef.class); + private final Boolean acrossJVM; - private final RemoteCommonServiceGrpc.RemoteCommonServiceBlockingStub stub; + private final RemoteCommonServiceGrpc.RemoteCommonServiceStub stub; + private StreamObserver streamObserver; private final AbstractRemoteWorker remoteWorker; public RemoteWorkerRef(Role role, AbstractRemoteWorker remoteWorker) { @@ -25,7 +32,8 @@ public class RemoteWorkerRef extends WorkerRef { super(role); this.remoteWorker = null; this.acrossJVM = true; - this.stub = RemoteCommonServiceGrpc.newBlockingStub(client.getChannel()); + this.stub = RemoteCommonServiceGrpc.newStub(client.getChannel()); + createStreamObserver(); } @Override @@ -36,7 +44,8 @@ public class RemoteWorkerRef extends WorkerRef { RemoteMessage.Builder builder = RemoteMessage.newBuilder(); builder.setWorkerRole(getRole().roleName()); builder.setRemoteData(remoteData); - stub.call(builder.build()); + + streamObserver.onNext(builder.build()); } else { remoteWorker.allocateJob(message); } @@ -45,4 +54,63 @@ public class RemoteWorkerRef extends WorkerRef { public Boolean isAcrossJVM() { return acrossJVM; } + + private void createStreamObserver() { + StreamStatus status = new StreamStatus(false); + streamObserver = stub.call(new StreamObserver() { + @Override public void onNext(Empty empty) { + } + + @Override public void onError(Throwable throwable) { + logger.error(throwable.getMessage(), throwable); + } + + @Override public void onCompleted() { + status.finished(); + } + }); + } + + class StreamStatus { + private volatile boolean status; + + public StreamStatus(boolean status) { + this.status = status; + } + + public boolean isFinish() { + return status; + } + + public void finished() { + this.status = true; + } + + /** + * @param maxTimeout max wait time, milliseconds. + */ + public void wait4Finish(long maxTimeout) { + long time = 0; + while (!status) { + if (time > maxTimeout) { + break; + } + try2Sleep(5); + time += 5; + } + } + + /** + * Try to sleep, and ignore the {@link InterruptedException} + * + * @param millis the length of time to sleep in milliseconds + */ + private void try2Sleep(long millis) { + try { + Thread.sleep(millis); + } catch (InterruptedException e) { + + } + } + } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java index 4cc4676d6..255d07016 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java @@ -51,7 +51,7 @@ public abstract class AggregationWorker extends AbstractLocalAsyncWorker { private void sendToNext() throws WorkerException { dataCache.switchPointer(); - while (dataCache.getLast().isHolding()) { + while (dataCache.getLast().isWriting()) { try { Thread.sleep(10); } catch (InterruptedException e) { @@ -66,17 +66,17 @@ public abstract class AggregationWorker extends AbstractLocalAsyncWorker { logger.error(e.getMessage(), e); } }); - dataCache.releaseLast(); + dataCache.finishReadingLast(); } protected final void aggregate(Object message) { Data data = (Data)message; - dataCache.hold(); + dataCache.writing(); if (dataCache.containsKey(data.id())) { getRole().dataDefine().mergeData(data, dataCache.get(data.id())); } else { dataCache.put(data.id(), data); } - dataCache.release(); + dataCache.finishWriting(); } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java index 875cba06b..c1e509c20 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java @@ -1,11 +1,13 @@ package org.skywalking.apm.collector.stream.worker.impl; -import java.util.ArrayList; +import java.util.LinkedList; import java.util.List; import java.util.Map; import org.skywalking.apm.collector.core.queue.EndOfBatchCommand; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.core.util.ObjectUtils; +import org.skywalking.apm.collector.storage.dao.DAOContainer; +import org.skywalking.apm.collector.storage.dao.IBatchDAO; import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorker; import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; @@ -35,24 +37,37 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { } @Override protected final void onWork(Object message) throws WorkerException { - if (message instanceof EndOfBatchCommand || message instanceof FlushAndSwitch) { - if (dataCache.trySwitchPointer()) { - dataCache.switchPointer(); - } - } else { - if (dataCache.currentCollectionSize() >= 1000) { + if (message instanceof FlushAndSwitch) { + try { if (dataCache.trySwitchPointer()) { dataCache.switchPointer(); } + } finally { + dataCache.trySwitchPointerFinally(); + } + } else if (message instanceof EndOfBatchCommand) { + } else { + if (dataCache.currentCollectionSize() >= 5000) { + try { + if (dataCache.trySwitchPointer()) { + dataCache.switchPointer(); + + List collection = buildBatchCollection(); + IBatchDAO dao = (IBatchDAO)DAOContainer.INSTANCE.get(IBatchDAO.class.getName()); + dao.batchPersistence(collection); + } + } finally { + dataCache.trySwitchPointerFinally(); + } } aggregate(message); } } public final List buildBatchCollection() throws WorkerException { - List batchCollection; + List batchCollection = new LinkedList<>(); try { - while (dataCache.getLast().isHolding()) { + while (dataCache.getLast().isWriting()) { try { Thread.sleep(10); } catch (InterruptedException e) { @@ -60,16 +75,18 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { } } - batchCollection = prepareBatch(dataCache.getLast().asMap()); + if (dataCache.getLast().asMap() != null) { + batchCollection = prepareBatch(dataCache.getLast().asMap()); + } } finally { - dataCache.releaseLast(); + dataCache.finishReadingLast(); } return batchCollection; } protected final List prepareBatch(Map dataMap) { - List insertBatchCollection = new ArrayList<>(); - List updateBatchCollection = new ArrayList<>(); + List insertBatchCollection = new LinkedList<>(); + List updateBatchCollection = new LinkedList<>(); dataMap.forEach((id, data) -> { if (needMergeDBData()) { Data dbData = persistenceDAO().get(id, getRole().dataDefine()); @@ -101,18 +118,16 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { } private void aggregate(Object message) { - dataCache.hold(); + dataCache.writing(); Data data = (Data)message; if (dataCache.containsKey(data.id())) { getRole().dataDefine().mergeData(data, dataCache.get(data.id())); } else { - if (dataCache.currentCollectionSize() < 1000) { - dataCache.put(data.id(), data); - } + dataCache.put(data.id(), data); } - dataCache.release(); + dataCache.finishWriting(); } protected abstract IPersistenceDAO persistenceDAO(); diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java index ff81c3ea8..459e7062f 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java @@ -21,16 +21,16 @@ public class DataCache extends Window { lockedDataCollection.put(id, data); } - public void hold() { - lockedDataCollection = getCurrentAndHold(); + public void writing() { + lockedDataCollection = getCurrentAndWriting(); } public int currentCollectionSize() { - return getCurrentAndHold().size(); + return getCurrent().size(); } - public void release() { - lockedDataCollection.release(); + public void finishWriting() { + lockedDataCollection.finishWriting(); lockedDataCollection = null; } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java index bc280d8d4..ee599d983 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java @@ -1,7 +1,7 @@ package org.skywalking.apm.collector.stream.worker.impl.data; -import java.util.HashMap; import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import org.skywalking.apm.collector.core.stream.Data; /** @@ -9,23 +9,37 @@ import org.skywalking.apm.collector.core.stream.Data; */ public class DataCollection { private Map data; - private volatile boolean isHold; + private volatile boolean writing; + private volatile boolean reading; public DataCollection() { - this.data = new HashMap<>(); - this.isHold = false; + this.data = new ConcurrentHashMap<>(); + this.writing = false; + this.reading = false; } - public void release() { - isHold = false; + public void finishWriting() { + writing = false; } - public void hold() { - isHold = true; + public void writing() { + writing = true; } - public boolean isHolding() { - return isHold; + public boolean isWriting() { + return writing; + } + + public void finishReading() { + reading = false; + } + + public void reading() { + reading = true; + } + + public boolean isReading() { + return reading; } public boolean containsKey(String key) { diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java index 887f24809..f2f729f7d 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java @@ -21,32 +21,40 @@ public abstract class Window { } public boolean trySwitchPointer() { - if (windowSwitch.incrementAndGet() == 1) { + if (windowSwitch.incrementAndGet() == 1 && !getLast().isReading()) { return true; } else { - windowSwitch.addAndGet(-1); return false; } } + public void trySwitchPointerFinally() { + windowSwitch.addAndGet(-1); + } + public void switchPointer() { if (pointer == windowDataA) { pointer = windowDataB; } else { pointer = windowDataA; } + getLast().reading(); } - protected DataCollection getCurrentAndHold() { + protected DataCollection getCurrentAndWriting() { if (pointer == windowDataA) { - windowDataA.hold(); + windowDataA.writing(); return windowDataA; } else { - windowDataB.hold(); + windowDataB.writing(); return windowDataB; } } + protected DataCollection getCurrent() { + return pointer; + } + public DataCollection getLast() { if (pointer == windowDataA) { return windowDataB; @@ -55,8 +63,8 @@ public abstract class Window { } } - public void releaseLast() { + public void finishReadingLast() { getLast().clear(); - windowSwitch.addAndGet(-1); + getLast().finishReading(); } } From 2178a736a91f5c8593013248c5521c701d658d17 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Sun, 3 Sep 2017 16:39:26 +0800 Subject: [PATCH 4/4] fixed check style error --- .../apm/collector/stream/worker/impl/data/Window.java | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java index f2f729f7d..851285f01 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java @@ -21,11 +21,7 @@ public abstract class Window { } public boolean trySwitchPointer() { - if (windowSwitch.incrementAndGet() == 1 && !getLast().isReading()) { - return true; - } else { - return false; - } + return windowSwitch.incrementAndGet() == 1 && !getLast().isReading(); } public void trySwitchPointerFinally() {