From 889ae9dbe8c317040cc4759e305a655493382022 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=90=B4=E6=99=9F=20Wu=20Sheng?= Date: Mon, 5 May 2025 22:53:09 +0800 Subject: [PATCH] Increase the idle check interval of the message queue to 200ms to reduce CPU usage under low load conditions (#13227) --- docs/en/changes/changes.md | 1 + .../core/analysis/worker/MetricsAggregateWorker.java | 2 +- .../core/analysis/worker/MetricsPersistentWorker.java | 2 +- .../oap/server/library/datacarrier/DataCarrier.java | 8 ++++---- 4 files changed, 7 insertions(+), 6 deletions(-) diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index baf6d20d4b..e8f2e06d58 100644 --- a/docs/en/changes/changes.md +++ b/docs/en/changes/changes.md @@ -14,6 +14,7 @@ * Support Flink monitoring. * BanyanDB: Support `@ShardingKey` for Measure tags and set to TopNAggregation group tag by default. * BanyanDB: Support cold stage data query for metrics/traces/logs. +* Increase the idle check interval of the message queue to 200ms to reduce CPU usage under low load conditions. #### UI diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsAggregateWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsAggregateWorker.java index 255aa44306..d9a64d94dd 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsAggregateWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsAggregateWorker.java @@ -73,7 +73,7 @@ public class MetricsAggregateWorker extends AbstractWorker { "MetricsAggregateWorker." + modelName, name, queueChannelSize, queueBufferSize, BufferStrategy.IF_POSSIBLE); BulkConsumePool.Creator creator = new BulkConsumePool.Creator( - name, BulkConsumePool.Creator.recommendMaxSize() * 2, 20); + name, BulkConsumePool.Creator.recommendMaxSize() * 2, 200); try { ConsumerPoolFactory.INSTANCE.createIfAbsent(name, creator); } catch (Exception e) { 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 349a60d2fd..47f4a4f309 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 @@ -127,7 +127,7 @@ public class MetricsPersistentWorker extends PersistenceWorker implemen if (size == 0) { size = 1; } - BulkConsumePool.Creator creator = new BulkConsumePool.Creator(name, size, 20); + BulkConsumePool.Creator creator = new BulkConsumePool.Creator(name, size, 200); try { ConsumerPoolFactory.INSTANCE.createIfAbsent(name, creator); } catch (Exception e) { diff --git a/oap-server/server-library/library-datacarrier-queue/src/main/java/org/apache/skywalking/oap/server/library/datacarrier/DataCarrier.java b/oap-server/server-library/library-datacarrier-queue/src/main/java/org/apache/skywalking/oap/server/library/datacarrier/DataCarrier.java index 5183cc74d0..86dd497f92 100644 --- a/oap-server/server-library/library-datacarrier-queue/src/main/java/org/apache/skywalking/oap/server/library/datacarrier/DataCarrier.java +++ b/oap-server/server-library/library-datacarrier-queue/src/main/java/org/apache/skywalking/oap/server/library/datacarrier/DataCarrier.java @@ -106,14 +106,14 @@ public class DataCarrier { } /** - * set consumeDriver to this Carrier. consumer begin to run when {@link DataCarrier#produce} begin to work with 20 + * set consumeDriver to this Carrier. consumer begins to run when {@link DataCarrier#produce} begin to work with 200 * millis consume cycle. * * @param consumerClass class of consumer * @param num number of consumer threads */ public DataCarrier consume(Class> consumerClass, int num) { - return this.consume(consumerClass, num, 20, new Properties()); + return this.consume(consumerClass, num, 200, new Properties()); } /** @@ -132,14 +132,14 @@ public class DataCarrier { } /** - * set consumeDriver to this Carrier. consumer begin to run when {@link DataCarrier#produce} begin to work with 20 + * set consumeDriver to this Carrier. consumer begin to run when {@link DataCarrier#produce} begin to work with 200 * millis consume cycle. * * @param consumer single instance of consumer, all consumer threads will all use this instance. * @param num number of consumer threads */ public DataCarrier consume(IConsumer consumer, int num) { - return this.consume(consumer, num, 20); + return this.consume(consumer, num, 200); } /**