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 846129b42..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 @@ -3,23 +3,23 @@ package org.skywalking.apm.collector.agentjvm.grpc.handler; import io.grpc.stub.StreamObserver; import java.util.List; import org.skywalking.apm.collector.agentjvm.worker.cpu.CpuMetricPersistenceWorker; -import org.skywalking.apm.collector.storage.define.jvm.CpuMetricDataDefine; import org.skywalking.apm.collector.agentjvm.worker.gc.GCMetricPersistenceWorker; -import org.skywalking.apm.collector.storage.define.jvm.GCMetricDataDefine; import org.skywalking.apm.collector.agentjvm.worker.heartbeat.InstHeartBeatPersistenceWorker; import org.skywalking.apm.collector.agentjvm.worker.heartbeat.define.InstanceHeartBeatDataDefine; import org.skywalking.apm.collector.agentjvm.worker.memory.MemoryMetricPersistenceWorker; -import org.skywalking.apm.collector.storage.define.jvm.MemoryMetricDataDefine; import org.skywalking.apm.collector.agentjvm.worker.memorypool.MemoryPoolMetricPersistenceWorker; -import org.skywalking.apm.collector.storage.define.jvm.MemoryPoolMetricDataDefine; 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; +import org.skywalking.apm.collector.storage.define.jvm.MemoryMetricDataDefine; +import org.skywalking.apm.collector.storage.define.jvm.MemoryPoolMetricDataDefine; 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; @@ -111,7 +111,7 @@ public class JVMMetricsServiceHandler extends JVMMetricsServiceGrpc.JVMMetricsSe memoryPools.forEach(memoryPool -> { MemoryPoolMetricDataDefine.MemoryPoolMetric memoryPoolMetric = new MemoryPoolMetricDataDefine.MemoryPoolMetric(); - memoryPoolMetric.setId(timeBucket + Const.ID_SPLIT + applicationInstanceId + Const.ID_SPLIT + String.valueOf(memoryPool.getType().getNumber())); + memoryPoolMetric.setId(timeBucket + Const.ID_SPLIT + applicationInstanceId + Const.ID_SPLIT + memoryPool.getIsHeap() + Const.ID_SPLIT + String.valueOf(memoryPool.getType().getNumber())); memoryPoolMetric.setApplicationInstanceId(applicationInstanceId); memoryPoolMetric.setPoolType(memoryPool.getType().getNumber()); memoryPoolMetric.setHeap(memoryPool.getIsHeap()); @@ -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-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java index c2362ff5e..b1f1284f8 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java @@ -10,6 +10,7 @@ import org.elasticsearch.action.admin.indices.delete.DeleteIndexResponse; import org.elasticsearch.action.admin.indices.exists.indices.IndicesExistsResponse; import org.elasticsearch.action.bulk.BulkRequestBuilder; import org.elasticsearch.action.get.GetRequestBuilder; +import org.elasticsearch.action.get.MultiGetRequestBuilder; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.search.SearchRequestBuilder; import org.elasticsearch.action.update.UpdateRequest; @@ -123,6 +124,10 @@ public class ElasticSearchClient implements Client { return client.prepareGet(indexName, "type", id); } + public MultiGetRequestBuilder prepareMultiGet() { + return client.prepareMultiGet(); + } + public BulkRequestBuilder prepareBulk() { return client.prepareBulk(); } 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/CpuMetricEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/CpuMetricEsDAO.java new file mode 100644 index 000000000..becd67d5a --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/CpuMetricEsDAO.java @@ -0,0 +1,49 @@ +package org.skywalking.apm.collector.ui.dao; + +import com.google.gson.JsonArray; +import org.elasticsearch.action.get.GetResponse; +import org.elasticsearch.action.get.MultiGetItemResponse; +import org.elasticsearch.action.get.MultiGetRequestBuilder; +import org.elasticsearch.action.get.MultiGetResponse; +import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.storage.define.jvm.CpuMetricTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; + +/** + * @author pengys5 + */ +public class CpuMetricEsDAO extends EsDAO implements ICpuMetricDAO { + + @Override public int getMetric(int instanceId, long timeBucket) { + String id = timeBucket + Const.ID_SPLIT + instanceId; + GetResponse getResponse = getClient().prepareGet(CpuMetricTable.TABLE, id).get(); + + if (getResponse.isExists()) { + return ((Number)getResponse.getSource().get(CpuMetricTable.COLUMN_USAGE_PERCENT)).intValue(); + } + return 0; + } + + @Override public JsonArray getMetric(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, CpuMetricTable.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()) { + metrics.add(((Number)response.getResponse().getSource().get(CpuMetricTable.COLUMN_USAGE_PERCENT)).intValue()); + } else { + metrics.add(0); + } + } + return metrics; + } +} 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 9f8a9081d..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 @@ -1,43 +1,62 @@ package org.skywalking.apm.collector.ui.dao; +import com.google.gson.JsonArray; +import com.google.gson.JsonObject; +import org.elasticsearch.action.get.GetResponse; +import org.elasticsearch.action.get.MultiGetItemResponse; +import org.elasticsearch.action.get.MultiGetRequestBuilder; +import org.elasticsearch.action.get.MultiGetResponse; import org.elasticsearch.action.search.SearchRequestBuilder; import org.elasticsearch.action.search.SearchResponse; 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.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +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); @@ -48,4 +67,69 @@ public class GCMetricEsDAO extends EsDAO implements IGCMetricDAO { return gcCount; } + + @Override public JsonObject getMetric(int instanceId, long timeBucket) { + JsonObject response = new JsonObject(); + + String youngId = timeBucket + Const.ID_SPLIT + GCPhrase.NEW_VALUE + instanceId; + GetResponse youngResponse = getClient().prepareGet(GCMetricTable.TABLE, youngId).get(); + if (youngResponse.isExists()) { + response.addProperty("ygc", ((Number)youngResponse.getSource().get(GCMetricTable.COLUMN_COUNT)).intValue()); + } + + String oldId = timeBucket + Const.ID_SPLIT + GCPhrase.OLD_VALUE + instanceId; + GetResponse oldResponse = getClient().prepareGet(GCMetricTable.TABLE, oldId).get(); + if (oldResponse.isExists()) { + response.addProperty("ogc", ((Number)oldResponse.getSource().get(GCMetricTable.COLUMN_COUNT)).intValue()); + } + + return response; + } + + @Override public JsonObject getMetric(int instanceId, long startTimeBucket, long endTimeBucket) { + JsonObject response = new JsonObject(); + + 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); + i++; + } + while (startTimeBucket + i <= endTimeBucket); + + JsonArray youngArray = new JsonArray(); + 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()); + } else { + youngArray.add(0); + } + } + response.add("ygc", youngArray); + + 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); + i++; + } + while (startTimeBucket + i <= endTimeBucket); + + JsonArray oldArray = new JsonArray(); + + multiGetResponse = oldPrepareMultiGet.get(); + for (MultiGetItemResponse itemResponse : multiGetResponse.getResponses()) { + if (itemResponse.getResponse().isExists()) { + oldArray.add(((Number)itemResponse.getResponse().getSource().get(CpuMetricTable.COLUMN_USAGE_PERCENT)).intValue()); + } else { + oldArray.add(0); + } + } + response.add("ogc", oldArray); + + return response; + } } 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/ICpuMetricDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ICpuMetricDAO.java new file mode 100644 index 000000000..b6671dada --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ICpuMetricDAO.java @@ -0,0 +1,12 @@ +package org.skywalking.apm.collector.ui.dao; + +import com.google.gson.JsonArray; + +/** + * @author pengys5 + */ +public interface ICpuMetricDAO { + int getMetric(int instanceId, long timeBucket); + + JsonArray getMetric(int instanceId, long startTimeBucket, long endTimeBucket); +} diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IGCMetricDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IGCMetricDAO.java index aa162e499..e9babdac2 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IGCMetricDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IGCMetricDAO.java @@ -1,5 +1,7 @@ package org.skywalking.apm.collector.ui.dao; +import com.google.gson.JsonObject; + /** * @author pengys5 */ @@ -7,6 +9,10 @@ public interface IGCMetricDAO { GCCount getGCCount(long timestamp, int instanceId); + JsonObject getMetric(int instanceId, long timeBucket); + + JsonObject getMetric(int instanceId, long startTimeBucket, long endTimeBucket); + class GCCount { private int young; private int old; 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 294d37f6f..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 @@ -1,12 +1,17 @@ package org.skywalking.apm.collector.ui.dao; +import com.google.gson.JsonArray; 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); + + JsonArray getMetric(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/IInstanceDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IInstanceDAO.java index 3d40c136b..c90ab98d9 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.define.register.InstanceDataDefine; /** * @author pengys5 @@ -12,6 +13,8 @@ public interface IInstanceDAO { List getApplications(long time); + InstanceDataDefine.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/IMemoryMetricDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IMemoryMetricDAO.java new file mode 100644 index 000000000..b3b299f95 --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IMemoryMetricDAO.java @@ -0,0 +1,12 @@ +package org.skywalking.apm.collector.ui.dao; + +import com.google.gson.JsonObject; + +/** + * @author pengys5 + */ +public interface IMemoryMetricDAO { + JsonObject getMetric(int instanceId, long timeBucket, boolean isHeap); + + JsonObject getMetric(int instanceId, long startTimeBucket, long endTimeBucket, boolean isHeap); +} diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IMemoryPoolMetricDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IMemoryPoolMetricDAO.java new file mode 100644 index 000000000..64525c9fc --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/IMemoryPoolMetricDAO.java @@ -0,0 +1,12 @@ +package org.skywalking.apm.collector.ui.dao; + +import com.google.gson.JsonObject; + +/** + * @author pengys5 + */ +public interface IMemoryPoolMetricDAO { + JsonObject getMetric(int instanceId, long timeBucket, boolean isHeap, int poolType); + + JsonObject getMetric(int instanceId, long startTimeBucket, long endTimeBucket, boolean isHeap, int poolType); +} 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 badc69f76..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 @@ -1,48 +1,103 @@ package org.skywalking.apm.collector.ui.dao; +import com.google.gson.JsonArray; import java.util.LinkedList; import java.util.List; +import org.elasticsearch.action.get.GetResponse; +import org.elasticsearch.action.get.MultiGetItemResponse; +import org.elasticsearch.action.get.MultiGetRequestBuilder; +import org.elasticsearch.action.get.MultiGetResponse; import org.elasticsearch.action.search.SearchRequestBuilder; import org.elasticsearch.action.search.SearchResponse; 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.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +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; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; /** * @author pengys5 */ 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; } + + @Override public int getMetric(int instanceId, long timeBucket) { + String id = timeBucket + Const.ID_SPLIT + instanceId; + GetResponse getResponse = getClient().prepareGet(InstPerformanceTable.TABLE, id).get(); + + if (getResponse.isExists()) { + return ((Number)getResponse.getSource().get(InstPerformanceTable.COLUMN_CALL_TIMES)).intValue(); + } + return 0; + } + + @Override public JsonArray getMetric(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); + i++; + } + while (startTimeBucket + i <= endTimeBucket); + + JsonArray metrics = new JsonArray(); + MultiGetResponse multiGetResponse = prepareMultiGet.get(); + for (MultiGetItemResponse response : multiGetResponse.getResponses()) { + if (response.getResponse().isExists()) { + metrics.add(((Number)response.getResponse().getSource().get(InstPerformanceTable.COLUMN_CALL_TIMES)).intValue()); + } 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 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/InstanceEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceEsDAO.java index 45faa26c3..6725798f9 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; @@ -15,8 +17,9 @@ import org.elasticsearch.search.aggregations.AggregationBuilders; 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.define.register.InstanceDataDefine; import org.skywalking.apm.collector.storage.define.register.InstanceTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -93,4 +96,21 @@ public class InstanceEsDAO extends EsDAO implements IInstanceDAO { } return applications; } + + @Override public InstanceDataDefine.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()) { + InstanceDataDefine.Instance instance = new InstanceDataDefine.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..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,6 +1,7 @@ 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; /** @@ -18,4 +19,8 @@ public class InstanceH2DAO extends H2DAO implements IInstanceDAO { @Override public List getApplications(long time) { return null; } + + @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/dao/MemoryMetricEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryMetricEsDAO.java new file mode 100644 index 000000000..3e481a587 --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryMetricEsDAO.java @@ -0,0 +1,62 @@ +package org.skywalking.apm.collector.ui.dao; + +import com.google.gson.JsonArray; +import com.google.gson.JsonObject; +import org.elasticsearch.action.get.GetResponse; +import org.elasticsearch.action.get.MultiGetItemResponse; +import org.elasticsearch.action.get.MultiGetRequestBuilder; +import org.elasticsearch.action.get.MultiGetResponse; +import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.storage.define.jvm.MemoryMetricTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; + +/** + * @author pengys5 + */ +public class MemoryMetricEsDAO extends EsDAO implements IMemoryMetricDAO { + + @Override public JsonObject getMetric(int instanceId, long timeBucket, boolean isHeap) { + String id = timeBucket + Const.ID_SPLIT + instanceId + Const.ID_SPLIT + isHeap; + GetResponse getResponse = getClient().prepareGet(MemoryMetricTable.TABLE, id).get(); + + JsonObject metric = new JsonObject(); + if (getResponse.isExists()) { + metric.addProperty("max", ((Number)getResponse.getSource().get(MemoryMetricTable.COLUMN_MAX)).intValue()); + metric.addProperty("init", ((Number)getResponse.getSource().get(MemoryMetricTable.COLUMN_INIT)).intValue()); + metric.addProperty("used", ((Number)getResponse.getSource().get(MemoryMetricTable.COLUMN_USED)).intValue()); + } else { + metric.addProperty("max", 0); + metric.addProperty("init", 0); + metric.addProperty("used", 0); + } + return metric; + } + + @Override public JsonObject getMetric(int instanceId, long startTimeBucket, long endTimeBucket, boolean isHeap) { + MultiGetRequestBuilder prepareMultiGet = getClient().prepareMultiGet(); + + int i = 0; + do { + String id = (startTimeBucket + i) + Const.ID_SPLIT + instanceId + Const.ID_SPLIT + isHeap; + prepareMultiGet.add(MemoryMetricTable.TABLE, MemoryMetricTable.TABLE_TYPE, id); + i++; + } + 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()) { + metric.addProperty("max", ((Number)response.getResponse().getSource().get(MemoryMetricTable.COLUMN_MAX)).intValue()); + metric.addProperty("init", ((Number)response.getResponse().getSource().get(MemoryMetricTable.COLUMN_INIT)).intValue()); + usedMetric.add(((Number)response.getResponse().getSource().get(MemoryMetricTable.COLUMN_USED)).intValue()); + } else { + usedMetric.add(0); + } + } + 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 new file mode 100644 index 000000000..a2cbbc10d --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/MemoryPoolMetricEsDAO.java @@ -0,0 +1,63 @@ +package org.skywalking.apm.collector.ui.dao; + +import com.google.gson.JsonArray; +import com.google.gson.JsonObject; +import org.elasticsearch.action.get.GetResponse; +import org.elasticsearch.action.get.MultiGetItemResponse; +import org.elasticsearch.action.get.MultiGetRequestBuilder; +import org.elasticsearch.action.get.MultiGetResponse; +import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.storage.define.jvm.MemoryPoolMetricTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; + +/** + * @author pengys5 + */ +public class MemoryPoolMetricEsDAO extends EsDAO implements IMemoryPoolMetricDAO { + + @Override public JsonObject getMetric(int instanceId, long timeBucket, boolean isHeap, int poolType) { + String id = timeBucket + Const.ID_SPLIT + instanceId + Const.ID_SPLIT + isHeap + Const.ID_SPLIT + poolType; + GetResponse getResponse = getClient().prepareGet(MemoryPoolMetricTable.TABLE, id).get(); + + JsonObject metric = new JsonObject(); + if (getResponse.isExists()) { + metric.addProperty("max", ((Number)getResponse.getSource().get(MemoryPoolMetricTable.COLUMN_MAX)).intValue()); + metric.addProperty("init", ((Number)getResponse.getSource().get(MemoryPoolMetricTable.COLUMN_INIT)).intValue()); + metric.addProperty("used", ((Number)getResponse.getSource().get(MemoryPoolMetricTable.COLUMN_USED)).intValue()); + } else { + metric.addProperty("max", 0); + metric.addProperty("init", 0); + metric.addProperty("used", 0); + } + return metric; + } + + @Override public JsonObject getMetric(int instanceId, long startTimeBucket, long endTimeBucket, boolean isHeap, + int poolType) { + MultiGetRequestBuilder prepareMultiGet = getClient().prepareMultiGet(); + + int i = 0; + do { + String id = (startTimeBucket + i) + Const.ID_SPLIT + instanceId + Const.ID_SPLIT + isHeap + Const.ID_SPLIT + poolType; + prepareMultiGet.add(MemoryPoolMetricTable.TABLE, MemoryPoolMetricTable.TABLE_TYPE, id); + i++; + } + 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()) { + metric.addProperty("max", ((Number)response.getResponse().getSource().get(MemoryPoolMetricTable.COLUMN_MAX)).intValue()); + metric.addProperty("init", ((Number)response.getResponse().getSource().get(MemoryPoolMetricTable.COLUMN_INIT)).intValue()); + usedMetric.add(((Number)response.getResponse().getSource().get(MemoryPoolMetricTable.COLUMN_USED)).intValue()); + } else { + usedMetric.add(0); + } + } + return metric; + } +} 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/instancemetric/InstanceMetricGetOneTimeBucketHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceMetricGetOneTimeBucketHandler.java new file mode 100644 index 000000000..3215a19a4 --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceMetricGetOneTimeBucketHandler.java @@ -0,0 +1,62 @@ +package org.skywalking.apm.collector.ui.jetty.handler.instancemetric; + +import com.google.gson.JsonElement; +import java.util.LinkedHashSet; +import java.util.Set; +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 InstanceMetricGetOneTimeBucketHandler extends JettyHandler { + + private final Logger logger = LoggerFactory.getLogger(InstanceMetricGetOneTimeBucketHandler.class); + + @Override public String pathSpec() { + return "/instance/jvm/instanceId/oneBucket"; + } + + private InstanceJVMService service = new InstanceJVMService(); + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + String timeBucketStr = req.getParameter("timeBucket"); + String instanceIdStr = req.getParameter("instanceId"); + String[] metricTypes = req.getParameterValues("metricTypes"); + + logger.debug("instance jvm metric get timeBucket: {}, instance id: {}, metric types: {}", timeBucketStr, instanceIdStr, metricTypes); + + long timeBucket; + try { + timeBucket = Long.parseLong(timeBucketStr); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("timeBucket must be long"); + } + + int instanceId; + try { + instanceId = Integer.parseInt(instanceIdStr); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("instance id must be integer"); + } + + if (metricTypes.length == 0) { + throw new ArgumentsParseException("at least one metric type"); + } + + Set metricTypeSet = new LinkedHashSet<>(); + for (String metricType : metricTypes) { + metricTypeSet.add(metricType); + } + + return service.getInstanceJvmMetric(instanceId, metricTypeSet, timeBucket); + } + + @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/InstanceMetricGetRangeTimeBucketHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceMetricGetRangeTimeBucketHandler.java new file mode 100644 index 000000000..f74be568e --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/instancemetric/InstanceMetricGetRangeTimeBucketHandler.java @@ -0,0 +1,70 @@ +package org.skywalking.apm.collector.ui.jetty.handler.instancemetric; + +import com.google.gson.JsonElement; +import java.util.LinkedHashSet; +import java.util.Set; +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 InstanceMetricGetRangeTimeBucketHandler extends JettyHandler { + + private final Logger logger = LoggerFactory.getLogger(InstanceMetricGetRangeTimeBucketHandler.class); + + @Override public String pathSpec() { + return "/instance/jvm/instanceId/rangeBucket"; + } + + private InstanceJVMService service = new InstanceJVMService(); + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + String startTimeBucketStr = req.getParameter("startTimeBucket"); + String endTimeBucketStr = req.getParameter("endTimeBucket"); + String instanceIdStr = req.getParameter("instanceId"); + String[] metricTypes = req.getParameterValues("metricTypes"); + + logger.debug("instance jvm metric get start timeBucket: {}, end timeBucket:{} , instance id: {}, metric types: {}", startTimeBucketStr, endTimeBucketStr, instanceIdStr, metricTypes); + + long startTimeBucket; + try { + startTimeBucket = Long.parseLong(startTimeBucketStr); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("start timeBucket must be long"); + } + + long endTimeBucket; + try { + endTimeBucket = Long.parseLong(endTimeBucketStr); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("end timeBucket must be long"); + } + + int instanceId; + try { + instanceId = Integer.parseInt(instanceIdStr); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("instance id must be integer"); + } + + if (metricTypes.length == 0) { + throw new ArgumentsParseException("at least one metric type"); + } + + Set metricTypeSet = new LinkedHashSet<>(); + for (String metricType : metricTypes) { + metricTypeSet.add(metricType); + } + + return service.getInstanceJvmMetrics(instanceId, metricTypeSet, startTimeBucket, endTimeBucket); + } + + @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/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()); 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..d84aee2dc --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceJVMService.java @@ -0,0 +1,152 @@ +package org.skywalking.apm.collector.ui.service; + +import com.google.gson.Gson; +import com.google.gson.JsonObject; +import java.util.Set; +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.define.register.InstanceDataDefine; +import org.skywalking.apm.collector.ui.dao.ICpuMetricDAO; +import org.skywalking.apm.collector.ui.dao.IGCMetricDAO; +import org.skywalking.apm.collector.ui.dao.IInstPerformanceDAO; +import org.skywalking.apm.collector.ui.dao.IInstanceDAO; +import org.skywalking.apm.collector.ui.dao.IMemoryMetricDAO; +import org.skywalking.apm.collector.ui.dao.IMemoryPoolMetricDAO; +import org.skywalking.apm.network.proto.PoolType; +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()); + InstanceDataDefine.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, 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"); + } + return metrics; + } + + 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"); + } + return metrics; + } + + public enum MetricType { + cpu, gc, tps, heapmemory, heappermgen, heapmetaspace, heapnewgen, + heapoldgen, heapsurvivor, nonheapmemory, nonheappermgen, nonheapmetaspace, + nonheapnewgen, nonheapoldgen, nonheapsurvivor + } +}