From 8361f4156b46e2da3da2bbd16f9210698b34215f Mon Sep 17 00:00:00 2001 From: peng-yongsheng <8082209@qq.com> Date: Tue, 9 Jan 2018 16:03:40 +0800 Subject: [PATCH] CPU metric pyramid aggregate test successful. --- .../handler/JVMMetricsServiceHandler.java | 2 +- .../JVMMetricServiceHandlerTestCase.java | 9 +++++ .../provider/service/CpuMetricService.java | 7 +++- .../worker/cpu/CpuDayMetricTransformNode.java | 7 ++-- .../cpu/CpuHourMetricTransformNode.java | 7 ++-- .../provider/worker/cpu/CpuMetricCopy.java | 40 +++++++++++++++++++ .../cpu/CpuMinuteMetricTransformNode.java | 7 ++-- .../cpu/CpuMonthMetricTransformNode.java | 7 ++-- .../apm/collector/storage/StorageModule.java | 8 ++++ .../storage/table/jvm/CpuMetric.java | 13 +++++- .../storage/table/jvm/CpuMetricTable.java | 2 +- .../storage/es/StorageModuleEsProvider.java | 12 ++++++ .../AbstractCpuMetricEsPersistenceDAO.java | 13 +++++- .../AbstractCpuMetricEsTableDefine.java} | 18 ++++----- .../define/cpu/CpuDayMetricEsTableDefine.java | 37 +++++++++++++++++ .../cpu/CpuHourMetricEsTableDefine.java | 37 +++++++++++++++++ .../cpu/CpuMinuteMetricEsTableDefine.java | 37 +++++++++++++++++ .../cpu/CpuMonthMetricEsTableDefine.java | 37 +++++++++++++++++ .../cpu/CpuSecondMetricEsTableDefine.java | 37 +++++++++++++++++ .../resources/META-INF/defines/storage.define | 8 +++- 20 files changed, 316 insertions(+), 29 deletions(-) create mode 100644 apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMetricCopy.java rename apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/{CpuMetricEsTableDefine.java => cpu/AbstractCpuMetricEsTableDefine.java} (69%) create mode 100644 apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuDayMetricEsTableDefine.java create mode 100644 apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuHourMetricEsTableDefine.java create mode 100644 apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuMinuteMetricEsTableDefine.java create mode 100644 apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuMonthMetricEsTableDefine.java create mode 100644 apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuSecondMetricEsTableDefine.java diff --git a/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/JVMMetricsServiceHandler.java b/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/JVMMetricsServiceHandler.java index e89788935..fa958ae70 100644 --- a/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/JVMMetricsServiceHandler.java +++ b/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/JVMMetricsServiceHandler.java @@ -67,7 +67,7 @@ public class JVMMetricsServiceHandler extends JVMMetricsServiceGrpc.JVMMetricsSe request.getMetricsList().forEach(metric -> { long time = TimeBucketUtils.INSTANCE.getSecondTimeBucket(metric.getTime()); // sendToInstanceHeartBeatService(instanceId, metric.getTime()); -// sendToCpuMetricService(instanceId, time, metric.getCpu()); + sendToCpuMetricService(instanceId, time, metric.getCpu()); sendToMemoryMetricService(instanceId, time, metric.getMemoryList()); sendToMemoryPoolMetricService(instanceId, time, metric.getMemoryPoolList()); sendToGCMetricService(instanceId, time, metric.getGcList()); diff --git a/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/JVMMetricServiceHandlerTestCase.java b/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/JVMMetricServiceHandlerTestCase.java index 079696574..7862fcaaf 100644 --- a/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/JVMMetricServiceHandlerTestCase.java +++ b/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/JVMMetricServiceHandlerTestCase.java @@ -20,6 +20,7 @@ package org.apache.skywalking.apm.collector.agent.grpc.provider.handler; import io.grpc.ManagedChannel; import io.grpc.ManagedChannelBuilder; +import org.apache.skywalking.apm.network.proto.CPU; import org.apache.skywalking.apm.network.proto.GC; import org.apache.skywalking.apm.network.proto.GCPhrase; import org.apache.skywalking.apm.network.proto.JVMMetric; @@ -44,6 +45,7 @@ public class JVMMetricServiceHandlerTestCase { JVMMetric.Builder metricBuilder = JVMMetric.newBuilder(); metricBuilder.setTime(System.currentTimeMillis()); + buildCPUMetric(metricBuilder); buildGCMetric(metricBuilder); buildMemoryMetric(metricBuilder); buildMemoryPoolMetric(metricBuilder); @@ -82,4 +84,11 @@ public class JVMMetricServiceHandlerTestCase { metricBuilder.addGc(builder); } + + private static void buildCPUMetric(JVMMetric.Builder metricBuilder) { + CPU.Builder builder = CPU.newBuilder(); + builder.setUsagePercent(20); + + metricBuilder.setCpu(builder.build()); + } } diff --git a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/service/CpuMetricService.java b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/service/CpuMetricService.java index 3bc966fdc..49b68e926 100644 --- a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/service/CpuMetricService.java +++ b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/service/CpuMetricService.java @@ -45,10 +45,15 @@ public class CpuMetricService implements ICpuMetricService { } @Override public void send(int instanceId, long timeBucket, double usagePercent) { + String metricId = String.valueOf(instanceId); + String id = timeBucket + Const.ID_SPLIT + metricId; + CpuMetric cpuMetric = new CpuMetric(); - cpuMetric.setId(timeBucket + Const.ID_SPLIT + instanceId); + cpuMetric.setId(id); + cpuMetric.setMetricId(metricId); cpuMetric.setInstanceId(instanceId); cpuMetric.setUsagePercent(usagePercent); + cpuMetric.setTimes(1L); cpuMetric.setTimeBucket(timeBucket); logger.debug("push to cpu metric graph, id: {}", cpuMetric.getId()); diff --git a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuDayMetricTransformNode.java b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuDayMetricTransformNode.java index 3128fc917..444e9c758 100644 --- a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuDayMetricTransformNode.java +++ b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuDayMetricTransformNode.java @@ -36,9 +36,10 @@ public class CpuDayMetricTransformNode implements NodeProcessor next) { long timeBucket = TimeBucketUtils.INSTANCE.secondToDay(cpuMetric.getTimeBucket()); - cpuMetric.setId(String.valueOf(timeBucket) + Const.ID_SPLIT + cpuMetric.getMetricId()); - cpuMetric.setTimeBucket(timeBucket); - next.execute(cpuMetric); + CpuMetric newCpuMetric = CpuMetricCopy.copy(cpuMetric); + newCpuMetric.setId(String.valueOf(timeBucket) + Const.ID_SPLIT + cpuMetric.getMetricId()); + newCpuMetric.setTimeBucket(timeBucket); + next.execute(newCpuMetric); } } diff --git a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuHourMetricTransformNode.java b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuHourMetricTransformNode.java index 48699714c..cbc98da8a 100644 --- a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuHourMetricTransformNode.java +++ b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuHourMetricTransformNode.java @@ -36,9 +36,10 @@ public class CpuHourMetricTransformNode implements NodeProcessor next) { long timeBucket = TimeBucketUtils.INSTANCE.secondToHour(cpuMetric.getTimeBucket()); - cpuMetric.setId(String.valueOf(timeBucket) + Const.ID_SPLIT + cpuMetric.getMetricId()); - cpuMetric.setTimeBucket(timeBucket); - next.execute(cpuMetric); + CpuMetric newCpuMetric = CpuMetricCopy.copy(cpuMetric); + newCpuMetric.setId(String.valueOf(timeBucket) + Const.ID_SPLIT + cpuMetric.getMetricId()); + newCpuMetric.setTimeBucket(timeBucket); + next.execute(newCpuMetric); } } diff --git a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMetricCopy.java b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMetricCopy.java new file mode 100644 index 000000000..067b07344 --- /dev/null +++ b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMetricCopy.java @@ -0,0 +1,40 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.analysis.jvm.provider.worker.cpu; + +import org.apache.skywalking.apm.collector.storage.table.jvm.CpuMetric; + +/** + * @author peng-yongsheng + */ +public class CpuMetricCopy { + + public static CpuMetric copy(CpuMetric cpuMetric) { + CpuMetric newCpuMetric = new CpuMetric(); + newCpuMetric.setId(cpuMetric.getId()); + newCpuMetric.setMetricId(cpuMetric.getMetricId()); + + newCpuMetric.setInstanceId(cpuMetric.getInstanceId()); + newCpuMetric.setUsagePercent(cpuMetric.getUsagePercent()); + newCpuMetric.setTimes(cpuMetric.getTimes()); + + newCpuMetric.setTimeBucket(cpuMetric.getTimeBucket()); + return newCpuMetric; + } +} diff --git a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMinuteMetricTransformNode.java b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMinuteMetricTransformNode.java index 6e276a9ee..e2d9f22f1 100644 --- a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMinuteMetricTransformNode.java +++ b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMinuteMetricTransformNode.java @@ -36,9 +36,10 @@ public class CpuMinuteMetricTransformNode implements NodeProcessor next) { long timeBucket = TimeBucketUtils.INSTANCE.secondToMinute(cpuMetric.getTimeBucket()); - cpuMetric.setId(String.valueOf(timeBucket) + Const.ID_SPLIT + cpuMetric.getMetricId()); - cpuMetric.setTimeBucket(timeBucket); - next.execute(cpuMetric); + CpuMetric newCpuMetric = CpuMetricCopy.copy(cpuMetric); + newCpuMetric.setId(String.valueOf(timeBucket) + Const.ID_SPLIT + cpuMetric.getMetricId()); + newCpuMetric.setTimeBucket(timeBucket); + next.execute(newCpuMetric); } } diff --git a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMonthMetricTransformNode.java b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMonthMetricTransformNode.java index 87ff41276..6c59f4756 100644 --- a/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMonthMetricTransformNode.java +++ b/apm-collector/apm-collector-analysis/analysis-jvm/jvm-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/jvm/provider/worker/cpu/CpuMonthMetricTransformNode.java @@ -36,9 +36,10 @@ public class CpuMonthMetricTransformNode implements NodeProcessor next) { long timeBucket = TimeBucketUtils.INSTANCE.secondToMonth(cpuMetric.getTimeBucket()); - cpuMetric.setId(String.valueOf(timeBucket) + Const.ID_SPLIT + cpuMetric.getMetricId()); - cpuMetric.setTimeBucket(timeBucket); - next.execute(cpuMetric); + CpuMetric newCpuMetric = CpuMetricCopy.copy(cpuMetric); + newCpuMetric.setId(String.valueOf(timeBucket) + Const.ID_SPLIT + cpuMetric.getMetricId()); + newCpuMetric.setTimeBucket(timeBucket); + next.execute(newCpuMetric); } } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/StorageModule.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/StorageModule.java index c607d8042..5ad1ae98e 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/StorageModule.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/StorageModule.java @@ -71,6 +71,10 @@ import org.apache.skywalking.apm.collector.storage.dao.cache.IApplicationCacheDA import org.apache.skywalking.apm.collector.storage.dao.cache.IInstanceCacheDAO; import org.apache.skywalking.apm.collector.storage.dao.cache.INetworkAddressCacheDAO; import org.apache.skywalking.apm.collector.storage.dao.cache.IServiceNameCacheDAO; +import org.apache.skywalking.apm.collector.storage.dao.cpump.ICpuDayMetricPersistenceDAO; +import org.apache.skywalking.apm.collector.storage.dao.cpump.ICpuHourMetricPersistenceDAO; +import org.apache.skywalking.apm.collector.storage.dao.cpump.ICpuMinuteMetricPersistenceDAO; +import org.apache.skywalking.apm.collector.storage.dao.cpump.ICpuMonthMetricPersistenceDAO; import org.apache.skywalking.apm.collector.storage.dao.cpump.ICpuSecondMetricPersistenceDAO; import org.apache.skywalking.apm.collector.storage.dao.gcmp.IGCDayMetricPersistenceDAO; import org.apache.skywalking.apm.collector.storage.dao.gcmp.IGCHourMetricPersistenceDAO; @@ -152,6 +156,10 @@ public class StorageModule extends Module { private void addPersistenceDAO(List classes) { classes.add(ICpuSecondMetricPersistenceDAO.class); + classes.add(ICpuMinuteMetricPersistenceDAO.class); + classes.add(ICpuHourMetricPersistenceDAO.class); + classes.add(ICpuDayMetricPersistenceDAO.class); + classes.add(ICpuMonthMetricPersistenceDAO.class); classes.add(IGCSecondMetricPersistenceDAO.class); classes.add(IGCMinuteMetricPersistenceDAO.class); diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetric.java index 73731f85f..67ecd6ecc 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetric.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetric.java @@ -35,6 +35,7 @@ public class CpuMetric extends StreamData { }; private static final Column[] LONG_COLUMNS = { + new Column(CpuMetricTable.COLUMN_TIMES, new AddOperation()), new Column(CpuMetricTable.COLUMN_TIME_BUCKET, new CoverOperation()), }; @@ -85,11 +86,19 @@ public class CpuMetric extends StreamData { setDataDouble(0, usagePercent); } - public Long getTimeBucket() { + public Long getTimes() { return getDataLong(0); } + public void setTimes(Long times) { + setDataLong(0, times); + } + + public Long getTimeBucket() { + return getDataLong(1); + } + public void setTimeBucket(Long timeBucket) { - setDataLong(0, timeBucket); + setDataLong(1, timeBucket); } } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetricTable.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetricTable.java index 31e83c564..1f4ef5798 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetricTable.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetricTable.java @@ -16,7 +16,6 @@ * */ - package org.apache.skywalking.apm.collector.storage.table.jvm; import org.apache.skywalking.apm.collector.core.data.CommonTable; @@ -28,4 +27,5 @@ public class CpuMetricTable extends CommonTable { public static final String TABLE = "cpu_metric"; public static final String COLUMN_INSTANCE_ID = "instance_id"; public static final String COLUMN_USAGE_PERCENT = "usage_percent"; + public static final String COLUMN_TIMES = "times"; } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/StorageModuleEsProvider.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/StorageModuleEsProvider.java index d4225309e..4ada2c763 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/StorageModuleEsProvider.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/StorageModuleEsProvider.java @@ -80,6 +80,10 @@ import org.apache.skywalking.apm.collector.storage.dao.cache.IApplicationCacheDA import org.apache.skywalking.apm.collector.storage.dao.cache.IInstanceCacheDAO; import org.apache.skywalking.apm.collector.storage.dao.cache.INetworkAddressCacheDAO; import org.apache.skywalking.apm.collector.storage.dao.cache.IServiceNameCacheDAO; +import org.apache.skywalking.apm.collector.storage.dao.cpump.ICpuDayMetricPersistenceDAO; +import org.apache.skywalking.apm.collector.storage.dao.cpump.ICpuHourMetricPersistenceDAO; +import org.apache.skywalking.apm.collector.storage.dao.cpump.ICpuMinuteMetricPersistenceDAO; +import org.apache.skywalking.apm.collector.storage.dao.cpump.ICpuMonthMetricPersistenceDAO; import org.apache.skywalking.apm.collector.storage.dao.cpump.ICpuSecondMetricPersistenceDAO; import org.apache.skywalking.apm.collector.storage.dao.gcmp.IGCDayMetricPersistenceDAO; import org.apache.skywalking.apm.collector.storage.dao.gcmp.IGCHourMetricPersistenceDAO; @@ -171,6 +175,10 @@ import org.apache.skywalking.apm.collector.storage.es.dao.cache.ApplicationEsCac import org.apache.skywalking.apm.collector.storage.es.dao.cache.InstanceEsCacheDAO; import org.apache.skywalking.apm.collector.storage.es.dao.cache.NetworkAddressEsCacheDAO; import org.apache.skywalking.apm.collector.storage.es.dao.cache.ServiceNameEsCacheDAO; +import org.apache.skywalking.apm.collector.storage.es.dao.cpump.CpuDayMetricEsPersistenceDAO; +import org.apache.skywalking.apm.collector.storage.es.dao.cpump.CpuHourMetricEsPersistenceDAO; +import org.apache.skywalking.apm.collector.storage.es.dao.cpump.CpuMinuteMetricEsPersistenceDAO; +import org.apache.skywalking.apm.collector.storage.es.dao.cpump.CpuMonthMetricEsPersistenceDAO; import org.apache.skywalking.apm.collector.storage.es.dao.cpump.CpuSecondMetricEsPersistenceDAO; import org.apache.skywalking.apm.collector.storage.es.dao.gcmp.GCDayMetricEsPersistenceDAO; import org.apache.skywalking.apm.collector.storage.es.dao.gcmp.GCHourMetricEsPersistenceDAO; @@ -302,6 +310,10 @@ public class StorageModuleEsProvider extends ModuleProvider { private void registerPersistenceDAO() throws ServiceNotProvidedException { this.registerServiceImplementation(ICpuSecondMetricPersistenceDAO.class, new CpuSecondMetricEsPersistenceDAO(elasticSearchClient)); + this.registerServiceImplementation(ICpuMinuteMetricPersistenceDAO.class, new CpuMinuteMetricEsPersistenceDAO(elasticSearchClient)); + this.registerServiceImplementation(ICpuHourMetricPersistenceDAO.class, new CpuHourMetricEsPersistenceDAO(elasticSearchClient)); + this.registerServiceImplementation(ICpuDayMetricPersistenceDAO.class, new CpuDayMetricEsPersistenceDAO(elasticSearchClient)); + this.registerServiceImplementation(ICpuMonthMetricPersistenceDAO.class, new CpuMonthMetricEsPersistenceDAO(elasticSearchClient)); this.registerServiceImplementation(IGCSecondMetricPersistenceDAO.class, new GCSecondMetricEsPersistenceDAO(elasticSearchClient)); this.registerServiceImplementation(IGCMinuteMetricPersistenceDAO.class, new GCMinuteMetricEsPersistenceDAO(elasticSearchClient)); diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/cpump/AbstractCpuMetricEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/cpump/AbstractCpuMetricEsPersistenceDAO.java index 824308e38..33ce89121 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/cpump/AbstractCpuMetricEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/cpump/AbstractCpuMetricEsPersistenceDAO.java @@ -39,7 +39,17 @@ public abstract class AbstractCpuMetricEsPersistenceDAO extends AbstractPersiste } @Override protected final CpuMetric esDataToStreamData(Map source) { - return null; + CpuMetric cpuMetric = new CpuMetric(); + cpuMetric.setId((String)source.get(CpuMetricTable.COLUMN_ID)); + cpuMetric.setMetricId((String)source.get(CpuMetricTable.COLUMN_METRIC_ID)); + + cpuMetric.setInstanceId(((Number)source.get(CpuMetricTable.COLUMN_INSTANCE_ID)).intValue()); + + cpuMetric.setUsagePercent(((Number)source.get(CpuMetricTable.COLUMN_USAGE_PERCENT)).doubleValue()); + cpuMetric.setTimes(((Number)source.get(CpuMetricTable.COLUMN_TIMES)).longValue()); + cpuMetric.setTimeBucket(((Number)source.get(CpuMetricTable.COLUMN_TIME_BUCKET)).longValue()); + + return cpuMetric; } @Override protected final Map esStreamDataToEsData(CpuMetric streamData) { @@ -49,6 +59,7 @@ public abstract class AbstractCpuMetricEsPersistenceDAO extends AbstractPersiste source.put(CpuMetricTable.COLUMN_INSTANCE_ID, streamData.getInstanceId()); source.put(CpuMetricTable.COLUMN_USAGE_PERCENT, streamData.getUsagePercent()); + source.put(CpuMetricTable.COLUMN_TIMES, streamData.getTimes()); source.put(CpuMetricTable.COLUMN_TIME_BUCKET, streamData.getTimeBucket()); return source; diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/CpuMetricEsTableDefine.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/AbstractCpuMetricEsTableDefine.java similarity index 69% rename from apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/CpuMetricEsTableDefine.java rename to apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/AbstractCpuMetricEsTableDefine.java index 0179ddfb7..61b39134f 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/CpuMetricEsTableDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/AbstractCpuMetricEsTableDefine.java @@ -16,8 +16,7 @@ * */ - -package org.apache.skywalking.apm.collector.storage.es.define; +package org.apache.skywalking.apm.collector.storage.es.define.cpu; import org.apache.skywalking.apm.collector.storage.es.base.define.ElasticSearchColumnDefine; import org.apache.skywalking.apm.collector.storage.es.base.define.ElasticSearchTableDefine; @@ -26,19 +25,18 @@ import org.apache.skywalking.apm.collector.storage.table.jvm.CpuMetricTable; /** * @author peng-yongsheng */ -public class CpuMetricEsTableDefine extends ElasticSearchTableDefine { +public abstract class AbstractCpuMetricEsTableDefine extends ElasticSearchTableDefine { - public CpuMetricEsTableDefine() { - super(CpuMetricTable.TABLE); + public AbstractCpuMetricEsTableDefine(String name) { + super(name); } - @Override public int refreshInterval() { - return 1; - } - - @Override public void initialize() { + @Override public final void initialize() { + addColumn(new ElasticSearchColumnDefine(CpuMetricTable.COLUMN_ID, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(CpuMetricTable.COLUMN_METRIC_ID, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(CpuMetricTable.COLUMN_INSTANCE_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(CpuMetricTable.COLUMN_USAGE_PERCENT, ElasticSearchColumnDefine.Type.Double.name())); + addColumn(new ElasticSearchColumnDefine(CpuMetricTable.COLUMN_TIMES, ElasticSearchColumnDefine.Type.Long.name())); addColumn(new ElasticSearchColumnDefine(CpuMetricTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuDayMetricEsTableDefine.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuDayMetricEsTableDefine.java new file mode 100644 index 000000000..6ca9215b0 --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuDayMetricEsTableDefine.java @@ -0,0 +1,37 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.storage.es.define.cpu; + +import org.apache.skywalking.apm.collector.core.storage.TimePyramid; +import org.apache.skywalking.apm.collector.core.util.Const; +import org.apache.skywalking.apm.collector.storage.table.jvm.CpuMetricTable; + +/** + * @author peng-yongsheng + */ +public class CpuDayMetricEsTableDefine extends AbstractCpuMetricEsTableDefine { + + public CpuDayMetricEsTableDefine() { + super(CpuMetricTable.TABLE + Const.ID_SPLIT + TimePyramid.Day.getName()); + } + + @Override public int refreshInterval() { + return 1; + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuHourMetricEsTableDefine.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuHourMetricEsTableDefine.java new file mode 100644 index 000000000..85d4cb32c --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuHourMetricEsTableDefine.java @@ -0,0 +1,37 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.storage.es.define.cpu; + +import org.apache.skywalking.apm.collector.core.storage.TimePyramid; +import org.apache.skywalking.apm.collector.core.util.Const; +import org.apache.skywalking.apm.collector.storage.table.jvm.CpuMetricTable; + +/** + * @author peng-yongsheng + */ +public class CpuHourMetricEsTableDefine extends AbstractCpuMetricEsTableDefine { + + public CpuHourMetricEsTableDefine() { + super(CpuMetricTable.TABLE + Const.ID_SPLIT + TimePyramid.Hour.getName()); + } + + @Override public int refreshInterval() { + return 1; + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuMinuteMetricEsTableDefine.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuMinuteMetricEsTableDefine.java new file mode 100644 index 000000000..37f9de474 --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuMinuteMetricEsTableDefine.java @@ -0,0 +1,37 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.storage.es.define.cpu; + +import org.apache.skywalking.apm.collector.core.storage.TimePyramid; +import org.apache.skywalking.apm.collector.core.util.Const; +import org.apache.skywalking.apm.collector.storage.table.jvm.CpuMetricTable; + +/** + * @author peng-yongsheng + */ +public class CpuMinuteMetricEsTableDefine extends AbstractCpuMetricEsTableDefine { + + public CpuMinuteMetricEsTableDefine() { + super(CpuMetricTable.TABLE + Const.ID_SPLIT + TimePyramid.Minute.getName()); + } + + @Override public int refreshInterval() { + return 1; + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuMonthMetricEsTableDefine.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuMonthMetricEsTableDefine.java new file mode 100644 index 000000000..017d988b7 --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuMonthMetricEsTableDefine.java @@ -0,0 +1,37 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.storage.es.define.cpu; + +import org.apache.skywalking.apm.collector.core.storage.TimePyramid; +import org.apache.skywalking.apm.collector.core.util.Const; +import org.apache.skywalking.apm.collector.storage.table.jvm.CpuMetricTable; + +/** + * @author peng-yongsheng + */ +public class CpuMonthMetricEsTableDefine extends AbstractCpuMetricEsTableDefine { + + public CpuMonthMetricEsTableDefine() { + super(CpuMetricTable.TABLE + Const.ID_SPLIT + TimePyramid.Month.getName()); + } + + @Override public int refreshInterval() { + return 1; + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuSecondMetricEsTableDefine.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuSecondMetricEsTableDefine.java new file mode 100644 index 000000000..c84e89c07 --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/cpu/CpuSecondMetricEsTableDefine.java @@ -0,0 +1,37 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.storage.es.define.cpu; + +import org.apache.skywalking.apm.collector.core.storage.TimePyramid; +import org.apache.skywalking.apm.collector.core.util.Const; +import org.apache.skywalking.apm.collector.storage.table.jvm.CpuMetricTable; + +/** + * @author peng-yongsheng + */ +public class CpuSecondMetricEsTableDefine extends AbstractCpuMetricEsTableDefine { + + public CpuSecondMetricEsTableDefine() { + super(CpuMetricTable.TABLE + Const.ID_SPLIT + TimePyramid.Second.getName()); + } + + @Override public int refreshInterval() { + return 1; + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/resources/META-INF/defines/storage.define b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/resources/META-INF/defines/storage.define index aae5c67bf..d016737fb 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/resources/META-INF/defines/storage.define +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/resources/META-INF/defines/storage.define @@ -81,4 +81,10 @@ org.apache.skywalking.apm.collector.storage.es.define.gc.GCSecondMetricEsTableDe org.apache.skywalking.apm.collector.storage.es.define.gc.GCMinuteMetricEsTableDefine org.apache.skywalking.apm.collector.storage.es.define.gc.GCHourMetricEsTableDefine org.apache.skywalking.apm.collector.storage.es.define.gc.GCDayMetricEsTableDefine -org.apache.skywalking.apm.collector.storage.es.define.gc.GCMonthMetricEsTableDefine \ No newline at end of file +org.apache.skywalking.apm.collector.storage.es.define.gc.GCMonthMetricEsTableDefine + +org.apache.skywalking.apm.collector.storage.es.define.cpu.CpuSecondMetricEsTableDefine +org.apache.skywalking.apm.collector.storage.es.define.cpu.CpuMinuteMetricEsTableDefine +org.apache.skywalking.apm.collector.storage.es.define.cpu.CpuHourMetricEsTableDefine +org.apache.skywalking.apm.collector.storage.es.define.cpu.CpuDayMetricEsTableDefine +org.apache.skywalking.apm.collector.storage.es.define.cpu.CpuMonthMetricEsTableDefine \ No newline at end of file