From 0ddbcdd04cdded1961e253512991ce01fb285bec Mon Sep 17 00:00:00 2001 From: mingri31164 <3116430062@qq.com> Date: Wed, 29 Oct 2025 14:17:34 +0800 Subject: [PATCH] Clarify responsibilities and eliminate redundancy in BlockingQueueManager --- .../docs/user_docs/dev_manual/queue-custom.md | 10 +++-- .../user_docs/dev_manual/queue-custom.md | 10 +++-- .../support/BlockingQueueManager.java | 42 +++---------------- .../support/BlockingQueueSpiTest.java | 20 ++++----- .../core/ServerThreadPoolDynamicRefresh.java | 4 +- 5 files changed, 29 insertions(+), 57 deletions(-) diff --git a/docs/docs/user_docs/dev_manual/queue-custom.md b/docs/docs/user_docs/dev_manual/queue-custom.md index 9a879620..90575a86 100644 --- a/docs/docs/user_docs/dev_manual/queue-custom.md +++ b/docs/docs/user_docs/dev_manual/queue-custom.md @@ -51,13 +51,15 @@ com.example.queue.MyArrayBlockingQueue ### 3.1 队列创建与验证 ```java -// 创建队列 -BlockingQueue q = BlockingQueueManager.createQueue(queueType, capacity); +// 创建队列 - 使用 BlockingQueueTypeEnum +BlockingQueue q = BlockingQueueTypeEnum.createBlockingQueue(queueType, capacity); +// 或者通过队列名称创建 +BlockingQueue q2 = BlockingQueueTypeEnum.createBlockingQueue("ArrayBlockingQueue", capacity); -// 验证队列配置 +// 验证队列配置 - 使用 BlockingQueueManager boolean valid = BlockingQueueManager.validateQueueConfig(queueType, capacity); -// 动态调整容量(仅 ResizableCapacityLinkedBlockingQueue 支持) +// 动态调整容量(仅 ResizableCapacityLinkedBlockingQueue 支持)- 使用 BlockingQueueManager boolean ok = BlockingQueueManager.changeQueueCapacity(executor.getQueue(), newCapacity); ``` diff --git a/docs/i18n/zh/docusaurus-plugin-content-docs/current/user_docs/dev_manual/queue-custom.md b/docs/i18n/zh/docusaurus-plugin-content-docs/current/user_docs/dev_manual/queue-custom.md index 9a879620..90575a86 100644 --- a/docs/i18n/zh/docusaurus-plugin-content-docs/current/user_docs/dev_manual/queue-custom.md +++ b/docs/i18n/zh/docusaurus-plugin-content-docs/current/user_docs/dev_manual/queue-custom.md @@ -51,13 +51,15 @@ com.example.queue.MyArrayBlockingQueue ### 3.1 队列创建与验证 ```java -// 创建队列 -BlockingQueue q = BlockingQueueManager.createQueue(queueType, capacity); +// 创建队列 - 使用 BlockingQueueTypeEnum +BlockingQueue q = BlockingQueueTypeEnum.createBlockingQueue(queueType, capacity); +// 或者通过队列名称创建 +BlockingQueue q2 = BlockingQueueTypeEnum.createBlockingQueue("ArrayBlockingQueue", capacity); -// 验证队列配置 +// 验证队列配置 - 使用 BlockingQueueManager boolean valid = BlockingQueueManager.validateQueueConfig(queueType, capacity); -// 动态调整容量(仅 ResizableCapacityLinkedBlockingQueue 支持) +// 动态调整容量(仅 ResizableCapacityLinkedBlockingQueue 支持)- 使用 BlockingQueueManager boolean ok = BlockingQueueManager.changeQueueCapacity(executor.getQueue(), newCapacity); ``` diff --git a/infra/common/src/main/java/cn/hippo4j/common/executor/support/BlockingQueueManager.java b/infra/common/src/main/java/cn/hippo4j/common/executor/support/BlockingQueueManager.java index e0919922..99e87146 100644 --- a/infra/common/src/main/java/cn/hippo4j/common/executor/support/BlockingQueueManager.java +++ b/infra/common/src/main/java/cn/hippo4j/common/executor/support/BlockingQueueManager.java @@ -23,49 +23,17 @@ import lombok.extern.slf4j.Slf4j; import java.util.Collection; import java.util.Objects; import java.util.concurrent.BlockingQueue; -import java.util.concurrent.ThreadPoolExecutor; /** - * Blocking queue manager for queue operations. - * Supports SPI extension, queue creation, capacity management, and type recognition. + * Blocking queue runtime manager. + * Provides queue management operations: capacity adjustment, type recognition, and configuration validation. + * + *

