From a6d9111d7cf47e9e141192419087e980a6c7fb81 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Sun, 20 Aug 2017 14:27:08 +0800 Subject: [PATCH 1/2] no message --- .../apm/collector/ui/dao/IInstanceDAO.java | 3 ++ .../apm/collector/ui/dao/InstanceEsDAO.java | 20 ++++++++ .../apm/collector/ui/dao/InstanceH2DAO.java | 5 ++ .../InstanceMetricGetHandler.java | 49 +++++++++++++++++++ .../InstanceOsInfoGetHandler.java | 41 ++++++++++++++++ .../ui/service/InstanceJVMService.java | 43 ++++++++++++++++ 6 files changed, 161 insertions(+) create mode 100644 apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceMetricGetHandler.java create mode 100644 apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceOsInfoGetHandler.java create mode 100644 apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceJVMService.java diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IInstanceDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IInstanceDAO.java index 3d40c136b..ce652e15b 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IInstanceDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IInstanceDAO.java @@ -1,6 +1,7 @@ package org.skywalking.apm.collector.ui.dao; import java.util.List; +import org.skywalking.apm.collector.storage.table.instance.Instance; /** * @author pengys5 @@ -12,6 +13,8 @@ public interface IInstanceDAO { List getApplications(long time); + Instance getInstance(int instanceId); + class Application { private final int applicationId; private final long count; diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceEsDAO.java index ecef9bf0e..b851fafd1 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceEsDAO.java @@ -2,6 +2,8 @@ package org.skywalking.apm.collector.ui.dao; import java.util.LinkedList; import java.util.List; +import org.elasticsearch.action.get.GetRequestBuilder; +import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.search.SearchRequestBuilder; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.action.search.SearchType; @@ -16,6 +18,7 @@ import org.elasticsearch.search.aggregations.bucket.terms.Terms; import org.elasticsearch.search.sort.SortBuilders; import org.elasticsearch.search.sort.SortMode; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.storage.table.instance.Instance; import org.skywalking.apm.collector.storage.table.register.InstanceTable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -93,4 +96,21 @@ public class InstanceEsDAO extends EsDAO implements IInstanceDAO { } return applications; } + + @Override public Instance getInstance(int instanceId) { + logger.debug("get instance info, instance id: {}", instanceId); + GetRequestBuilder requestBuilder = getClient().prepareGet(InstanceTable.TABLE, String.valueOf(instanceId)); + GetResponse getResponse = requestBuilder.get(); + if (getResponse.isExists()) { + Instance instance = new Instance(); + instance.setId(String.valueOf(instanceId)); + instance.setApplicationId(((Number)getResponse.getSource().get(InstanceTable.COLUMN_APPLICATION_ID)).intValue()); + instance.setAgentUUID((String)getResponse.getSource().get(InstanceTable.COLUMN_AGENT_UUID)); + instance.setRegisterTime(((Number)getResponse.getSource().get(InstanceTable.COLUMN_REGISTER_TIME)).longValue()); + instance.setHeartBeatTime(((Number)getResponse.getSource().get(InstanceTable.COLUMN_HEARTBEAT_TIME)).longValue()); + instance.setOsInfo((String)getResponse.getSource().get(InstanceTable.COLUMN_OS_INFO)); + return instance; + } + return null; + } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java index 2e1badda7..0ba423247 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java @@ -2,6 +2,7 @@ package org.skywalking.apm.collector.ui.dao; import java.util.List; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.table.instance.Instance; /** * @author pengys5 @@ -18,4 +19,8 @@ public class InstanceH2DAO extends H2DAO implements IInstanceDAO { @Override public List getApplications(long time) { return null; } + + @Override public Instance getInstance(int instanceId) { + return null; + } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceMetricGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceMetricGetHandler.java new file mode 100644 index 000000000..5e32d090f --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceMetricGetHandler.java @@ -0,0 +1,49 @@ +package org.skywalking.apm.collector.ui.jetty.handler.instancemetric; + +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.service.InstanceHealthService; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class InstanceMetricGetHandler extends JettyHandler { + + private final Logger logger = LoggerFactory.getLogger(InstanceMetricGetHandler.class); + + @Override public String pathSpec() { + return "/instance/jvm/instanceId"; + } + + private InstanceHealthService service = new InstanceHealthService(); + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + String timestampStr = req.getParameter("timestamp"); + String applicationIdStr = req.getParameter("applicationId"); + logger.debug("instance health get timestamp: {}", timestampStr); + + long timestamp; + try { + timestamp = Long.parseLong(timestampStr); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("timestamp must be long"); + } + + int applicationId; + try { + applicationId = Integer.parseInt(applicationIdStr); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("application id must be integer"); + } + + return service.getInstances(timestamp, applicationId); + } + + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { + throw new UnsupportedOperationException(); + } +} diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceOsInfoGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceOsInfoGetHandler.java new file mode 100644 index 000000000..4d4180705 --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceOsInfoGetHandler.java @@ -0,0 +1,41 @@ +package org.skywalking.apm.collector.ui.jetty.handler.instancemetric; + +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.service.InstanceJVMService; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class InstanceOsInfoGetHandler extends JettyHandler { + + private final Logger logger = LoggerFactory.getLogger(InstanceOsInfoGetHandler.class); + + @Override public String pathSpec() { + return "/instance/os/instanceId"; + } + + private InstanceJVMService service = new InstanceJVMService(); + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + String instanceIdStr = req.getParameter("instanceId"); + logger.debug("instance os info get, instance id: {}", instanceIdStr); + + int instanceId; + try { + instanceId = Integer.parseInt(instanceIdStr); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("instance id must be integer"); + } + + return service.getInstanceOsInfo(instanceId); + } + + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { + throw new UnsupportedOperationException(); + } +} 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 new file mode 100644 index 000000000..5924bf367 --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceJVMService.java @@ -0,0 +1,43 @@ +package org.skywalking.apm.collector.ui.service; + +import com.google.gson.Gson; +import com.google.gson.JsonObject; +import org.skywalking.apm.collector.core.framework.UnexpectedException; +import org.skywalking.apm.collector.core.util.ObjectUtils; +import org.skywalking.apm.collector.storage.dao.DAOContainer; +import org.skywalking.apm.collector.storage.table.instance.Instance; +import org.skywalking.apm.collector.ui.dao.IInstanceDAO; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class InstanceJVMService { + + private final Logger logger = LoggerFactory.getLogger(InstanceJVMService.class); + + private Gson gson = new Gson(); + + public JsonObject getInstanceOsInfo(int instanceId) { + IInstanceDAO instanceDAO = (IInstanceDAO)DAOContainer.INSTANCE.get(IInstanceDAO.class.getName()); + Instance instance = instanceDAO.getInstance(instanceId); + if (ObjectUtils.isEmpty(instance)) { + throw new UnexpectedException("instance id: " + instance + " not exist."); + } + + JsonObject response = gson.fromJson(instance.getOsInfo(), JsonObject.class); + return response; + } + + public JsonObject getInstanceJvmMetric(int instanceId, String metricType) { + + return null; + } + + public enum MetricType { + cpu, gc, tps, heapmemory, heappermgen, heapmetaspace, heapnewgen, + heapoldgen, heapsurvivor, nonheap, memory, nonheappermgen, nonheapmetaspace, + nonheapnewgen, nonheapoldgen, nonheapsurvivor + } +} From 166112074bed442264e74ad723adb1fbfbc52576 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Mon, 21 Aug 2017 23:59:03 +0800 Subject: [PATCH 2/2] Instance health web ui working with collector success. #365 --- .../handler/JVMMetricsServiceHandler.java | 3 +- .../worker/cpu/dao/CpuMetricEsDAO.java | 9 ++++- .../agentjvm/worker/gc/dao/GCMetricEsDAO.java | 7 ++-- .../gc/define/GCMetricEsTableDefine.java | 3 +- .../gc/define/GCMetricH2TableDefine.java | 1 + .../JVMMetricsServiceHandlerTestCase.java | 35 +++++++++++------- .../global/GlobalTraceSpanListener.java | 2 +- .../InstPerformanceSpanListener.java | 5 ++- .../performance/dao/InstPerformanceEsDAO.java | 18 ++++++--- .../define/InstPerformanceEsTableDefine.java | 3 +- .../define/InstPerformanceH2TableDefine.java | 3 +- .../component/NodeComponentSpanListener.java | 2 +- .../node/mapping/NodeMappingSpanListener.java | 2 +- .../reference/NodeRefSpanListener.java | 2 +- .../summary/NodeRefSumSpanListener.java | 2 +- .../segment/cost/SegmentCostSpanListener.java | 2 +- .../entry/ServiceEntrySpanListener.java | 2 +- .../reference/ServiceRefSpanListener.java | 2 +- .../agentstream/mock/SegmentPost.java | 35 ++++++++++++++---- .../src/main/resources/logback.xml | 16 ++++++++ .../core/storage/StorageInstaller.java | 5 ++- .../apm/collector/core/stream/Data.java | 1 + .../collector/core}/util/TimeBucketUtils.java | 13 ++++++- .../core/utils/TimeBucketUtilsTestCase.java | 23 ++++++++++++ .../collector/storage/define/CommonTable.java | 1 + .../collector/storage/define/DataDefine.java | 20 ++++++---- .../instance/InstPerformanceDataDefine.java | 18 ++++++++- .../define/jvm/GCMetricDataDefine.java | 26 ++++++++++--- .../stream/worker/impl/PersistenceWorker.java | 20 ++++++++-- .../worker/util/TimeBucketUtilsTestCase.java | 2 +- .../collector/ui/cache/ApplicationCache.java | 10 +++++ .../apm/collector/ui/dao/GCMetricEsDAO.java | 27 ++++++++++---- .../apm/collector/ui/dao/GCMetricH2DAO.java | 9 +++++ .../collector/ui/dao/IInstPerformanceDAO.java | 2 +- .../ui/dao/InstPerformanceEsDAO.java | 37 +++++++++++++------ .../ui/dao/InstPerformanceH2DAO.java | 9 +++++ .../apm/collector/ui/dao/InstanceH2DAO.java | 4 +- .../ApplicationsGetHandler.java | 10 ++--- .../InstanceHealthGetHandler.java | 28 ++++++++++---- .../time/AllInstanceLastTimeGetHandler.java | 10 ++++- .../ui/service/InstanceHealthService.java | 28 +++++++------- 41 files changed, 339 insertions(+), 118 deletions(-) create mode 100644 apm-collector/apm-collector-boot/src/main/resources/logback.xml rename apm-collector/{apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker => apm-collector-core/src/main/java/org/skywalking/apm/collector/core}/util/TimeBucketUtils.java (82%) create mode 100644 apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/utils/TimeBucketUtilsTestCase.java diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandler.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandler.java index 201a1fb7c..3b6f8a4f5 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandler.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandler.java @@ -10,6 +10,7 @@ import org.skywalking.apm.collector.agentjvm.worker.memory.MemoryMetricPersisten import org.skywalking.apm.collector.agentjvm.worker.memorypool.MemoryPoolMetricPersistenceWorker; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.server.grpc.GRPCHandler; import org.skywalking.apm.collector.storage.define.jvm.CpuMetricDataDefine; import org.skywalking.apm.collector.storage.define.jvm.GCMetricDataDefine; @@ -19,7 +20,6 @@ import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; import org.skywalking.apm.network.proto.CPU; import org.skywalking.apm.network.proto.Downstream; import org.skywalking.apm.network.proto.GC; @@ -139,6 +139,7 @@ public class JVMMetricsServiceHandler extends JVMMetricsServiceGrpc.JVMMetricsSe gcMetric.setCount(gc.getCount()); gcMetric.setTime(gc.getTime()); gcMetric.setTimeBucket(timeBucket); + gcMetric.setS5TimeBucket(TimeBucketUtils.INSTANCE.getFiveSecondTimeBucket(timeBucket)); try { logger.debug("send to gc metric persistence worker, id: {}", gcMetric.getId()); context.getClusterWorkerContext().lookup(GCMetricPersistenceWorker.WorkerRole.INSTANCE).tell(gcMetric.toData()); diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/dao/CpuMetricEsDAO.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/dao/CpuMetricEsDAO.java index 2812ccb80..6a00e2aa1 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/dao/CpuMetricEsDAO.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/dao/CpuMetricEsDAO.java @@ -4,17 +4,21 @@ import java.util.HashMap; import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.skywalking.apm.collector.core.stream.Data; +import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.jvm.CpuMetricTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; -import org.skywalking.apm.collector.core.stream.Data; -import org.skywalking.apm.collector.storage.define.DataDefine; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class CpuMetricEsDAO extends EsDAO implements ICpuMetricDAO, IPersistenceDAO { + private final Logger logger = LoggerFactory.getLogger(CpuMetricEsDAO.class); + @Override public Data get(String id, DataDefine dataDefine) { return null; } @@ -25,6 +29,7 @@ public class CpuMetricEsDAO extends EsDAO implements ICpuMetricDAO, IPersistence source.put(CpuMetricTable.COLUMN_USAGE_PERCENT, data.getDataDouble(0)); source.put(CpuMetricTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + logger.debug("prepare cpu metric batch insert, id: {}", data.getDataString(0)); return getClient().prepareIndex(CpuMetricTable.TABLE, data.getDataString(0)).setSource(source); } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/dao/GCMetricEsDAO.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/dao/GCMetricEsDAO.java index bf8af650c..6c645c0a6 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/dao/GCMetricEsDAO.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/dao/GCMetricEsDAO.java @@ -4,11 +4,11 @@ import java.util.HashMap; import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.skywalking.apm.collector.core.stream.Data; +import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.jvm.GCMetricTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; -import org.skywalking.apm.collector.core.stream.Data; -import org.skywalking.apm.collector.storage.define.DataDefine; /** * @author pengys5 @@ -22,10 +22,11 @@ public class GCMetricEsDAO extends EsDAO implements IGCMetricDAO, IPersistenceDA @Override public IndexRequestBuilder prepareBatchInsert(Data data) { Map source = new HashMap<>(); source.put(GCMetricTable.COLUMN_APPLICATION_INSTANCE_ID, data.getDataInteger(0)); - source.put(GCMetricTable.COLUMN_PHRASE, data.getDataInteger(0)); + source.put(GCMetricTable.COLUMN_PHRASE, data.getDataInteger(1)); source.put(GCMetricTable.COLUMN_COUNT, data.getDataLong(0)); source.put(GCMetricTable.COLUMN_TIME, data.getDataLong(1)); source.put(GCMetricTable.COLUMN_TIME_BUCKET, data.getDataLong(2)); + source.put(GCMetricTable.COLUMN_5S_TIME_BUCKET, data.getDataLong(3)); return getClient().prepareIndex(GCMetricTable.TABLE, data.getDataString(0)).setSource(source); } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricEsTableDefine.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricEsTableDefine.java index ea37e3f1b..58933d2e4 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricEsTableDefine.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricEsTableDefine.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentjvm.worker.gc.define; +import org.skywalking.apm.collector.storage.define.jvm.GCMetricTable; 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.jvm.GCMetricTable; /** * @author pengys5 @@ -31,5 +31,6 @@ public class GCMetricEsTableDefine extends ElasticSearchTableDefine { addColumn(new ElasticSearchColumnDefine(GCMetricTable.COLUMN_COUNT, ElasticSearchColumnDefine.Type.Long.name())); addColumn(new ElasticSearchColumnDefine(GCMetricTable.COLUMN_TIME, ElasticSearchColumnDefine.Type.Long.name())); addColumn(new ElasticSearchColumnDefine(GCMetricTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); + addColumn(new ElasticSearchColumnDefine(GCMetricTable.COLUMN_5S_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); } } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricH2TableDefine.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricH2TableDefine.java index e3744a21b..1d0859952 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricH2TableDefine.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricH2TableDefine.java @@ -20,5 +20,6 @@ public class GCMetricH2TableDefine extends H2TableDefine { addColumn(new H2ColumnDefine(GCMetricTable.COLUMN_COUNT, H2ColumnDefine.Type.Bigint.name())); addColumn(new H2ColumnDefine(GCMetricTable.COLUMN_TIME, H2ColumnDefine.Type.Bigint.name())); addColumn(new H2ColumnDefine(GCMetricTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); + addColumn(new H2ColumnDefine(GCMetricTable.COLUMN_5S_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); } } 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 3b2121e85..6183381fb 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 @@ -2,6 +2,8 @@ package org.skywalking.apm.collector.agentjvm.grpc.handler; import io.grpc.ManagedChannel; import io.grpc.ManagedChannelBuilder; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import org.skywalking.apm.network.proto.CPU; import org.skywalking.apm.network.proto.GC; import org.skywalking.apm.network.proto.GCPhrase; @@ -27,17 +29,16 @@ public class JVMMetricsServiceHandlerTestCase { ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 11800).usePlaintext(true).build(); stub = JVMMetricsServiceGrpc.newBlockingStub(channel); - buildJvmMetric(2); - buildJvmMetric(3); - - try { - Thread.sleep(2000); - } catch (InterruptedException e) { - e.printStackTrace(); - } + final long timeInterval = 1; + Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate(() -> multiInstanceJvmSend(), 1, timeInterval, TimeUnit.SECONDS); } - public static void buildJvmMetric(int instanceId) { + public static void multiInstanceJvmSend() { + buildJvmMetric(2); + buildJvmMetric(3); + } + + private static void buildJvmMetric(int instanceId) { JVMMetrics.Builder jvmMetricsBuilder = JVMMetrics.newBuilder(); jvmMetricsBuilder.setApplicationInstanceId(instanceId); @@ -88,10 +89,16 @@ public class JVMMetricsServiceHandlerTestCase { } private static void buildGcMetric(JVMMetric.Builder jvmMetric) { - GC.Builder gcBuilder = GC.newBuilder(); - gcBuilder.setPhrase(GCPhrase.NEW); - gcBuilder.setCount(2); - gcBuilder.setTime(100); - jvmMetric.addGc(gcBuilder.build()); + GC.Builder newGcBuilder = GC.newBuilder(); + newGcBuilder.setPhrase(GCPhrase.NEW); + newGcBuilder.setCount(2); + newGcBuilder.setTime(100); + jvmMetric.addGc(newGcBuilder.build()); + + GC.Builder oldGcBuilder = GC.newBuilder(); + oldGcBuilder.setPhrase(GCPhrase.OLD); + oldGcBuilder.setCount(2); + oldGcBuilder.setTime(100); + jvmMetric.addGc(oldGcBuilder.build()); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/GlobalTraceSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/GlobalTraceSpanListener.java index 29b6c9495..095a21ae1 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/GlobalTraceSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/GlobalTraceSpanListener.java @@ -5,7 +5,7 @@ import java.util.List; import org.skywalking.apm.collector.storage.define.global.GlobalTraceDataDefine; import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.GlobalTraceIdsListener; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/InstPerformanceSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/InstPerformanceSpanListener.java index 7c9051b18..110fd4238 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/InstPerformanceSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/InstPerformanceSpanListener.java @@ -1,15 +1,15 @@ package org.skywalking.apm.collector.agentstream.worker.instance.performance; -import org.skywalking.apm.collector.storage.define.instance.InstPerformanceDataDefine; import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; +import org.skywalking.apm.collector.storage.define.instance.InstPerformanceDataDefine; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; import org.skywalking.apm.network.proto.SpanObject; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -48,6 +48,7 @@ public class InstPerformanceSpanListener implements EntrySpanListener, FirstSpan instPerformance.setCallTimes(1); instPerformance.setCostTotal(cost); instPerformance.setTimeBucket(timeBucket); + instPerformance.setS5TimeBucket(TimeBucketUtils.INSTANCE.getFiveSecondTimeBucket(timeBucket)); try { logger.debug("send to instance performance persistence worker, id: {}", instPerformance.getId()); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/dao/InstPerformanceEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/dao/InstPerformanceEsDAO.java index 763ae1f1f..8ba89a0ca 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/dao/InstPerformanceEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/dao/InstPerformanceEsDAO.java @@ -5,27 +5,33 @@ import java.util.Map; import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; -import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; -import org.skywalking.apm.collector.storage.define.instance.InstPerformanceTable; -import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; +import org.skywalking.apm.collector.storage.define.instance.InstPerformanceTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class InstPerformanceEsDAO extends EsDAO implements IInstPerformanceDAO, IPersistenceDAO { + private final Logger logger = LoggerFactory.getLogger(InstPerformanceEsDAO.class); + @Override public Data get(String id, DataDefine dataDefine) { GetResponse getResponse = getClient().prepareGet(InstPerformanceTable.TABLE, id).get(); if (getResponse.isExists()) { + logger.debug("id: {} is exist", id); Data data = dataDefine.build(id); Map source = getResponse.getSource(); data.setDataInteger(0, (Integer)source.get(InstPerformanceTable.COLUMN_APPLICATION_ID)); data.setDataInteger(1, (Integer)source.get(InstPerformanceTable.COLUMN_INSTANCE_ID)); data.setDataInteger(2, (Integer)source.get(InstPerformanceTable.COLUMN_CALL_TIMES)); - data.setDataLong(0, (Long)source.get(InstPerformanceTable.COLUMN_COST_TOTAL)); - data.setDataLong(1, (Long)source.get(InstPerformanceTable.COLUMN_TIME_BUCKET)); + data.setDataLong(0, ((Number)source.get(InstPerformanceTable.COLUMN_COST_TOTAL)).longValue()); + data.setDataLong(1, ((Number)source.get(InstPerformanceTable.COLUMN_TIME_BUCKET)).longValue()); + data.setDataLong(2, ((Number)source.get(InstPerformanceTable.COLUMN_5S_TIME_BUCKET)).longValue()); return data; } else { return null; @@ -39,6 +45,7 @@ public class InstPerformanceEsDAO extends EsDAO implements IInstPerformanceDAO, source.put(InstPerformanceTable.COLUMN_CALL_TIMES, data.getDataInteger(2)); source.put(InstPerformanceTable.COLUMN_COST_TOTAL, data.getDataLong(0)); source.put(InstPerformanceTable.COLUMN_TIME_BUCKET, data.getDataLong(1)); + source.put(InstPerformanceTable.COLUMN_5S_TIME_BUCKET, data.getDataLong(2)); return getClient().prepareIndex(InstPerformanceTable.TABLE, data.getDataString(0)).setSource(source); } @@ -50,6 +57,7 @@ public class InstPerformanceEsDAO extends EsDAO implements IInstPerformanceDAO, source.put(InstPerformanceTable.COLUMN_CALL_TIMES, data.getDataInteger(2)); source.put(InstPerformanceTable.COLUMN_COST_TOTAL, data.getDataLong(0)); source.put(InstPerformanceTable.COLUMN_TIME_BUCKET, data.getDataLong(1)); + source.put(InstPerformanceTable.COLUMN_5S_TIME_BUCKET, data.getDataLong(2)); return getClient().prepareUpdate(InstPerformanceTable.TABLE, data.getDataString(0)).setDoc(source); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java index 123f63053..32a51a46d 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.instance.performance.define; +import org.skywalking.apm.collector.storage.define.instance.InstPerformanceTable; 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.instance.InstPerformanceTable; /** * @author pengys5 @@ -31,5 +31,6 @@ public class InstPerformanceEsTableDefine extends ElasticSearchTableDefine { addColumn(new ElasticSearchColumnDefine(InstPerformanceTable.COLUMN_CALL_TIMES, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(InstPerformanceTable.COLUMN_COST_TOTAL, ElasticSearchColumnDefine.Type.Long.name())); addColumn(new ElasticSearchColumnDefine(InstPerformanceTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); + addColumn(new ElasticSearchColumnDefine(InstPerformanceTable.COLUMN_5S_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceH2TableDefine.java index 65f22af4a..c0e9bd2d7 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceH2TableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceH2TableDefine.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.instance.performance.define; +import org.skywalking.apm.collector.storage.define.instance.InstPerformanceTable; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; -import org.skywalking.apm.collector.storage.define.instance.InstPerformanceTable; /** * @author pengys5 @@ -20,5 +20,6 @@ public class InstPerformanceH2TableDefine extends H2TableDefine { addColumn(new H2ColumnDefine(InstPerformanceTable.COLUMN_CALL_TIMES, H2ColumnDefine.Type.Int.name())); addColumn(new H2ColumnDefine(InstPerformanceTable.COLUMN_COST_TOTAL, H2ColumnDefine.Type.Bigint.name())); addColumn(new H2ColumnDefine(InstPerformanceTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); + addColumn(new H2ColumnDefine(InstPerformanceTable.COLUMN_5S_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java index 03d5bd6ff..7866954f4 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java @@ -9,7 +9,7 @@ import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.LocalSpanListener; import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java index 54853ff34..608c0d267 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java @@ -7,7 +7,7 @@ import org.skywalking.apm.collector.storage.define.node.NodeMappingDataDefine; import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener; import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java index b97994140..31f6ca23d 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java @@ -10,7 +10,7 @@ import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener; import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java index ea91b9db8..4742fbbed 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java @@ -10,7 +10,7 @@ import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener; import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostSpanListener.java index 321d9cce3..4aecaa94e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostSpanListener.java @@ -8,7 +8,7 @@ import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener import org.skywalking.apm.collector.agentstream.worker.segment.LocalSpanListener; import org.skywalking.apm.collector.storage.define.segment.SegmentCostDataDefine; import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntrySpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntrySpanListener.java index d00328b82..b37d32e88 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntrySpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntrySpanListener.java @@ -6,7 +6,7 @@ import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener; import org.skywalking.apm.collector.storage.define.service.ServiceEntryDataDefine; import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/ServiceRefSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/ServiceRefSpanListener.java index cfd1dcbaf..fa7029761 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/ServiceRefSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/ServiceRefSpanListener.java @@ -10,7 +10,7 @@ import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener; import org.skywalking.apm.collector.storage.define.serviceref.ServiceRefDataDefine; import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; 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 8db43b593..506b0d0ee 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 @@ -1,14 +1,16 @@ package org.skywalking.apm.collector.agentstream.mock; +import com.google.gson.JsonArray; import com.google.gson.JsonElement; +import com.google.gson.JsonObject; import java.io.IOException; import org.skywalking.apm.collector.agentstream.HttpClientTools; -import org.skywalking.apm.collector.storage.define.register.ApplicationDataDefine; import org.skywalking.apm.collector.agentstream.worker.register.application.dao.ApplicationEsDAO; -import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; import org.skywalking.apm.collector.agentstream.worker.register.instance.dao.InstanceEsDAO; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; import org.skywalking.apm.collector.core.CollectorException; +import org.skywalking.apm.collector.storage.define.register.ApplicationDataDefine; +import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; /** * @author pengys5 @@ -35,14 +37,31 @@ public class SegmentPost { ApplicationDataDefine.Application providerApplication = new ApplicationDataDefine.Application("3", "dubbox-provider", 3); applicationEsDAO.save(providerApplication); - JsonElement consumer = JsonFileReader.INSTANCE.read("json/segment/normal/dubbox-consumer.json"); - HttpClientTools.INSTANCE.post("http://localhost:12800/segments", consumer.toString()); + while (true) { + JsonElement consumer = JsonFileReader.INSTANCE.read("json/segment/normal/dubbox-consumer.json"); + modifyTime(consumer); + HttpClientTools.INSTANCE.post("http://localhost:12800/segments", consumer.toString()); - Thread.sleep(5000); + JsonElement provider = JsonFileReader.INSTANCE.read("json/segment/normal/dubbox-provider.json"); + modifyTime(provider); + HttpClientTools.INSTANCE.post("http://localhost:12800/segments", provider.toString()); - JsonElement provider = JsonFileReader.INSTANCE.read("json/segment/normal/dubbox-provider.json"); - HttpClientTools.INSTANCE.post("http://localhost:12800/segments", provider.toString()); + Thread.sleep(200); + } + } - Thread.sleep(5000); + private static void modifyTime(JsonElement jsonElement) { + JsonArray segmentArray = jsonElement.getAsJsonArray(); + for (JsonElement element : segmentArray) { + JsonObject segmentObj = element.getAsJsonObject(); + JsonArray spans = segmentObj.get("sg").getAsJsonObject().get("ss").getAsJsonArray(); + 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)); + } + } } } diff --git a/apm-collector/apm-collector-boot/src/main/resources/logback.xml b/apm-collector/apm-collector-boot/src/main/resources/logback.xml new file mode 100644 index 000000000..a0ca18fb1 --- /dev/null +++ b/apm-collector/apm-collector-boot/src/main/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-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageInstaller.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageInstaller.java index bb0a94511..01abbf72c 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageInstaller.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageInstaller.java @@ -20,13 +20,14 @@ public abstract class StorageInstaller { defineFilter(tableDefines); for (TableDefine tableDefine : tableDefines) { + tableDefine.initialize(); if (!isExists(client, tableDefine)) { logger.info("table: {} not exists", tableDefine.getName()); - tableDefine.initialize(); createTable(client, tableDefine); } else { logger.info("table: {} exists", tableDefine.getName()); -// deleteTable(client, tableDefine); + deleteTable(client, tableDefine); + createTable(client, tableDefine); } } } catch (DefineException e) { diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/stream/Data.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/stream/Data.java index 37fbc4eef..6789e26e2 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/stream/Data.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/stream/Data.java @@ -21,6 +21,7 @@ public class Data extends AbstractHashMessage { int booleanCapacity, int byteCapacity) { super(id); this.dataStrings = new String[stringCapacity]; + this.dataStrings[0] = id; this.dataLongs = new Long[longCapacity]; this.dataDoubles = new Double[doubleCapacity]; this.dataIntegers = new Integer[integerCapacity]; diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/util/TimeBucketUtils.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/TimeBucketUtils.java similarity index 82% rename from apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/util/TimeBucketUtils.java rename to apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/TimeBucketUtils.java index 99bd64033..7d2728ae5 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/util/TimeBucketUtils.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/TimeBucketUtils.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.stream.worker.util; +package org.skywalking.apm.collector.core.util; import java.text.SimpleDateFormat; import java.util.Calendar; @@ -43,6 +43,17 @@ public enum TimeBucketUtils { return Long.valueOf(timeStr); } + public long getFiveSecondTimeBucket(long secondTimeBucket) { + long mantissa = secondTimeBucket % 10; + if (mantissa < 5) { + return (secondTimeBucket / 10) * 10; + } else if (mantissa == 5) { + return secondTimeBucket; + } else { + return ((secondTimeBucket / 10) + 1) * 10; + } + } + public long changeToUTCTimeBucket(long timeBucket) { String timeBucketStr = String.valueOf(timeBucket); diff --git a/apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/utils/TimeBucketUtilsTestCase.java b/apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/utils/TimeBucketUtilsTestCase.java new file mode 100644 index 000000000..43e7b70dd --- /dev/null +++ b/apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/utils/TimeBucketUtilsTestCase.java @@ -0,0 +1,23 @@ +package org.skywalking.apm.collector.core.utils; + +import org.junit.Assert; +import org.junit.Test; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; + +/** + * @author pengys5 + */ +public class TimeBucketUtilsTestCase { + + @Test + public void testGetFiveSecondTimeBucket() { + long fiveSecondTimeBucket = TimeBucketUtils.INSTANCE.getFiveSecondTimeBucket(20170804224812L); + Assert.assertEquals(20170804224810L, fiveSecondTimeBucket); + + fiveSecondTimeBucket = TimeBucketUtils.INSTANCE.getFiveSecondTimeBucket(20170804224818L); + Assert.assertEquals(20170804224820L, fiveSecondTimeBucket); + + fiveSecondTimeBucket = TimeBucketUtils.INSTANCE.getFiveSecondTimeBucket(20170804224815L); + Assert.assertEquals(20170804224815L, fiveSecondTimeBucket); + } +} diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/CommonTable.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/CommonTable.java index eb2eb953e..05db7ae6e 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/CommonTable.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/CommonTable.java @@ -8,4 +8,5 @@ public class CommonTable { public static final String COLUMN_ID = "id"; public static final String COLUMN_AGG = "agg"; public static final String COLUMN_TIME_BUCKET = "time_bucket"; + public static final String COLUMN_5S_TIME_BUCKET = "s5_time_bucket"; } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/DataDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/DataDefine.java index 3dcdd0677..5c897a987 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/DataDefine.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/DataDefine.java @@ -61,22 +61,28 @@ public abstract class DataDefine { for (int i = 0; i < initialCapacity(); i++) { Attribute attribute = attributes[i]; if (AttributeType.STRING.equals(attribute.getType())) { - attribute.getOperation().operate(newData.getDataString(stringPosition), oldData.getDataString(stringPosition)); + String stringData = attribute.getOperation().operate(newData.getDataString(stringPosition), oldData.getDataString(stringPosition)); + newData.setDataString(stringPosition, stringData); stringPosition++; } else if (AttributeType.LONG.equals(attribute.getType())) { - attribute.getOperation().operate(newData.getDataLong(longPosition), oldData.getDataLong(longPosition)); + Long longData = attribute.getOperation().operate(newData.getDataLong(longPosition), oldData.getDataLong(longPosition)); + newData.setDataLong(longPosition, longData); longPosition++; } else if (AttributeType.DOUBLE.equals(attribute.getType())) { - attribute.getOperation().operate(newData.getDataDouble(doublePosition), oldData.getDataDouble(doublePosition)); + Double doubleData = attribute.getOperation().operate(newData.getDataDouble(doublePosition), oldData.getDataDouble(doublePosition)); + newData.setDataDouble(doublePosition, doubleData); doublePosition++; } else if (AttributeType.INTEGER.equals(attribute.getType())) { - attribute.getOperation().operate(newData.getDataInteger(integerPosition), oldData.getDataInteger(integerPosition)); + Integer integerData = attribute.getOperation().operate(newData.getDataInteger(integerPosition), oldData.getDataInteger(integerPosition)); + newData.setDataInteger(integerPosition, integerData); integerPosition++; } else if (AttributeType.BOOLEAN.equals(attribute.getType())) { - attribute.getOperation().operate(newData.getDataBoolean(booleanPosition), oldData.getDataBoolean(booleanPosition)); - integerPosition++; + Boolean booleanData = attribute.getOperation().operate(newData.getDataBoolean(booleanPosition), oldData.getDataBoolean(booleanPosition)); + newData.setDataBoolean(booleanPosition, booleanData); + booleanPosition++; } else if (AttributeType.BYTE.equals(attribute.getType())) { - attribute.getOperation().operate(newData.getDataBytes(bytePosition), oldData.getDataBytes(integerPosition)); + byte[] byteData = attribute.getOperation().operate(newData.getDataBytes(bytePosition), oldData.getDataBytes(integerPosition)); + newData.setDataBytes(bytePosition, byteData); bytePosition++; } } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/instance/InstPerformanceDataDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/instance/InstPerformanceDataDefine.java index 37257e7ce..79f9201b6 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/instance/InstPerformanceDataDefine.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/instance/InstPerformanceDataDefine.java @@ -16,7 +16,7 @@ import org.skywalking.apm.collector.storage.define.DataDefine; public class InstPerformanceDataDefine extends DataDefine { @Override protected int initialCapacity() { - return 6; + return 7; } @Override protected void attributeDefine() { @@ -26,6 +26,7 @@ public class InstPerformanceDataDefine extends DataDefine { addAttribute(3, new Attribute(InstPerformanceTable.COLUMN_CALL_TIMES, AttributeType.INTEGER, new AddOperation())); addAttribute(4, new Attribute(InstPerformanceTable.COLUMN_COST_TOTAL, AttributeType.LONG, new AddOperation())); addAttribute(5, new Attribute(InstPerformanceTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); + addAttribute(6, new Attribute(InstPerformanceTable.COLUMN_5S_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); } @Override public Object deserialize(RemoteData remoteData) { @@ -43,15 +44,18 @@ public class InstPerformanceDataDefine extends DataDefine { private int callTimes; private long costTotal; private long timeBucket; + private long s5TimeBucket; public InstPerformance(String id, int applicationId, int instanceId, int callTimes, long costTotal, - long timeBucket) { + long timeBucket, + long s5TimeBucket) { this.id = id; this.applicationId = applicationId; this.instanceId = instanceId; this.callTimes = callTimes; this.costTotal = costTotal; this.timeBucket = timeBucket; + this.s5TimeBucket = s5TimeBucket; } public InstPerformance() { @@ -66,6 +70,7 @@ public class InstPerformanceDataDefine extends DataDefine { data.setDataInteger(2, this.callTimes); data.setDataLong(0, this.costTotal); data.setDataLong(1, this.timeBucket); + data.setDataLong(2, this.s5TimeBucket); return data; } @@ -76,6 +81,7 @@ public class InstPerformanceDataDefine extends DataDefine { this.callTimes = data.getDataInteger(2); this.costTotal = data.getDataLong(0); this.timeBucket = data.getDataLong(1); + this.s5TimeBucket = data.getDataLong(2); return this; } @@ -126,5 +132,13 @@ public class InstPerformanceDataDefine extends DataDefine { public void setApplicationId(int applicationId) { this.applicationId = applicationId; } + + public long getS5TimeBucket() { + return s5TimeBucket; + } + + public void setS5TimeBucket(long s5TimeBucket) { + this.s5TimeBucket = s5TimeBucket; + } } } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/jvm/GCMetricDataDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/jvm/GCMetricDataDefine.java index 44f068efa..5b30f3579 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/jvm/GCMetricDataDefine.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/jvm/GCMetricDataDefine.java @@ -1,14 +1,14 @@ package org.skywalking.apm.collector.storage.define.jvm; import org.skywalking.apm.collector.core.framework.UnexpectedException; -import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; -import org.skywalking.apm.collector.storage.define.Attribute; -import org.skywalking.apm.collector.storage.define.AttributeType; import org.skywalking.apm.collector.core.stream.Data; -import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.core.stream.Transform; import org.skywalking.apm.collector.core.stream.operate.CoverOperation; import org.skywalking.apm.collector.core.stream.operate.NonOperation; +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.storage.define.Attribute; +import org.skywalking.apm.collector.storage.define.AttributeType; +import org.skywalking.apm.collector.storage.define.DataDefine; /** * @author pengys5 @@ -16,7 +16,7 @@ import org.skywalking.apm.collector.core.stream.operate.NonOperation; public class GCMetricDataDefine extends DataDefine { @Override protected int initialCapacity() { - return 6; + return 7; } @Override protected void attributeDefine() { @@ -26,6 +26,7 @@ public class GCMetricDataDefine extends DataDefine { addAttribute(3, new Attribute(GCMetricTable.COLUMN_COUNT, AttributeType.LONG, new CoverOperation())); addAttribute(4, new Attribute(GCMetricTable.COLUMN_TIME, AttributeType.LONG, new CoverOperation())); addAttribute(5, new Attribute(GCMetricTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); + addAttribute(6, new Attribute(GCMetricTable.COLUMN_5S_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); } @Override public Object deserialize(RemoteData remoteData) { @@ -43,14 +44,17 @@ public class GCMetricDataDefine extends DataDefine { private long count; private long time; private long timeBucket; + private long s5TimeBucket; - public GCMetric(String id, int applicationInstanceId, int phrase, long count, long time, long timeBucket) { + public GCMetric(String id, int applicationInstanceId, int phrase, long count, long time, long timeBucket, + long s5TimeBucket) { this.id = id; this.applicationInstanceId = applicationInstanceId; this.phrase = phrase; this.count = count; this.time = time; this.timeBucket = timeBucket; + this.s5TimeBucket = s5TimeBucket; } public GCMetric() { @@ -65,6 +69,7 @@ public class GCMetricDataDefine extends DataDefine { data.setDataLong(0, this.count); data.setDataLong(1, this.time); data.setDataLong(2, this.timeBucket); + data.setDataLong(3, this.s5TimeBucket); return data; } @@ -75,6 +80,7 @@ public class GCMetricDataDefine extends DataDefine { this.count = data.getDataLong(0); this.time = data.getDataLong(1); this.timeBucket = data.getDataLong(2); + this.s5TimeBucket = data.getDataLong(3); return this; } @@ -125,5 +131,13 @@ public class GCMetricDataDefine extends DataDefine { public void setTime(long time) { this.time = time; } + + public long getS5TimeBucket() { + return s5TimeBucket; + } + + public void setS5TimeBucket(long s5TimeBucket) { + this.s5TimeBucket = s5TimeBucket; + } } } 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 8296d978c..875cba06b 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 @@ -4,6 +4,7 @@ import java.util.ArrayList; 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.stream.worker.AbstractLocalAsyncWorker; import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; @@ -11,7 +12,6 @@ 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.stream.worker.impl.dao.IPersistenceDAO; -import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.stream.worker.impl.data.DataCache; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -75,12 +75,24 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { Data dbData = persistenceDAO().get(id, getRole().dataDefine()); if (ObjectUtils.isNotEmpty(dbData)) { getRole().dataDefine().mergeData(data, dbData); - updateBatchCollection.add(persistenceDAO().prepareBatchUpdate(data)); + try { + updateBatchCollection.add(persistenceDAO().prepareBatchUpdate(data)); + } catch (Throwable t) { + logger.error(t.getMessage(), t); + } } else { - insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data)); + try { + insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data)); + } catch (Throwable t) { + logger.error(t.getMessage(), t); + } } } else { - insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data)); + try { + insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data)); + } catch (Throwable t) { + logger.error(t.getMessage(), t); + } } }); diff --git a/apm-collector/apm-collector-stream/src/test/java/org/skywalking/apm/collector/stream/worker/util/TimeBucketUtilsTestCase.java b/apm-collector/apm-collector-stream/src/test/java/org/skywalking/apm/collector/stream/worker/util/TimeBucketUtilsTestCase.java index 661e6a7d5..a53e2416c 100644 --- a/apm-collector/apm-collector-stream/src/test/java/org/skywalking/apm/collector/stream/worker/util/TimeBucketUtilsTestCase.java +++ b/apm-collector/apm-collector-stream/src/test/java/org/skywalking/apm/collector/stream/worker/util/TimeBucketUtilsTestCase.java @@ -4,7 +4,7 @@ import java.util.Calendar; import java.util.TimeZone; import org.junit.Assert; import org.junit.Test; -import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; /** * @author pengys5 diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/cache/ApplicationCache.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/cache/ApplicationCache.java index 1b40d82b0..602405d05 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/cache/ApplicationCache.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/cache/ApplicationCache.java @@ -24,4 +24,14 @@ public class ApplicationCache { return Const.EXCEPTION; } } + + public static String getForUI(int applicationId) { + String applicationCode = get(applicationId); + if (applicationCode.equals("Unknown")) { + IApplicationDAO dao = (IApplicationDAO)DAOContainer.INSTANCE.get(IApplicationDAO.class.getName()); + applicationCode = dao.getApplicationCode(applicationId); + CACHE.put(applicationId, applicationCode); + } + return applicationCode; + } } 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 708a2228e..0e6ce7a45 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 @@ -12,40 +12,51 @@ import org.elasticsearch.action.search.SearchType; import org.elasticsearch.index.query.BoolQueryBuilder; import org.elasticsearch.index.query.MatchQueryBuilder; import org.elasticsearch.index.query.QueryBuilders; -import org.elasticsearch.search.SearchHit; +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; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class GCMetricEsDAO extends EsDAO implements IGCMetricDAO { - @Override public GCCount getGCCount(long timestamp, int instanceId) { + private final Logger logger = LoggerFactory.getLogger(GCMetricEsDAO.class); + + @Override public GCCount getGCCount(long s5TimeBucket, int instanceId) { + logger.debug("get gc count, s5TimeBucket: {}, instanceId: {}", s5TimeBucket, instanceId); SearchRequestBuilder searchRequestBuilder = getClient().prepareSearch(GCMetricTable.TABLE); searchRequestBuilder.setTypes(GCMetricTable.TABLE_TYPE); searchRequestBuilder.setSearchType(SearchType.DFS_QUERY_THEN_FETCH); BoolQueryBuilder boolQuery = QueryBuilders.boolQuery(); MatchQueryBuilder matchApplicationId = QueryBuilders.matchQuery(GCMetricTable.COLUMN_APPLICATION_INSTANCE_ID, instanceId); - MatchQueryBuilder matchTimeBucket = QueryBuilders.matchQuery(GCMetricTable.COLUMN_TIME_BUCKET, timestamp); + MatchQueryBuilder matchTimeBucket = QueryBuilders.matchQuery(GCMetricTable.COLUMN_5S_TIME_BUCKET, s5TimeBucket); boolQuery.must().add(matchApplicationId); boolQuery.must().add(matchTimeBucket); searchRequestBuilder.setQuery(boolQuery); - searchRequestBuilder.setSize(100); + searchRequestBuilder.setSize(0); + searchRequestBuilder.addAggregation( + AggregationBuilders.terms(GCMetricTable.COLUMN_PHRASE).field(GCMetricTable.COLUMN_PHRASE) + .subAggregation(AggregationBuilders.sum(GCMetricTable.COLUMN_COUNT).field(GCMetricTable.COLUMN_COUNT))); SearchResponse searchResponse = searchRequestBuilder.execute().actionGet(); - SearchHit[] searchHits = searchResponse.getHits().getHits(); GCCount gcCount = new GCCount(); - for (SearchHit searchHit : searchHits) { - int phrase = (Integer)searchHit.getSource().get(GCMetricTable.COLUMN_PHRASE); - int count = (Integer)searchHit.getSource().get(GCMetricTable.COLUMN_COUNT); + Terms phraseAggregation = searchResponse.getAggregations().get(GCMetricTable.COLUMN_PHRASE); + for (Terms.Bucket phraseBucket : phraseAggregation.getBuckets()) { + int phrase = phraseBucket.getKeyAsNumber().intValue(); + Sum sumAggregation = phraseBucket.getAggregations().get(GCMetricTable.COLUMN_COUNT); + int count = (int)sumAggregation.getValue(); if (phrase == GCPhrase.NEW_VALUE) { gcCount.setYoung(count); diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricH2DAO.java index 42f0d0924..e8f39776a 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricH2DAO.java @@ -1,5 +1,6 @@ package org.skywalking.apm.collector.ui.dao; +import com.google.gson.JsonObject; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; /** @@ -10,4 +11,12 @@ public class GCMetricH2DAO extends H2DAO implements IGCMetricDAO { @Override public GCCount getGCCount(long timestamp, int instanceId) { return null; } + + @Override public JsonObject getMetric(int instanceId, long timeBucket) { + return null; + } + + @Override public JsonObject getMetric(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/IInstPerformanceDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IInstPerformanceDAO.java index 672d6964a..488df87b0 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 @@ -7,7 +7,7 @@ import java.util.List; * @author pengys5 */ public interface IInstPerformanceDAO { - List getMultiple(long timestamp, int applicationId); + List getMultiple(long timeBucket, int applicationId); int getMetric(int instanceId, long timeBucket); 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 f35f8b5d3..f307512a3 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 @@ -13,7 +13,9 @@ import org.elasticsearch.action.search.SearchType; import org.elasticsearch.index.query.BoolQueryBuilder; import org.elasticsearch.index.query.MatchQueryBuilder; import org.elasticsearch.index.query.QueryBuilders; -import org.elasticsearch.search.SearchHit; +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.instance.InstPerformanceTable; import org.skywalking.apm.collector.storage.define.jvm.CpuMetricTable; @@ -24,30 +26,43 @@ import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; */ public class InstPerformanceEsDAO extends EsDAO implements IInstPerformanceDAO { - @Override public List getMultiple(long timestamp, int applicationId) { + @Override public List getMultiple(long timeBucket, int applicationId) { SearchRequestBuilder searchRequestBuilder = getClient().prepareSearch(InstPerformanceTable.TABLE); searchRequestBuilder.setTypes(InstPerformanceTable.TABLE_TYPE); searchRequestBuilder.setSearchType(SearchType.DFS_QUERY_THEN_FETCH); BoolQueryBuilder boolQuery = QueryBuilders.boolQuery(); MatchQueryBuilder matchApplicationId = QueryBuilders.matchQuery(InstPerformanceTable.COLUMN_APPLICATION_ID, applicationId); - MatchQueryBuilder matchTimeBucket = QueryBuilders.matchQuery(InstPerformanceTable.COLUMN_TIME_BUCKET, timestamp); + MatchQueryBuilder matchTimeBucket = QueryBuilders.matchQuery(InstPerformanceTable.COLUMN_5S_TIME_BUCKET, timeBucket); boolQuery.must().add(matchApplicationId); boolQuery.must().add(matchTimeBucket); searchRequestBuilder.setQuery(boolQuery); - searchRequestBuilder.setSize(100); + searchRequestBuilder.setSize(0); + + searchRequestBuilder.addAggregation( + AggregationBuilders.terms(InstPerformanceTable.COLUMN_INSTANCE_ID).field(InstPerformanceTable.COLUMN_INSTANCE_ID) + .subAggregation( + AggregationBuilders.terms(InstPerformanceTable.COLUMN_5S_TIME_BUCKET).field(InstPerformanceTable.COLUMN_5S_TIME_BUCKET) + .subAggregation(AggregationBuilders.sum(InstPerformanceTable.COLUMN_CALL_TIMES).field(InstPerformanceTable.COLUMN_CALL_TIMES)) + .subAggregation(AggregationBuilders.sum(InstPerformanceTable.COLUMN_COST_TOTAL).field(InstPerformanceTable.COLUMN_COST_TOTAL)))); SearchResponse searchResponse = searchRequestBuilder.execute().actionGet(); - List instPerformances = new LinkedList<>(); - SearchHit[] searchHits = searchResponse.getHits().getHits(); - for (SearchHit searchHit : searchHits) { - int instanceId = (Integer)searchHit.getSource().get(InstPerformanceTable.COLUMN_INSTANCE_ID); - int callTimes = (Integer)searchHit.getSource().get(InstPerformanceTable.COLUMN_CALL_TIMES); - long costTotal = ((Number)searchHit.getSource().get(InstPerformanceTable.COLUMN_COST_TOTAL)).longValue(); - instPerformances.add(new InstPerformance(instanceId, callTimes, costTotal)); + Terms instanceTerms = searchResponse.getAggregations().get(InstPerformanceTable.COLUMN_INSTANCE_ID); + List instPerformances = new LinkedList<>(); + for (Terms.Bucket instanceBucket : instanceTerms.getBuckets()) { + int instanceId = instanceBucket.getKeyAsNumber().intValue(); + Terms timeBucketTerms = instanceBucket.getAggregations().get(InstPerformanceTable.COLUMN_5S_TIME_BUCKET); + for (Terms.Bucket timeBucketBucket : timeBucketTerms.getBuckets()) { + long count = timeBucketBucket.getDocCount(); + Sum sumCallTimes = timeBucketBucket.getAggregations().get(InstPerformanceTable.COLUMN_CALL_TIMES); + Sum sumCostTotal = timeBucketBucket.getAggregations().get(InstPerformanceTable.COLUMN_COST_TOTAL); + int avgCallTimes = (int)(sumCallTimes.getValue() / count); + int avgCost = (int)(sumCostTotal.getValue() / count); + instPerformances.add(new InstPerformance(instanceId, avgCallTimes, avgCost)); + } } return instPerformances; 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 311f4d5fc..a480c7ce3 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 @@ -1,5 +1,6 @@ package org.skywalking.apm.collector.ui.dao; +import com.google.gson.JsonArray; import java.util.List; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; @@ -11,4 +12,12 @@ public class InstPerformanceH2DAO extends H2DAO implements IInstPerformanceDAO { @Override public List getMultiple(long timestamp, int applicationId) { return null; } + + @Override public int getMetric(int instanceId, long timeBucket) { + return 0; + } + + @Override public JsonArray getMetric(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/InstanceH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java index 0ba423247..2a7413de0 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.ui.dao; import java.util.List; +import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; -import org.skywalking.apm.collector.storage.table.instance.Instance; /** * @author pengys5 @@ -20,7 +20,7 @@ public class InstanceH2DAO extends H2DAO implements IInstanceDAO { return null; } - @Override public Instance getInstance(int instanceId) { + @Override public InstanceDataDefine.Instance getInstance(int instanceId) { return null; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancehealth/ApplicationsGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancehealth/ApplicationsGetHandler.java index be2de6f68..ba50ccccf 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancehealth/ApplicationsGetHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancehealth/ApplicationsGetHandler.java @@ -22,17 +22,17 @@ public class ApplicationsGetHandler extends JettyHandler { private InstanceHealthService service = new InstanceHealthService(); @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { - String timestamp = req.getParameter("timestamp"); - logger.debug("instance health applications get time: {}", timestamp); + String timestampStr = req.getParameter("timestamp"); + logger.debug("instance health applications get time: {}", timestampStr); - long time; + long timestamp; try { - time = Long.parseLong(timestamp); + timestamp = Long.parseLong(timestampStr); } catch (NumberFormatException e) { throw new ArgumentsParseException("time must be long"); } - return service.getApplications(time); + return service.getApplications(timestamp); } @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 107730e13..57832ba20 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 @@ -1,6 +1,8 @@ package org.skywalking.apm.collector.ui.jetty.handler.instancehealth; +import com.google.gson.JsonArray; import com.google.gson.JsonElement; +import com.google.gson.JsonObject; import javax.servlet.http.HttpServletRequest; import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; import org.skywalking.apm.collector.server.jetty.JettyHandler; @@ -23,8 +25,8 @@ public class InstanceHealthGetHandler extends JettyHandler { @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { String timestampStr = req.getParameter("timestamp"); - String applicationIdStr = req.getParameter("applicationId"); - logger.debug("instance health get timestamp: {}", timestampStr); + String[] applicationIdsStr = req.getParameterValues("applicationIds"); + logger.debug("instance health get timestamp: {}, applicationIdsStr: {}", timestampStr, applicationIdsStr); long timestamp; try { @@ -33,14 +35,24 @@ public class InstanceHealthGetHandler extends JettyHandler { throw new ArgumentsParseException("timestamp must be long"); } - int applicationId; - try { - applicationId = Integer.parseInt(applicationIdStr); - } catch (NumberFormatException e) { - throw new ArgumentsParseException("application id must be integer"); + int[] applicationIds = new int[applicationIdsStr.length]; + for (int i = 0; i < applicationIdsStr.length; i++) { + try { + applicationIds[i] = Integer.parseInt(applicationIdsStr[i]); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("application id must be integer"); + } } - return service.getInstances(timestamp, applicationId); + JsonObject response = new JsonObject(); + response.addProperty("timestamp", timestamp); + JsonArray appInstances = new JsonArray(); + response.add("appInstances", appInstances); + + for (int applicationId : applicationIds) { + appInstances.add(service.getInstances(timestamp, applicationId)); + } + return response; } @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/time/AllInstanceLastTimeGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/time/AllInstanceLastTimeGetHandler.java index df09f1327..96b56b81c 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/time/AllInstanceLastTimeGetHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/time/AllInstanceLastTimeGetHandler.java @@ -2,6 +2,7 @@ package org.skywalking.apm.collector.ui.jetty.handler.time; import com.google.gson.JsonElement; import com.google.gson.JsonObject; +import java.util.Calendar; import javax.servlet.http.HttpServletRequest; import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; import org.skywalking.apm.collector.server.jetty.JettyHandler; @@ -25,8 +26,15 @@ public class AllInstanceLastTimeGetHandler extends JettyHandler { @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { Long time = service.allInstanceLastTime(); logger.debug("all instance last time: {}", time); + Calendar calendar = Calendar.getInstance(); + calendar.setTimeInMillis(time); + + int second = calendar.get(Calendar.SECOND); + second = (second % 5) * 5; + calendar.set(Calendar.SECOND, second); + JsonObject timeJson = new JsonObject(); - timeJson.addProperty("time", time); + timeJson.addProperty("time", calendar.getTimeInMillis()); 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 66ac5c3bd..73321386a 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 @@ -3,6 +3,7 @@ package org.skywalking.apm.collector.ui.service; import com.google.gson.JsonArray; import com.google.gson.JsonObject; import java.util.List; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.DAOContainer; import org.skywalking.apm.collector.ui.cache.ApplicationCache; import org.skywalking.apm.collector.ui.dao.IGCMetricDAO; @@ -18,14 +19,14 @@ public class InstanceHealthService { private final Logger logger = LoggerFactory.getLogger(InstanceHealthService.class); - public JsonObject getApplications(long timestamp) { + public JsonObject getApplications(long timeBucket) { IInstanceDAO instanceDAO = (IInstanceDAO)DAOContainer.INSTANCE.get(IInstanceDAO.class.getName()); - List applications = instanceDAO.getApplications(timestamp); + List applications = instanceDAO.getApplications(timeBucket); JsonObject response = new JsonObject(); JsonArray applicationArray = new JsonArray(); - response.addProperty("timestamp", timestamp); + response.addProperty("timeBucket", timeBucket); response.add("applicationList", applicationArray); applications.forEach(application -> { @@ -43,15 +44,16 @@ public class InstanceHealthService { public JsonObject getInstances(long timestamp, int applicationId) { JsonObject response = new JsonObject(); - IInstPerformanceDAO instPerformanceDAO = (IInstPerformanceDAO)DAOContainer.INSTANCE.get(IInstPerformanceDAO.class.getName()); - List performances = instPerformanceDAO.getMultiple(timestamp, applicationId); + long secondTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(timestamp); + long s5TimeBucket = TimeBucketUtils.INSTANCE.getFiveSecondTimeBucket(secondTimeBucket); - response.addProperty("timestamp", timestamp); + IInstPerformanceDAO instPerformanceDAO = (IInstPerformanceDAO)DAOContainer.INSTANCE.get(IInstPerformanceDAO.class.getName()); + List performances = instPerformanceDAO.getMultiple(s5TimeBucket, applicationId); JsonArray instances = new JsonArray(); - response.addProperty("applicationCode", ApplicationCache.get(applicationId)); + response.addProperty("applicationCode", ApplicationCache.getForUI(applicationId)); response.addProperty("applicationId", applicationId); - response.add("appInstances", instances); + response.add("instances", instances); IGCMetricDAO gcMetricDAO = (IGCMetricDAO)DAOContainer.INSTANCE.get(IGCMetricDAO.class.getName()); performances.forEach(instance -> { @@ -59,14 +61,14 @@ public class InstanceHealthService { instanceJson.addProperty("id", instance.getInstanceId()); instanceJson.addProperty("tps", instance.getCallTimes()); - int avg = (int)((instance.getCostTotal() / instance.getCallTimes()) / 1000); + int avg = (int)(instance.getCostTotal() / instance.getCallTimes()); instanceJson.addProperty("avg", avg); - if (avg > 5) { + if (avg > 5000) { instanceJson.addProperty("healthLevel", 0); - } else if (avg > 3 && avg <= 5) { + } else if (avg > 3000 && avg <= 5000) { instanceJson.addProperty("healthLevel", 1); - } else if (avg > 1 && avg <= 3) { + } else if (avg > 1000 && avg <= 3000) { instanceJson.addProperty("healthLevel", 2); } else { instanceJson.addProperty("healthLevel", 3); @@ -74,7 +76,7 @@ public class InstanceHealthService { instanceJson.addProperty("status", 0); - IGCMetricDAO.GCCount gcCount = gcMetricDAO.getGCCount(timestamp, instance.getInstanceId()); + IGCMetricDAO.GCCount gcCount = gcMetricDAO.getGCCount(s5TimeBucket, instance.getInstanceId()); instanceJson.addProperty("ygc", gcCount.getYoung()); instanceJson.addProperty("ogc", gcCount.getOld());