diff --git a/CHANGES.md b/CHANGES.md index fc932f6205..e399545494 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -88,8 +88,13 @@ Release Notes. * Remove `syncBulkActions` in ElasticSearch storage option. * Increase the default bulkActions(env, SW_STORAGE_ES_BULK_ACTIONS) to 5000(from 1000). * Increase the flush interval of ElasticSearch indices to 15s(from 10s) +* Provide distinct for elements of metadata lists. Due to the more aggressive asynchronous flush, metadata lists have + more chances including duplicate elements. Don't need this as indicate anymore. +* Reduce the flush period of hour and day level metrics, only run in 4 times of regular persistent period. This means + default flush period of hour and day level metrics are 25s * 4. #### UI + * Fix the date component for log conditions. * Fix selector keys for duplicate options. * Add Python celery plugin. @@ -99,6 +104,7 @@ Release Notes. * Fix chart types for setting metrics configure. #### Documentation + * Add FAQ about `Elasticsearch exception type=version_conflict_engine_exception since 8.7.0` All issues and pull requests are [here](https://github.com/apache/skywalking/milestone/90?closed=1) diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsPersistentWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsPersistentWorker.java index 5dd29db2ad..444ec21725 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsPersistentWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsPersistentWorker.java @@ -67,6 +67,16 @@ public class MetricsPersistentWorker extends PersistenceWorker { private final boolean supportUpdate; private long sessionTimeout; private CounterMetrics aggregationCounter; + /** + * The counter for the round of persistent. + */ + private int persistentCounter; + /** + * The mod value to control persistent. The MetricsPersistentWorker is driven by the {@link + * org.apache.skywalking.oap.server.core.storage.PersistenceTimer}. The down sampling level workers only execute in + * every {@link #persistentMod} periods. And minute level workers execute every time. + */ + private int persistentMod; MetricsPersistentWorker(ModuleDefineHolder moduleDefineHolder, Model model, IMetricsDAO metricsDAO, AbstractWorker nextAlarmWorker, AbstractWorker nextExportWorker, @@ -82,6 +92,8 @@ public class MetricsPersistentWorker extends PersistenceWorker { this.transWorker = Optional.ofNullable(transWorker); this.supportUpdate = supportUpdate; this.sessionTimeout = storageSessionTimeout; + this.persistentCounter = 0; + this.persistentMod = 1; String name = "METRICS_L2_AGGREGATION"; int size = BulkConsumePool.Creator.recommendMaxSize() / 8; @@ -122,6 +134,8 @@ public class MetricsPersistentWorker extends PersistenceWorker { // And add offset according to worker creation sequence, to avoid context clear overlap, // eventually optimize load of IDs reading. this.sessionTimeout = this.sessionTimeout * 4 + SESSION_TIMEOUT_OFFSITE_COUNTER * 200; + // The down sampling level worker executes every 4 periods. + this.persistentMod = 4; } /** @@ -135,6 +149,10 @@ public class MetricsPersistentWorker extends PersistenceWorker { @Override public void prepareBatch(Collection lastCollection, List prepareRequests) { + if (persistentCounter++ % persistentMod != 0) { + return; + } + long start = System.currentTimeMillis(); if (lastCollection.size() == 0) { return; diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/MetadataQueryService.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/MetadataQueryService.java index 7f6d67442f..6459208c56 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/MetadataQueryService.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/MetadataQueryService.java @@ -54,30 +54,34 @@ public class MetadataQueryService implements org.apache.skywalking.oap.server.li if (service.getGroup() == null) { service.setGroup(Const.EMPTY_STRING); } - }).collect(Collectors.toList()); + }) + .distinct() + .collect(Collectors.toList()); } public List getAllBrowserServices() throws IOException { - return getMetadataQueryDAO().getAllBrowserServices(); + return getMetadataQueryDAO().getAllBrowserServices().stream().distinct().collect(Collectors.toList()); } public List getAllDatabases() throws IOException { - return getMetadataQueryDAO().getAllDatabases(); + return getMetadataQueryDAO().getAllDatabases().stream().distinct().collect(Collectors.toList()); } public List searchServices(final long startTimestamp, final long endTimestamp, final String keyword) throws IOException { - return getMetadataQueryDAO().searchServices(keyword); + return getMetadataQueryDAO().searchServices(keyword).stream().distinct().collect(Collectors.toList()); } public List getServiceInstances(final long startTimestamp, final long endTimestamp, final String serviceId) throws IOException { - return getMetadataQueryDAO().getServiceInstances(startTimestamp, endTimestamp, serviceId); + return getMetadataQueryDAO().getServiceInstances(startTimestamp, endTimestamp, serviceId) + .stream().distinct().collect(Collectors.toList()); } public List searchEndpoint(final String keyword, final String serviceId, final int limit) throws IOException { - return getMetadataQueryDAO().searchEndpoint(keyword, serviceId, limit); + return getMetadataQueryDAO().searchEndpoint(keyword, serviceId, limit) + .stream().distinct().collect(Collectors.toList()); } public Service searchService(final String serviceCode) throws IOException {