Note: For queue creation, use {@link BlockingQueueTypeEnum#createBlockingQueue(int, Integer)} + * or {@link BlockingQueueTypeEnum#createBlockingQueue(String, Integer)} directly.

*/ @Slf4j public class BlockingQueueManager { - static { - ServiceLoaderRegistry.register(CustomBlockingQueue.class); - } - - /** - * Create blocking queue by type and capacity - * - * @param queueType queue type - * @param capacity queue capacity - * @param queue element type - * @return blocking queue instance - */ - public static BlockingQueue createQueue(Integer queueType, Integer capacity) { - if (queueType == null) { - queueType = BlockingQueueTypeEnum.LINKED_BLOCKING_QUEUE.getType(); - } - return BlockingQueueTypeEnum.createBlockingQueue(queueType, capacity); - } - - /** - * Create blocking queue by name and capacity - * - * @param queueName queue name - * @param capacity queue capacity - * @param queue element type - * @return blocking queue instance - */ - public static BlockingQueue createQueue(String queueName, Integer capacity) { - if (queueName == null || queueName.isEmpty()) { - queueName = BlockingQueueTypeEnum.LINKED_BLOCKING_QUEUE.getName(); - } - return BlockingQueueTypeEnum.createBlockingQueue(queueName, capacity); - } - /** * Check if queue capacity can be dynamically changed * diff --git a/infra/common/src/test/java/cn/hippo4j/common/executor/support/BlockingQueueSpiTest.java b/infra/common/src/test/java/cn/hippo4j/common/executor/support/BlockingQueueSpiTest.java index 15280579..428bd84c 100644 --- a/infra/common/src/test/java/cn/hippo4j/common/executor/support/BlockingQueueSpiTest.java +++ b/infra/common/src/test/java/cn/hippo4j/common/executor/support/BlockingQueueSpiTest.java @@ -143,28 +143,28 @@ public class BlockingQueueSpiTest { } /** - * Test Case 3: BlockingQueueManager unified creation entry + * Test Case 3: Queue creation via BlockingQueueTypeEnum */ @Test - public void testBlockingQueueManagerCreation() { - System.out.println("\n========== Test Case 3: BlockingQueueManager unified creation entry =========="); + public void testBlockingQueueCreation() { + System.out.println("\n========== Test Case 3: Queue creation via BlockingQueueTypeEnum =========="); - // Create built-in queue via BlockingQueueManager - BlockingQueue queue1 = BlockingQueueManager.createQueue(1, 512); + // Create built-in queue via BlockingQueueTypeEnum + BlockingQueue queue1 = BlockingQueueTypeEnum.createBlockingQueue(1, 512); Assert.assertNotNull("Should successfully create queue", queue1); Assert.assertTrue("Should create ArrayBlockingQueue", queue1 instanceof ArrayBlockingQueue); // Create by type name - BlockingQueue queue2 = BlockingQueueManager.createQueue("ArrayBlockingQueue", 1024); + BlockingQueue queue2 = BlockingQueueTypeEnum.createBlockingQueue("ArrayBlockingQueue", 1024); Assert.assertNotNull("Should successfully create queue by name", queue2); Assert.assertTrue("Should create ArrayBlockingQueue", queue2 instanceof ArrayBlockingQueue); - // Test default queue (null type) - BlockingQueue defaultQueue = BlockingQueueManager.createQueue("", 1024); - Assert.assertNotNull("Null type should create default queue", defaultQueue); + // Test default queue with null type - falls back to LinkedBlockingQueue + BlockingQueue defaultQueue = BlockingQueueTypeEnum.createBlockingQueue("", 1024); + Assert.assertNotNull("Empty name should create default queue", defaultQueue); System.out.println("Default queue type: " + defaultQueue.getClass().getSimpleName()); - System.out.println("Passed: BlockingQueueManager unified entry works"); + System.out.println("Passed: BlockingQueueTypeEnum queue creation works"); } /** diff --git a/starters/threadpool/server/src/main/java/cn/hippo4j/springboot/starter/core/ServerThreadPoolDynamicRefresh.java b/starters/threadpool/server/src/main/java/cn/hippo4j/springboot/starter/core/ServerThreadPoolDynamicRefresh.java index 909abd22..0c2669d4 100644 --- a/starters/threadpool/server/src/main/java/cn/hippo4j/springboot/starter/core/ServerThreadPoolDynamicRefresh.java +++ b/starters/threadpool/server/src/main/java/cn/hippo4j/springboot/starter/core/ServerThreadPoolDynamicRefresh.java @@ -160,9 +160,9 @@ public class ServerThreadPoolDynamicRefresh implements ThreadPoolDynamicRefresh if (BlockingQueueManager.canChangeCapacity(executor.getQueue())) { boolean success = BlockingQueueManager.changeQueueCapacity(executor.getQueue(), parameter.getCapacity()); if (success) { - log.info("Queue capacity changed to: {}", parameter.getCapacity()); + log.info("Queue capacity changed to: {} for thread pool: {}", parameter.getCapacity(), parameter.getTpId()); } else { - log.warn("Failed to change queue capacity to: {}", parameter.getCapacity()); + log.warn("Failed to change queue capacity to: {} for thread pool: {}", parameter.getCapacity(), parameter.getTpId()); } } else { log.warn("Queue capacity cannot be changed for current queue type: {}. " +