[feat]: migrate blocking executors to virtual threads - #4062
Merged
Conversation
Contributor
There was a problem hiding this comment.
Pull request overview
This PR modernizes HertzBeat’s concurrency infrastructure by introducing virtual-thread aware executors and configuration (hertzbeat.vthreads), then migrating multiple modules’ worker pools and scheduled tasks to use these executors with improved shutdown semantics and concurrency controls.
Changes:
- Added
ManagedExecutor/ManagedExecutors+AdmissionMode, and introducedVirtualThreadProperties(hertzbeat.vthreads) to configure virtual-thread execution and admission limits. - Migrated worker pools / schedulers across collector, manager, alerter, warehouse, remoting, and log modules to virtual-thread executors; added “long-running” execution lanes and cleanup on shutdown.
- Added docs + example YAML blocks, plus extensive tests validating virtual-thread execution and concurrency limiting behavior.
Reviewed changes
Copilot reviewed 71 out of 71 changed files in this pull request and generated 10 comments.
Show a summary per file
| File | Description |
|---|---|
| script/docker-compose/hertzbeat-postgresql-victoria-metrics/conf/application.yml | Adds hertzbeat.vthreads example overrides for compose deployments. |
| script/docker-compose/hertzbeat-postgresql-greptimedb/conf/application.yml | Adds hertzbeat.vthreads example overrides for compose deployments. |
| script/docker-compose/hertzbeat-mysql-victoria-metrics/conf/application.yml | Adds hertzbeat.vthreads example overrides for compose deployments. |
| script/docker-compose/hertzbeat-mysql-tdengine/conf/application.yml | Adds hertzbeat.vthreads example overrides for compose deployments. |
| script/docker-compose/hertzbeat-mysql-iotdb/conf/application.yml | Adds hertzbeat.vthreads example overrides for compose deployments. |
| script/application.yml | Adds hertzbeat.vthreads example overrides for script-based deployments. |
| home/i18n/zh-cn/docusaurus-plugin-content-docs/current/start/virtual-thread.md | New Chinese documentation for virtual-thread defaults, tuning, and rollback. |
| home/i18n/zh-cn/docusaurus-plugin-content-docs/current/start/custom-config.md | Links to the new virtual-thread documentation page (zh-cn). |
| home/i18n/zh-cn/docusaurus-plugin-content-docs/current/help/issue.md | Updates collector task-limit FAQ to reflect virtual-thread concurrency defaults (zh-cn). |
| home/docs/start/virtual-thread.md | New English documentation for virtual-thread defaults, tuning, and rollback. |
| home/docs/start/custom-config.md | Links to the new virtual-thread documentation page (en). |
| home/docs/help/issue.md | Updates collector task-limit FAQ to reflect virtual-thread concurrency defaults (en). |
| hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/duckdb/DuckdbDatabaseDataStorageTest.java | Adds tests for DuckDB cleaner virtual-thread execution + non-concurrency behavior. |
| hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/WarehouseWorkerPoolTest.java | Updates tests for warehouse pool: virtual threads + rejection behavior. |
| hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/duckdb/DuckdbDatabaseDataStorage.java | Refactors DuckDB cleaner scheduling/execution; adds vthread executor + shutdown. |
| hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/doris/DorisDataStorage.java | Moves long-running flush to executeLongRunning. |
| hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/DataStorageDispatch.java | Runs long-lived storage consumers via executeLongRunning. |
| hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/WarehouseWorkerPool.java | Migrates to ManagedExecutor, adds long-running lane, and implements DisposableBean. |
| hertzbeat-startup/src/test/java/org/apache/hertzbeat/startup/AsyncConfigTest.java | Adds tests for @Async executor virtual-thread usage and rejection behavior. |
| hertzbeat-startup/src/main/resources/application.yml | Adds hertzbeat.vthreads example overrides for startup module config. |
| hertzbeat-startup/src/main/resources/application-test.yml | Adds hertzbeat.vthreads overrides for startup test config. |
| hertzbeat-startup/src/main/java/org/apache/hertzbeat/startup/AsyncConfig.java | Introduces dedicated taskExecutor bean with virtual-thread + concurrency controls. |
| hertzbeat-remoting/src/main/java/org/apache/hertzbeat/remoting/netty/NettyRemotingServer.java | Runs server bootstrap as long-running task. |
| hertzbeat-remoting/src/main/java/org/apache/hertzbeat/remoting/netty/NettyRemotingClient.java | Runs client bootstrap as long-running task. |
| hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/scheduler/netty/ManageServerTest.java | Adds tests for channel health-check virtual-thread execution and deduping. |
| hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/scheduler/ManagerWorkerPoolTest.java | Adds tests for manager pool virtual threads, rejection, and platform long-running lane. |
| hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/component/status/CalculateStatusTest.java | Adds tests for status calculation virtual threads + non-concurrency behavior. |
| hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/scheduler/netty/ManageServer.java | Runs channel health-checks via virtual-thread executor with deduping + cleanup. |
| hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/scheduler/ManagerWorkerPool.java | Migrates manager pool to ManagedExecutor with long-running lane + cleanup. |
| hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/component/status/CalculateStatus.java | Refactors schedulers to dispatch onto virtual threads with deduping + cleanup. |
| hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/component/sd/ServiceDiscoveryWorker.java | Runs service discovery loop via executeLongRunning. |
| hertzbeat-log/src/test/java/org/apache/hertzbeat/log/notice/LogSseManagerTest.java | Adds tests for SSE virtual-thread sending + queue size cap behavior. |
| hertzbeat-log/src/main/java/org/apache/hertzbeat/log/notice/LogSseManager.java | Migrates SSE sender to virtual-thread-per-task executor. |
| hertzbeat-common-spring/src/test/java/org/apache/hertzbeat/common/support/CommonThreadPoolTest.java | Updates tests to validate virtual-thread usage, rejection, and long-running behavior. |
| hertzbeat-common-spring/src/test/java/org/apache/hertzbeat/common/config/VirtualThreadPropertiesTest.java | Adds tests for VirtualThreadProperties.defaults() safety expectations. |
| hertzbeat-common-spring/src/main/java/org/apache/hertzbeat/common/support/CommonThreadPool.java | Migrates common pool to ManagedExecutor with long-running lane + cleanup. |
| hertzbeat-common-spring/src/main/java/org/apache/hertzbeat/common/config/VirtualThreadProperties.java | Introduces hertzbeat.vthreads configuration model + defaults helpers. |
| hertzbeat-common-spring/src/main/java/org/apache/hertzbeat/common/config/CommonConfig.java | Registers VirtualThreadProperties as configuration properties. |
| hertzbeat-common-core/src/test/java/org/apache/hertzbeat/common/concurrent/ManagedExecutorsTest.java | Adds tests for managed executors: virtual thread, queued, discard-oldest, rejection. |
| hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/concurrent/ManagedExecutors.java | Adds managed executor factories + admission/queue semantics. |
| hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/concurrent/ManagedExecutor.java | Adds ManagedExecutor interface (named, closeable). |
| hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/concurrent/AdmissionMode.java | Adds admission mode enum for executor admission control. |
| hertzbeat-collector/hertzbeat-collector-rocketmq/src/test/java/org/apache/hertzbeat/collector/collect/rocketmq/RocketmqSingleCollectTest.java | Adds tests for consumer tasks on virtual threads and discard-oldest behavior. |
| hertzbeat-collector/hertzbeat-collector-rocketmq/src/main/java/org/apache/hertzbeat/collector/collect/rocketmq/RocketmqSingleCollectImpl.java | Replaces custom pool with managed discard-oldest virtual executor + cleanup. |
| hertzbeat-collector/hertzbeat-collector-rocketmq/pom.xml | Removes redundant module-level compiler properties. |
| hertzbeat-collector/hertzbeat-collector-common/src/test/java/org/apache/hertzbeat/collector/collect/common/http/CommonHttpClientVirtualThreadTest.java | Adds tests for HTTP pool cleanup dispatch on virtual threads + deduping. |
| hertzbeat-collector/hertzbeat-collector-common/src/test/java/org/apache/hertzbeat/collector/collect/common/cache/GlobalConnectionCacheTest.java | Adds tests for connection cache cleanup dispatch on virtual threads + deduping. |
| hertzbeat-collector/hertzbeat-collector-common/src/main/java/org/apache/hertzbeat/collector/dispatch/entrance/CollectServer.java | Moves heartbeat scheduling to virtual-thread executor with deduping + cleanup. |
| hertzbeat-collector/hertzbeat-collector-common/src/main/java/org/apache/hertzbeat/collector/dispatch/WorkerPool.java | Migrates collector worker pool to ManagedExecutor + long-running lane + cleanup. |
| hertzbeat-collector/hertzbeat-collector-common/src/main/java/org/apache/hertzbeat/collector/collect/common/http/CommonHttpClient.java | Refactors connection pool cleanup to virtual-thread dispatch with deduping + shutdown. |
| hertzbeat-collector/hertzbeat-collector-common/src/main/java/org/apache/hertzbeat/collector/collect/common/cache/GlobalConnectionCache.java | Refactors cache cleanup to virtual-thread dispatch with deduping + shutdown hook. |
| hertzbeat-collector/hertzbeat-collector-common/pom.xml | Adds dependency on hertzbeat-common-spring for config injection. |
| hertzbeat-collector/hertzbeat-collector-collector/src/main/resources/application.yml | Adds hertzbeat.vthreads example overrides for standalone collector config. |
| hertzbeat-collector/hertzbeat-collector-collector/src/main/java/org/apache/hertzbeat/collector/dispatch/CommonDispatcher.java | Runs dispatcher loop via executeLongRunning. |
| hertzbeat-collector/hertzbeat-collector-basic/src/test/java/org/apache/hertzbeat/collector/dispatch/entrance/CollectServerTest.java | Adds tests for heartbeat dispatch virtual-thread execution + deduping. |
| hertzbeat-collector/hertzbeat-collector-basic/src/test/java/org/apache/hertzbeat/collector/dispatch/WorkerPoolTest.java | Updates tests to validate collector worker pool virtual threads + rejection. |
| hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/reduce/AlarmInhibitReduceTest.java | Adds tests for inhibit cleanup dispatch virtual threads + deduping. |
| hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/reduce/AlarmGroupReduceTest.java | Adds tests for group check dispatch virtual threads + deduping. |
| hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/reduce/AlarmCommonReduceTest.java | Adds tests for reduce execution on vthreads + queued behavior under limits. |
| hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/notice/AlertNoticeDispatchTest.java | Verifies notify dispatch routes through typed per-channel execution. |
| hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/calculate/realtime/window/AlarmEvaluatorTest.java | Adds tests for evaluator vthread execution + queued behavior. |
| hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/calculate/periodic/PeriodicAlertRuleSchedulerTest.java | Adds tests for periodic scheduler vthreads, deduping, concurrency limit, cancellation. |
| hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/AlerterWorkerPoolTest.java | Adds tests for notify/log vthreads and rejection behavior including per-channel limits. |
| hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/reduce/AlarmInhibitReduce.java | Adds vthread-backed cleanup dispatch + shutdown hooks for inhibit cache cleanup. |
| hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/reduce/AlarmGroupReduce.java | Adds vthread-backed scheduled group checks + shutdown semantics. |
| hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/reduce/AlarmCommonReduce.java | Migrates reduce executor to managed queued virtual executor + cleanup. |
| hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/notice/AlertNoticeDispatch.java | Dispatches notify tasks via per-channel typed execution with rejection handling. |
| hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/calculate/realtime/window/AlarmEvaluator.java | Migrates evaluator executor to managed queued virtual executor + cleanup. |
| hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/calculate/periodic/PeriodicAlertRuleScheduler.java | Adds global periodic concurrency limit, vthread execution, deduping, and cleanup. |
| hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/AlerterWorkerPool.java | Migrates notify/log pools to managed executors + per-channel semaphore limits + cleanup. |
| hertzbeat-alerter/pom.xml | Adds hertzbeat-common-spring dependency for config injection. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
This pull request introduces significant improvements to thread management and concurrency control in the alerting and monitoring modules. The main focus is on enabling support for virtual threads, improving resource management, and enhancing per-channel concurrency limits for notification tasks. Additionally, some refactoring and dependency updates are included to support these changes.
Thread Management and Virtual Thread Support:
Refactored
AlerterWorkerPoolto useManagedExecutorandManagedExecutors, adding support for virtual threads and improved resource cleanup by implementingDisposableBean. Introduced per-channel concurrency limits for notification tasks usingSemaphore, and ensured proper shutdown of executors. (hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/AlerterWorkerPool.java) [1] [2] [3] [4] [5]Updated
PeriodicAlertRuleSchedulerto optionally use virtual threads for periodic alert rule execution, with concurrency control viaSemaphore. Added aScheduledTaskStateclass to manage task state, support cancellation, and handle pending executions. Also implementsDisposableBeanfor proper shutdown. (hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/calculate/periodic/PeriodicAlertRuleScheduler.java) [1] [2] [3] [4] [5]Added dependency on
hertzbeat-common-springto support new thread and configuration utilities. (hertzbeat-alerter/pom.xml)Refactoring and Type Updates:
Updated usage of parameter definition types in
MonitorToolsImplto useParamDefineInfoinstead ofParamDefinefor compatibility with changes in the parameter definition API. (hertzbeat-ai/src/main/java/org/apache/hertzbeat/ai/tools/impl/MonitorToolsImpl.java) [1] [2] [3] [4]Changed construction of
MonitorDtoinMonitorToolsImplfrom builder pattern to setter methods for improved clarity and compatibility. (hertzbeat-ai/src/main/java/org/apache/hertzbeat/ai/tools/impl/MonitorToolsImpl.java)Preparatory Changes for Virtual Thread Support:
AlarmEvaluatorto prepare for virtual thread support and managed executors by importing relevant classes and interfaces. (hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/calculate/realtime/window/AlarmEvaluator.java)