From a9eaa2eecf9297f277385a1ab1104ae9284f63b2 Mon Sep 17 00:00:00 2001 From: peng-yongsheng <8082209@qq.com> Date: Tue, 21 Nov 2017 21:58:26 +0800 Subject: [PATCH] 1. The collectors provider vote about leader which collector execute the delete operation. 2. Delete the history data before the setting that you can find the property name of history_delete_before_days in the application.yml file. --- .../src/main/resources/application.yml | 1 + .../elasticsearch/ElasticSearchClient.java | 6 + .../collector/core/util/CollectionUtils.java | 9 ++ .../storage/base/dao/IPersistenceDAO.java | 2 + .../collector-storage-es-provider/pom.xml | 5 + .../storage/es/HistoryDataDeleteTimer.java | 128 ++++++++++++++++++ .../es/StorageModuleEsNamingListener.java | 42 ++++++ .../storage/es/StorageModuleEsProvider.java | 24 +++- .../es/StorageModuleEsRegistration.java | 40 ++++++ .../es/dao/CpuMetricEsPersistenceDAO.java | 15 ++ .../es/dao/GCMetricEsPersistenceDAO.java | 19 +++ .../es/dao/GlobalTraceEsPersistenceDAO.java | 15 ++ .../dao/InstPerformanceEsPersistenceDAO.java | 15 ++ .../InstanceHeartBeatEsPersistenceDAO.java | 3 + .../es/dao/MemoryMetricEsPersistenceDAO.java | 19 +++ .../dao/MemoryPoolMetricEsPersistenceDAO.java | 19 +++ .../es/dao/NodeComponentEsPersistenceDAO.java | 19 +++ .../es/dao/NodeMappingEsPersistenceDAO.java | 19 +++ .../es/dao/NodeReferenceEsPersistenceDAO.java | 19 +++ .../es/dao/SegmentCostEsPersistenceDAO.java | 15 ++ .../es/dao/SegmentEsPersistenceDAO.java | 15 ++ .../es/dao/ServiceEntryEsPersistenceDAO.java | 3 + .../dao/ServiceReferenceEsPersistenceDAO.java | 15 ++ .../h2/dao/CpuMetricH2PersistenceDAO.java | 3 + .../h2/dao/GCMetricH2PersistenceDAO.java | 3 + .../h2/dao/GlobalTraceH2PersistenceDAO.java | 3 + .../dao/InstPerformanceH2PersistenceDAO.java | 3 + .../InstanceHeartBeatH2PersistenceDAO.java | 3 + .../h2/dao/MemoryMetricH2PersistenceDAO.java | 3 + .../dao/MemoryPoolMetricH2PersistenceDAO.java | 3 + .../h2/dao/NodeComponentH2PersistenceDAO.java | 3 + .../h2/dao/NodeMappingH2PersistenceDAO.java | 3 + .../h2/dao/NodeReferenceH2PersistenceDAO.java | 3 + .../h2/dao/SegmentCostH2PersistenceDAO.java | 3 + .../h2/dao/SegmentH2PersistenceDAO.java | 3 + .../h2/dao/ServiceEntryH2PersistenceDAO.java | 3 + .../dao/ServiceReferenceH2PersistenceDAO.java | 3 + 37 files changed, 506 insertions(+), 3 deletions(-) create mode 100644 apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/HistoryDataDeleteTimer.java create mode 100644 apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsNamingListener.java create mode 100644 apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsRegistration.java diff --git a/apm-collector/apm-collector-boot/src/main/resources/application.yml b/apm-collector/apm-collector-boot/src/main/resources/application.yml index c5f2008b7..4cf075db6 100644 --- a/apm-collector/apm-collector-boot/src/main/resources/application.yml +++ b/apm-collector/apm-collector-boot/src/main/resources/application.yml @@ -37,3 +37,4 @@ ui: # cluster_nodes: localhost:9300 # index_shards_number: 2 # index_replicas_number: 0 +# history_delete_before_days: 3 \ No newline at end of file diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java index fa5df5e9a..f4843c7cd 100644 --- a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java @@ -37,6 +37,8 @@ import org.elasticsearch.client.IndicesAdminClient; import org.elasticsearch.common.settings.Settings; import org.elasticsearch.common.transport.InetSocketTransportAddress; import org.elasticsearch.common.xcontent.XContentBuilder; +import org.elasticsearch.index.reindex.DeleteByQueryAction; +import org.elasticsearch.index.reindex.DeleteByQueryRequestBuilder; import org.elasticsearch.transport.client.PreBuiltTransportClient; import org.skywalking.apm.collector.client.Client; import org.skywalking.apm.collector.client.ClientException; @@ -146,6 +148,10 @@ public class ElasticSearchClient implements Client { return client.prepareGet(indexName, "type", id); } + public DeleteByQueryRequestBuilder prepareDelete() { + return DeleteByQueryAction.INSTANCE.newRequestBuilder(client); + } + public MultiGetRequestBuilder prepareMultiGet() { return client.prepareMultiGet(); } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java index a0dd89bdb..6946b7acd 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java @@ -20,6 +20,7 @@ package org.skywalking.apm.collector.core.util; import java.util.List; import java.util.Map; +import java.util.Set; /** * @author peng-yongsheng @@ -34,10 +35,18 @@ public class CollectionUtils { return list == null || list.size() == 0; } + public static boolean isEmpty(Set set) { + return set == null || set.size() == 0; + } + public static boolean isNotEmpty(List list) { return !isEmpty(list); } + public static boolean isNotEmpty(Set set) { + return !isEmpty(set); + } + public static boolean isNotEmpty(Map map) { return !isEmpty(map); } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/base/dao/IPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/base/dao/IPersistenceDAO.java index 8219b5c99..071ecf3ad 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/base/dao/IPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/base/dao/IPersistenceDAO.java @@ -29,4 +29,6 @@ public interface IPersistenceDAO extends Insert prepareBatchInsert(DataImpl data); Update prepareBatchUpdate(DataImpl data); + + void deleteHistory(Long startTimestamp, Long endTimestamp); } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/pom.xml b/apm-collector/apm-collector-storage/collector-storage-es-provider/pom.xml index 1585a938a..fef09433c 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/pom.xml +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/pom.xml @@ -36,5 +36,10 @@ collector-storage-define ${project.version} + + org.skywalking + collector-cluster-define + ${project.version} + diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/HistoryDataDeleteTimer.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/HistoryDataDeleteTimer.java new file mode 100644 index 000000000..0900a78ec --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/HistoryDataDeleteTimer.java @@ -0,0 +1,128 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.storage.es; + +import java.util.Calendar; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import org.skywalking.apm.collector.core.module.ModuleManager; +import org.skywalking.apm.collector.core.util.CollectionUtils; +import org.skywalking.apm.collector.storage.StorageModule; +import org.skywalking.apm.collector.storage.dao.ICpuMetricPersistenceDAO; +import org.skywalking.apm.collector.storage.dao.IGCMetricPersistenceDAO; +import org.skywalking.apm.collector.storage.dao.IGlobalTracePersistenceDAO; +import org.skywalking.apm.collector.storage.dao.IInstPerformancePersistenceDAO; +import org.skywalking.apm.collector.storage.dao.IMemoryMetricPersistenceDAO; +import org.skywalking.apm.collector.storage.dao.IMemoryPoolMetricPersistenceDAO; +import org.skywalking.apm.collector.storage.dao.INodeComponentPersistenceDAO; +import org.skywalking.apm.collector.storage.dao.INodeMappingPersistenceDAO; +import org.skywalking.apm.collector.storage.dao.INodeReferencePersistenceDAO; +import org.skywalking.apm.collector.storage.dao.ISegmentCostPersistenceDAO; +import org.skywalking.apm.collector.storage.dao.ISegmentPersistenceDAO; +import org.skywalking.apm.collector.storage.dao.IServiceReferencePersistenceDAO; + +/** + * @author peng-yongsheng + */ +public class HistoryDataDeleteTimer { + + private final ModuleManager moduleManager; + private final StorageModuleEsNamingListener namingListener; + private final String selfAddress; + private final int daysBefore; + + public HistoryDataDeleteTimer(ModuleManager moduleManager, + StorageModuleEsNamingListener namingListener, String selfAddress, int daysBefore) { + this.moduleManager = moduleManager; + this.namingListener = namingListener; + this.selfAddress = selfAddress; + this.daysBefore = daysBefore; + } + + public void start() { + Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate(this::delete, 1, 8, TimeUnit.HOURS); + } + + private void tryDelete() { + if (CollectionUtils.isNotEmpty(namingListener.getAddresses())) { + String firstAddress = namingListener.getAddresses().iterator().next(); + if (firstAddress.equals(selfAddress)) { + delete(); + } + } + } + + private void delete() { + Calendar calendar = Calendar.getInstance(); + calendar.setTimeInMillis(System.currentTimeMillis()); + calendar.set(Calendar.DAY_OF_MONTH, -daysBefore); + calendar.set(Calendar.HOUR_OF_DAY, 0); + calendar.set(Calendar.MINUTE, 0); + calendar.set(Calendar.SECOND, 0); + + long startTimestamp = calendar.getTimeInMillis(); + + calendar.set(Calendar.MINUTE, 59); + calendar.set(Calendar.SECOND, 59); + long endTimestamp = calendar.getTimeInMillis(); + + deleteJVMMetricData(startTimestamp, endTimestamp); + deleteTraceMetricData(startTimestamp, endTimestamp); + } + + private void deleteJVMMetricData(long startTimestamp, long endTimestamp) { + ICpuMetricPersistenceDAO cpuMetricPersistenceDAO = moduleManager.find(StorageModule.NAME).getService(ICpuMetricPersistenceDAO.class); + cpuMetricPersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + + IGCMetricPersistenceDAO gcMetricPersistenceDAO = moduleManager.find(StorageModule.NAME).getService(IGCMetricPersistenceDAO.class); + gcMetricPersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + + IMemoryMetricPersistenceDAO memoryMetricPersistenceDAO = moduleManager.find(StorageModule.NAME).getService(IMemoryMetricPersistenceDAO.class); + memoryMetricPersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + + IMemoryPoolMetricPersistenceDAO memoryPoolMetricPersistenceDAO = moduleManager.find(StorageModule.NAME).getService(IMemoryPoolMetricPersistenceDAO.class); + memoryPoolMetricPersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + } + + private void deleteTraceMetricData(long startTimestamp, long endTimestamp) { + IGlobalTracePersistenceDAO globalTracePersistenceDAO = moduleManager.find(StorageModule.NAME).getService(IGlobalTracePersistenceDAO.class); + globalTracePersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + + IInstPerformancePersistenceDAO instPerformancePersistenceDAO = moduleManager.find(StorageModule.NAME).getService(IInstPerformancePersistenceDAO.class); + instPerformancePersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + + INodeComponentPersistenceDAO nodeComponentPersistenceDAO = moduleManager.find(StorageModule.NAME).getService(INodeComponentPersistenceDAO.class); + nodeComponentPersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + + INodeMappingPersistenceDAO nodeMappingPersistenceDAO = moduleManager.find(StorageModule.NAME).getService(INodeMappingPersistenceDAO.class); + nodeMappingPersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + + INodeReferencePersistenceDAO nodeReferencePersistenceDAO = moduleManager.find(StorageModule.NAME).getService(INodeReferencePersistenceDAO.class); + nodeReferencePersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + + ISegmentCostPersistenceDAO segmentCostPersistenceDAO = moduleManager.find(StorageModule.NAME).getService(ISegmentCostPersistenceDAO.class); + segmentCostPersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + + ISegmentPersistenceDAO segmentPersistenceDAO = moduleManager.find(StorageModule.NAME).getService(ISegmentPersistenceDAO.class); + segmentPersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + + IServiceReferencePersistenceDAO serviceReferencePersistenceDAO = moduleManager.find(StorageModule.NAME).getService(IServiceReferencePersistenceDAO.class); + serviceReferencePersistenceDAO.deleteHistory(startTimestamp, endTimestamp); + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsNamingListener.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsNamingListener.java new file mode 100644 index 000000000..eb8558c4a --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsNamingListener.java @@ -0,0 +1,42 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.storage.es; + +import org.skywalking.apm.collector.cluster.ClusterModuleListener; +import org.skywalking.apm.collector.storage.StorageModule; + +/** + * @author peng-yongsheng + */ +public class StorageModuleEsNamingListener extends ClusterModuleListener { + + public static final String PATH = "/" + StorageModule.NAME + "/" + StorageModuleEsProvider.NAME; + + @Override public String path() { + return PATH; + } + + @Override public void serverJoinNotify(String serverAddress) { + + } + + @Override public void serverQuitNotify(String serverAddress) { + + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsProvider.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsProvider.java index 43e0e7c0e..b66507344 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsProvider.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsProvider.java @@ -19,8 +19,12 @@ package org.skywalking.apm.collector.storage.es; import java.util.Properties; +import java.util.UUID; import org.skywalking.apm.collector.client.ClientException; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.cluster.ClusterModule; +import org.skywalking.apm.collector.cluster.service.ModuleListenerService; +import org.skywalking.apm.collector.cluster.service.ModuleRegisterService; import org.skywalking.apm.collector.core.module.Module; import org.skywalking.apm.collector.core.module.ModuleProvider; import org.skywalking.apm.collector.core.module.ServiceNotProvidedException; @@ -107,16 +111,19 @@ public class StorageModuleEsProvider extends ModuleProvider { private final Logger logger = LoggerFactory.getLogger(StorageModuleEsProvider.class); + public static final String NAME = "elasticsearch"; private static final String CLUSTER_NAME = "cluster_name"; private static final String CLUSTER_TRANSPORT_SNIFFER = "cluster_transport_sniffer"; private static final String CLUSTER_NODES = "cluster_nodes"; private static final String INDEX_SHARDS_NUMBER = "index_shards_number"; private static final String INDEX_REPLICAS_NUMBER = "index_replicas_number"; + private static final String HISTORY_DELETE_BEFORE_DAYS = "history_delete_before_days"; private ElasticSearchClient elasticSearchClient; + private HistoryDataDeleteTimer deleteTimer; @Override public String name() { - return "elasticsearch"; + return NAME; } @Override public Class module() { @@ -147,14 +154,25 @@ public class StorageModuleEsProvider extends ModuleProvider { } catch (ClientException | StorageException e) { logger.error(e.getMessage(), e); } + + String uuId = UUID.randomUUID().toString(); + ModuleRegisterService moduleRegisterService = getManager().find(ClusterModule.NAME).getService(ModuleRegisterService.class); + moduleRegisterService.register(StorageModule.NAME, this.name(), new StorageModuleEsRegistration(uuId, 0)); + + StorageModuleEsNamingListener namingListener = new StorageModuleEsNamingListener(); + ModuleListenerService moduleListenerService = getManager().find(ClusterModule.NAME).getService(ModuleListenerService.class); + moduleListenerService.addListener(namingListener); + + Integer beforeDay = (Integer)config.getOrDefault(HISTORY_DELETE_BEFORE_DAYS, 3); + deleteTimer = new HistoryDataDeleteTimer(getManager(), namingListener, uuId + 0, beforeDay); } @Override public void notifyAfterCompleted() throws ServiceNotProvidedException { - + deleteTimer.start(); } @Override public String[] requiredModules() { - return new String[0]; + return new String[] {ClusterModule.NAME}; } private void registerCacheDAO() throws ServiceNotProvidedException { diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsRegistration.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsRegistration.java new file mode 100644 index 000000000..4ba3dc4d9 --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/StorageModuleEsRegistration.java @@ -0,0 +1,40 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.storage.es; + +import org.skywalking.apm.collector.cluster.ModuleRegistration; +import org.skywalking.apm.collector.core.util.Const; + +/** + * @author peng-yongsheng + */ +public class StorageModuleEsRegistration extends ModuleRegistration { + + private final String virtualHost; + private final int virtualPort; + + StorageModuleEsRegistration(String virtualHost, int virtualPort) { + this.virtualHost = virtualHost; + this.virtualPort = virtualPort; + } + + @Override public Value buildValue() { + return new Value(this.virtualHost, virtualPort, Const.EMPTY_STRING); + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/CpuMetricEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/CpuMetricEsPersistenceDAO.java index 18b43116f..3d2ffb891 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/CpuMetricEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/CpuMetricEsPersistenceDAO.java @@ -22,7 +22,10 @@ import java.util.HashMap; import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.ICpuMetricPersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.jvm.CpuMetric; @@ -58,4 +61,16 @@ public class CpuMetricEsPersistenceDAO extends EsDAO implements ICpuMetricPersis @Override public UpdateRequestBuilder prepareBatchUpdate(CpuMetric cpuMetric) { return null; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(CpuMetricTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(CpuMetricTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, CpuMetricTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/GCMetricEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/GCMetricEsPersistenceDAO.java index 643480bd3..9cd87a535 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/GCMetricEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/GCMetricEsPersistenceDAO.java @@ -22,17 +22,24 @@ import java.util.HashMap; import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.IGCMetricPersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.jvm.GCMetric; import org.skywalking.apm.collector.storage.table.jvm.GCMetricTable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author peng-yongsheng */ public class GCMetricEsPersistenceDAO extends EsDAO implements IGCMetricPersistenceDAO { + private final Logger logger = LoggerFactory.getLogger(GCMetricEsPersistenceDAO.class); + public GCMetricEsPersistenceDAO(ElasticSearchClient client) { super(client); } @@ -55,4 +62,16 @@ public class GCMetricEsPersistenceDAO extends EsDAO implements IGCMetricPersiste @Override public UpdateRequestBuilder prepareBatchUpdate(GCMetric gcMetric) { return null; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(GCMetricTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(GCMetricTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, GCMetricTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/GlobalTraceEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/GlobalTraceEsPersistenceDAO.java index 27c4aa3bb..412b40cdc 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/GlobalTraceEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/GlobalTraceEsPersistenceDAO.java @@ -22,8 +22,11 @@ import java.util.HashMap; import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; import org.skywalking.apm.collector.core.UnexpectedException; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.IGlobalTracePersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.global.GlobalTrace; @@ -58,4 +61,16 @@ public class GlobalTraceEsPersistenceDAO extends EsDAO implements IGlobalTracePe logger.debug("global trace source: {}", source.toString()); return getClient().prepareIndex(GlobalTraceTable.TABLE, data.getId()).setSource(source); } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(GlobalTraceTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(GlobalTraceTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, GlobalTraceTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/InstPerformanceEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/InstPerformanceEsPersistenceDAO.java index 59aa39a69..06e811adf 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/InstPerformanceEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/InstPerformanceEsPersistenceDAO.java @@ -23,7 +23,10 @@ import java.util.Map; import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.IInstPerformancePersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.instance.InstPerformance; @@ -80,4 +83,16 @@ public class InstPerformanceEsPersistenceDAO extends EsDAO implements IInstPerfo return getClient().prepareUpdate(InstPerformanceTable.TABLE, data.getId()).setDoc(source); } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(InstPerformanceTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(InstPerformanceTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, InstPerformanceTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/InstanceHeartBeatEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/InstanceHeartBeatEsPersistenceDAO.java index 1cd0756c3..1fe3b15dc 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/InstanceHeartBeatEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/InstanceHeartBeatEsPersistenceDAO.java @@ -67,4 +67,7 @@ public class InstanceHeartBeatEsPersistenceDAO extends EsDAO implements IInstanc source.put(InstanceTable.COLUMN_HEARTBEAT_TIME, data.getHeartBeatTime()); return getClient().prepareUpdate(InstanceTable.TABLE, data.getId()).setDoc(source); } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/MemoryMetricEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/MemoryMetricEsPersistenceDAO.java index 9e81860f9..252f50c33 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/MemoryMetricEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/MemoryMetricEsPersistenceDAO.java @@ -22,17 +22,24 @@ import java.util.HashMap; import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.IMemoryMetricPersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.jvm.MemoryMetric; import org.skywalking.apm.collector.storage.table.jvm.MemoryMetricTable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author peng-yongsheng */ public class MemoryMetricEsPersistenceDAO extends EsDAO implements IMemoryMetricPersistenceDAO { + private final Logger logger = LoggerFactory.getLogger(MemoryMetricEsPersistenceDAO.class); + public MemoryMetricEsPersistenceDAO(ElasticSearchClient client) { super(client); } @@ -57,4 +64,16 @@ public class MemoryMetricEsPersistenceDAO extends EsDAO implements IMemoryMetric @Override public UpdateRequestBuilder prepareBatchUpdate(MemoryMetric data) { return null; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(MemoryMetricTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(MemoryMetricTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, MemoryMetricTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/MemoryPoolMetricEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/MemoryPoolMetricEsPersistenceDAO.java index aae64b469..15d5f909d 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/MemoryPoolMetricEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/MemoryPoolMetricEsPersistenceDAO.java @@ -22,17 +22,24 @@ import java.util.HashMap; import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.IMemoryPoolMetricPersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.jvm.MemoryPoolMetric; import org.skywalking.apm.collector.storage.table.jvm.MemoryPoolMetricTable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author peng-yongsheng */ public class MemoryPoolMetricEsPersistenceDAO extends EsDAO implements IMemoryPoolMetricPersistenceDAO { + private final Logger logger = LoggerFactory.getLogger(MemoryPoolMetricEsPersistenceDAO.class); + public MemoryPoolMetricEsPersistenceDAO(ElasticSearchClient client) { super(client); } @@ -57,4 +64,16 @@ public class MemoryPoolMetricEsPersistenceDAO extends EsDAO implements IMemoryPo @Override public UpdateRequestBuilder prepareBatchUpdate(MemoryPoolMetric data) { return null; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getSecondTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(MemoryPoolMetricTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(MemoryPoolMetricTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, MemoryPoolMetricTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeComponentEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeComponentEsPersistenceDAO.java index 5dad91001..d1fde779c 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeComponentEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeComponentEsPersistenceDAO.java @@ -23,17 +23,24 @@ import java.util.Map; import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.INodeComponentPersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.node.NodeComponent; import org.skywalking.apm.collector.storage.table.node.NodeComponentTable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author peng-yongsheng */ public class NodeComponentEsPersistenceDAO extends EsDAO implements INodeComponentPersistenceDAO { + private final Logger logger = LoggerFactory.getLogger(NodeComponentEsPersistenceDAO.class); + public NodeComponentEsPersistenceDAO(ElasticSearchClient client) { super(client); } @@ -69,4 +76,16 @@ public class NodeComponentEsPersistenceDAO extends EsDAO implements INodeCompone return getClient().prepareUpdate(NodeComponentTable.TABLE, data.getId()).setDoc(source); } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(NodeComponentTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(NodeComponentTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, NodeComponentTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeMappingEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeMappingEsPersistenceDAO.java index 0372faadb..08d341339 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeMappingEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeMappingEsPersistenceDAO.java @@ -23,17 +23,24 @@ import java.util.Map; import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.INodeMappingPersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.node.NodeMapping; import org.skywalking.apm.collector.storage.table.node.NodeMappingTable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author peng-yongsheng */ public class NodeMappingEsPersistenceDAO extends EsDAO implements INodeMappingPersistenceDAO { + private final Logger logger = LoggerFactory.getLogger(NodeMappingEsPersistenceDAO.class); + public NodeMappingEsPersistenceDAO(ElasticSearchClient client) { super(client); } @@ -68,4 +75,16 @@ public class NodeMappingEsPersistenceDAO extends EsDAO implements INodeMappingPe source.put(NodeMappingTable.COLUMN_TIME_BUCKET, data.getTimeBucket()); return getClient().prepareUpdate(NodeMappingTable.TABLE, data.getId()).setDoc(source); } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(NodeMappingTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(NodeMappingTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, NodeMappingTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeReferenceEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeReferenceEsPersistenceDAO.java index c58f94061..784acf513 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeReferenceEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/NodeReferenceEsPersistenceDAO.java @@ -23,17 +23,24 @@ import java.util.Map; import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.INodeReferencePersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.noderef.NodeReference; import org.skywalking.apm.collector.storage.table.noderef.NodeReferenceTable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author peng-yongsheng */ public class NodeReferenceEsPersistenceDAO extends EsDAO implements INodeReferencePersistenceDAO { + private final Logger logger = LoggerFactory.getLogger(NodeReferenceEsPersistenceDAO.class); + public NodeReferenceEsPersistenceDAO(ElasticSearchClient client) { super(client); } @@ -87,4 +94,16 @@ public class NodeReferenceEsPersistenceDAO extends EsDAO implements INodeReferen return getClient().prepareUpdate(NodeReferenceTable.TABLE, data.getId()).setDoc(source); } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(NodeReferenceTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(NodeReferenceTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, NodeReferenceTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/SegmentCostEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/SegmentCostEsPersistenceDAO.java index 8e1b0eeac..56f2ce801 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/SegmentCostEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/SegmentCostEsPersistenceDAO.java @@ -22,7 +22,10 @@ import java.util.HashMap; import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.ISegmentCostPersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.segment.SegmentCost; @@ -63,4 +66,16 @@ public class SegmentCostEsPersistenceDAO extends EsDAO implements ISegmentCostPe logger.debug("segment cost source: {}", source.toString()); return getClient().prepareIndex(SegmentCostTable.TABLE, data.getId()).setSource(source); } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(SegmentCostTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(SegmentCostTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, SegmentCostTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/SegmentEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/SegmentEsPersistenceDAO.java index a1011e0ec..ea08369d6 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/SegmentEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/SegmentEsPersistenceDAO.java @@ -23,7 +23,10 @@ import java.util.HashMap; import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.ISegmentPersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.segment.Segment; @@ -57,4 +60,16 @@ public class SegmentEsPersistenceDAO extends EsDAO implements ISegmentPersistenc logger.debug("segment source: {}", source.toString()); return getClient().prepareIndex(SegmentTable.TABLE, data.getId()).setSource(source); } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(SegmentTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(SegmentTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, SegmentTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/ServiceEntryEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/ServiceEntryEsPersistenceDAO.java index a4f16d612..e04c367d2 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/ServiceEntryEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/ServiceEntryEsPersistenceDAO.java @@ -74,4 +74,7 @@ public class ServiceEntryEsPersistenceDAO extends EsDAO implements IServiceEntry return getClient().prepareUpdate(ServiceEntryTable.TABLE, data.getId()).setDoc(source); } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/ServiceReferenceEsPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/ServiceReferenceEsPersistenceDAO.java index ff96008b4..f39f3ca44 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/ServiceReferenceEsPersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/skywalking/apm/collector/storage/es/dao/ServiceReferenceEsPersistenceDAO.java @@ -23,7 +23,10 @@ import java.util.Map; import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.storage.dao.IServiceReferencePersistenceDAO; import org.skywalking.apm.collector.storage.es.base.dao.EsDAO; import org.skywalking.apm.collector.storage.table.serviceref.ServiceReference; @@ -97,4 +100,16 @@ public class ServiceReferenceEsPersistenceDAO extends EsDAO implements IServiceR return getClient().prepareUpdate(ServiceReferenceTable.TABLE, data.getId()).setDoc(source); } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + long startTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(startTimestamp); + long endTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(endTimestamp); + BulkByScrollResponse response = getClient().prepareDelete() + .filter(QueryBuilders.rangeQuery(ServiceReferenceTable.COLUMN_TIME_BUCKET).gte(startTimeBucket).lte(endTimeBucket)) + .source(ServiceReferenceTable.TABLE) + .get(); + + long deleted = response.getDeleted(); + logger.info("Delete {} rows history from {} index.", deleted, ServiceReferenceTable.TABLE); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/CpuMetricH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/CpuMetricH2PersistenceDAO.java index 4d146d192..5faa595ee 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/CpuMetricH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/CpuMetricH2PersistenceDAO.java @@ -63,4 +63,7 @@ public class CpuMetricH2PersistenceDAO extends H2DAO implements ICpuMetricPersis @Override public H2SqlEntity prepareBatchUpdate(CpuMetric data) { return null; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/GCMetricH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/GCMetricH2PersistenceDAO.java index 9cca8ddaa..81616f249 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/GCMetricH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/GCMetricH2PersistenceDAO.java @@ -60,4 +60,7 @@ public class GCMetricH2PersistenceDAO extends H2DAO implements IGCMetricPersiste @Override public H2SqlEntity prepareBatchUpdate(GCMetric data) { return null; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/GlobalTraceH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/GlobalTraceH2PersistenceDAO.java index 805670fca..837226ae5 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/GlobalTraceH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/GlobalTraceH2PersistenceDAO.java @@ -64,4 +64,7 @@ public class GlobalTraceH2PersistenceDAO extends H2DAO implements IGlobalTracePe entity.setParams(source.values().toArray(new Object[0])); return entity; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/InstPerformanceH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/InstPerformanceH2PersistenceDAO.java index a42773530..f3aabc4e3 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/InstPerformanceH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/InstPerformanceH2PersistenceDAO.java @@ -97,4 +97,7 @@ public class InstPerformanceH2PersistenceDAO extends H2DAO implements IInstPerfo entity.setParams(values.toArray(new Object[0])); return entity; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/InstanceHeartBeatH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/InstanceHeartBeatH2PersistenceDAO.java index fd8e3982e..db94d82a7 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/InstanceHeartBeatH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/InstanceHeartBeatH2PersistenceDAO.java @@ -81,4 +81,7 @@ public class InstanceHeartBeatH2PersistenceDAO extends H2DAO implements IInstanc entity.setParams(params.toArray(new Object[0])); return entity; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/MemoryMetricH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/MemoryMetricH2PersistenceDAO.java index 0b28d6ba0..404343462 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/MemoryMetricH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/MemoryMetricH2PersistenceDAO.java @@ -62,4 +62,7 @@ public class MemoryMetricH2PersistenceDAO extends H2DAO implements IMemoryMetric @Override public H2SqlEntity prepareBatchUpdate(MemoryMetric data) { return null; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/MemoryPoolMetricH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/MemoryPoolMetricH2PersistenceDAO.java index c5fe14eeb..aea2d72fc 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/MemoryPoolMetricH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/MemoryPoolMetricH2PersistenceDAO.java @@ -62,4 +62,7 @@ public class MemoryPoolMetricH2PersistenceDAO extends H2DAO implements IMemoryPo @Override public H2SqlEntity prepareBatchUpdate(MemoryPoolMetric data) { return null; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeComponentH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeComponentH2PersistenceDAO.java index 50f059840..15c33d0b7 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeComponentH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeComponentH2PersistenceDAO.java @@ -94,4 +94,7 @@ public class NodeComponentH2PersistenceDAO extends H2DAO implements INodeCompone entity.setParams(values.toArray(new Object[0])); return entity; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeMappingH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeMappingH2PersistenceDAO.java index 7d170c4a8..bb8ee9fb3 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeMappingH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeMappingH2PersistenceDAO.java @@ -92,4 +92,7 @@ public class NodeMappingH2PersistenceDAO extends H2DAO implements INodeMappingPe entity.setParams(values.toArray(new Object[0])); return entity; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeReferenceH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeReferenceH2PersistenceDAO.java index d690aebb3..c5cc58d16 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeReferenceH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/NodeReferenceH2PersistenceDAO.java @@ -110,4 +110,7 @@ public class NodeReferenceH2PersistenceDAO extends H2DAO implements INodeReferen entity.setParams(values.toArray(new Object[0])); return entity; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/SegmentCostH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/SegmentCostH2PersistenceDAO.java index eae890307..31939c078 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/SegmentCostH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/SegmentCostH2PersistenceDAO.java @@ -69,4 +69,7 @@ public class SegmentCostH2PersistenceDAO extends H2DAO implements ISegmentCostPe @Override public H2SqlEntity prepareBatchUpdate(SegmentCost data) { return null; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/SegmentH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/SegmentH2PersistenceDAO.java index c0e0d85e3..4d207ca4e 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/SegmentH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/SegmentH2PersistenceDAO.java @@ -62,4 +62,7 @@ public class SegmentH2PersistenceDAO extends H2DAO implements ISegmentPersistenc @Override public H2SqlEntity prepareBatchUpdate(Segment data) { return null; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/ServiceEntryH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/ServiceEntryH2PersistenceDAO.java index d74abfda5..da6311fd8 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/ServiceEntryH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/ServiceEntryH2PersistenceDAO.java @@ -97,4 +97,7 @@ public class ServiceEntryH2PersistenceDAO extends H2DAO implements IServiceEntry entity.setParams(values.toArray(new Object[0])); return entity; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/ServiceReferenceH2PersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/ServiceReferenceH2PersistenceDAO.java index 9b18bb5aa..bdcd189fb 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/ServiceReferenceH2PersistenceDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/skywalking/apm/collector/storage/h2/dao/ServiceReferenceH2PersistenceDAO.java @@ -120,4 +120,7 @@ public class ServiceReferenceH2PersistenceDAO extends H2DAO implements IServiceR entity.setParams(values.toArray(new Object[0])); return entity; } + + @Override public void deleteHistory(Long startTimestamp, Long endTimestamp) { + } }