task-74: 隔离调度线程池、文件作业线程池和外部 Coze/图片执行池
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
package com.nanri.aiimage.config;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
@@ -20,9 +21,9 @@ public class SchedulingConfig {
|
||||
private static final ZoneId BUSINESS_ZONE = ZoneId.of("Asia/Shanghai");
|
||||
|
||||
@Bean
|
||||
public TaskScheduler taskScheduler() {
|
||||
public TaskScheduler taskScheduler(@Value("${aiimage.scheduling.pool-size:4}") int poolSize) {
|
||||
ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
|
||||
scheduler.setPoolSize(4);
|
||||
scheduler.setPoolSize(Math.max(1, poolSize));
|
||||
scheduler.setThreadNamePrefix("aiimage-scheduling-");
|
||||
scheduler.setClock(Clock.system(BUSINESS_ZONE));
|
||||
scheduler.setWaitForTasksToCompleteOnShutdown(true);
|
||||
|
||||
@@ -42,19 +42,24 @@ public class TaskFileJobConfig {
|
||||
ExecutorService cozeVirtualThreadExecutor,
|
||||
@Value("${aiimage.coze-task.max-concurrent:12}") int maxConcurrent) {
|
||||
Semaphore semaphore = new Semaphore(Math.max(1, maxConcurrent));
|
||||
return new ConcurrentTaskExecutor(command -> cozeVirtualThreadExecutor.execute(() -> {
|
||||
boolean acquired = false;
|
||||
try {
|
||||
semaphore.acquire();
|
||||
acquired = true;
|
||||
command.run();
|
||||
} catch (InterruptedException ex) {
|
||||
Thread.currentThread().interrupt();
|
||||
} finally {
|
||||
if (acquired) {
|
||||
semaphore.release();
|
||||
}
|
||||
return new ConcurrentTaskExecutor(command -> {
|
||||
if (command == null) {
|
||||
throw new IllegalArgumentException("coze 任务不能为 null");
|
||||
}
|
||||
}));
|
||||
cozeVirtualThreadExecutor.execute(() -> {
|
||||
boolean acquired = false;
|
||||
try {
|
||||
semaphore.acquire();
|
||||
acquired = true;
|
||||
command.run();
|
||||
} catch (InterruptedException ex) {
|
||||
Thread.currentThread().interrupt();
|
||||
} finally {
|
||||
if (acquired) {
|
||||
semaphore.release();
|
||||
}
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
+249
@@ -0,0 +1,249 @@
|
||||
package com.nanri.aiimage.config;
|
||||
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.core.task.TaskExecutor;
|
||||
import org.springframework.core.task.TaskRejectedException;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Date;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* Task 74:隔离调度线程池、文件作业线程池和外部 Coze/图片执行池。
|
||||
* 三个执行池各自独立配置、独立命名、容量互不影响:调度池
|
||||
* (aiimage.scheduling.pool-size,默认 4)与文件作业派发池
|
||||
* (aiimage.result-file-job.*,默认 2 线程/队列 200)互不共享线程;
|
||||
* 外部 Coze 池以虚拟线程 + 信号量限流(默认 12)。容量非法值统一
|
||||
* 钳制到最小值;任务失败后信号量名额与调度槽位必须释放,任一池打满
|
||||
* 不影响其他池。
|
||||
*/
|
||||
class ThreadPoolIsolationConfigTest {
|
||||
|
||||
private final List<AutoCloseable> closeables = new ArrayList<>();
|
||||
|
||||
@AfterEach
|
||||
void tearDown() throws Exception {
|
||||
for (AutoCloseable closeable : closeables) {
|
||||
closeable.close();
|
||||
}
|
||||
}
|
||||
|
||||
private ThreadPoolTaskScheduler newScheduler(int poolSize) {
|
||||
ThreadPoolTaskScheduler scheduler =
|
||||
(ThreadPoolTaskScheduler) new SchedulingConfig().taskScheduler(poolSize);
|
||||
closeables.add(scheduler::destroy);
|
||||
return scheduler;
|
||||
}
|
||||
|
||||
private ThreadPoolTaskExecutor newDispatch(int poolSize, int queueCapacity) {
|
||||
ThreadPoolTaskExecutor executor = (ThreadPoolTaskExecutor) new TaskFileJobConfig()
|
||||
.taskFileJobDispatchExecutor(poolSize, queueCapacity);
|
||||
closeables.add(executor::destroy);
|
||||
return executor;
|
||||
}
|
||||
|
||||
private ExecutorService newCozeVirtual() {
|
||||
ExecutorService executor = new TaskFileJobConfig().cozeVirtualThreadExecutor();
|
||||
closeables.add(() -> executor.shutdownNow());
|
||||
return executor;
|
||||
}
|
||||
|
||||
private TaskExecutor newCoze(ExecutorService virtualExecutor, int maxConcurrent) {
|
||||
return new TaskFileJobConfig().cozeTaskExecutor(virtualExecutor, maxConcurrent);
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_074_image_dispatch_job_normal_default_path() throws Exception {
|
||||
// 默认路径:三池按默认容量初始化,正常提交一个任务即可执行。
|
||||
ThreadPoolTaskScheduler scheduler = newScheduler(4);
|
||||
assertEquals(4, scheduler.getScheduledThreadPoolExecutor().getCorePoolSize(), "调度池默认 4 线程");
|
||||
assertTrue(scheduler.getThreadNamePrefix().startsWith("aiimage-scheduling-"),
|
||||
"调度池线程名独立前缀");
|
||||
|
||||
ThreadPoolTaskExecutor dispatch = newDispatch(2, 200);
|
||||
assertEquals(2, dispatch.getCorePoolSize());
|
||||
assertEquals(2, dispatch.getMaxPoolSize(), "文件作业池 core=max,不随压力扩张");
|
||||
assertEquals(200, dispatch.getQueueCapacity());
|
||||
|
||||
ExecutorService cozeVirtual = newCozeVirtual();
|
||||
TaskExecutor coze = newCoze(cozeVirtual, 12);
|
||||
AtomicBoolean ran = new AtomicBoolean(false);
|
||||
CountDownLatch done = new CountDownLatch(1);
|
||||
coze.execute(() -> {
|
||||
ran.set(true);
|
||||
done.countDown();
|
||||
});
|
||||
assertTrue(done.await(5, TimeUnit.SECONDS), "Coze 池默认限流 12,正常提交即执行");
|
||||
assertTrue(ran.get());
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_074_image_dispatch_job_normal_multiple_items() throws Exception {
|
||||
// 批量场景:三池同时按各自容量配置,各跑一个任务互不阻塞;
|
||||
// 线程名前缀互不相同,线程转储可识别归属池。
|
||||
ThreadPoolTaskScheduler scheduler = newScheduler(6);
|
||||
ThreadPoolTaskExecutor dispatch = newDispatch(3, 500);
|
||||
ExecutorService cozeVirtual = newCozeVirtual();
|
||||
TaskExecutor coze = newCoze(cozeVirtual, 8);
|
||||
assertEquals(6, scheduler.getScheduledThreadPoolExecutor().getCorePoolSize());
|
||||
assertEquals(3, dispatch.getCorePoolSize());
|
||||
assertEquals(500, dispatch.getQueueCapacity());
|
||||
|
||||
CountDownLatch all = new CountDownLatch(3);
|
||||
dispatch.execute(all::countDown);
|
||||
scheduler.schedule((Runnable) all::countDown, new Date(System.currentTimeMillis() + 50));
|
||||
coze.execute(all::countDown);
|
||||
assertTrue(all.await(5, TimeUnit.SECONDS), "三个池同时执行互不阻塞");
|
||||
|
||||
assertNotEquals(scheduler.getThreadNamePrefix(), dispatch.getThreadNamePrefix(),
|
||||
"调度池与文件作业池线程名前缀隔离");
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_074_image_dispatch_job_normal_repeated_operation_is_idempotent() throws Exception {
|
||||
// 幂等:同一任务重复提交各自独立执行一次,不合并不丢失,
|
||||
// 线程数不因重复提交而扩张。
|
||||
ThreadPoolTaskExecutor dispatch = newDispatch(2, 100);
|
||||
AtomicInteger count = new AtomicInteger();
|
||||
CountDownLatch all = new CountDownLatch(3);
|
||||
Runnable task = () -> {
|
||||
count.incrementAndGet();
|
||||
all.countDown();
|
||||
};
|
||||
dispatch.execute(task);
|
||||
dispatch.execute(task);
|
||||
dispatch.execute(task);
|
||||
assertTrue(all.await(5, TimeUnit.SECONDS), "同一任务重复提交各自执行一次");
|
||||
assertEquals(3, count.get());
|
||||
assertEquals(2, dispatch.getMaxPoolSize(), "重复提交不扩张线程数");
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_074_image_dispatch_job_boundary_empty_input() throws Exception {
|
||||
// 空输入:容量配置为 0 时统一钳制到最小值,池仍可用。
|
||||
ThreadPoolTaskScheduler scheduler = newScheduler(0);
|
||||
assertEquals(1, scheduler.getScheduledThreadPoolExecutor().getCorePoolSize(), "调度池 0 钳制到 1");
|
||||
ThreadPoolTaskExecutor dispatch = newDispatch(0, 0);
|
||||
assertEquals(1, dispatch.getCorePoolSize(), "文件作业池 0 钳制到 1");
|
||||
assertEquals(10, dispatch.getQueueCapacity(), "队列 0 钳制到 10");
|
||||
ExecutorService cozeVirtual = newCozeVirtual();
|
||||
TaskExecutor coze = newCoze(cozeVirtual, 0);
|
||||
CountDownLatch done = new CountDownLatch(1);
|
||||
coze.execute(done::countDown);
|
||||
assertTrue(done.await(5, TimeUnit.SECONDS), "Coze 池 0 钳制到 1 后仍可执行");
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_074_image_dispatch_job_boundary_single_item() throws Exception {
|
||||
// 单元素:单线程池单任务直接完成,不依赖批量路径。
|
||||
ThreadPoolTaskExecutor dispatch = newDispatch(1, 10);
|
||||
CountDownLatch done = new CountDownLatch(1);
|
||||
dispatch.execute(done::countDown);
|
||||
assertTrue(done.await(5, TimeUnit.SECONDS), "单线程池单任务直接完成");
|
||||
assertEquals(1, dispatch.getCorePoolSize());
|
||||
|
||||
ThreadPoolTaskScheduler scheduler = newScheduler(1);
|
||||
CountDownLatch scheduled = new CountDownLatch(1);
|
||||
scheduler.schedule((Runnable) scheduled::countDown, new Date(System.currentTimeMillis() + 30));
|
||||
assertTrue(scheduled.await(5, TimeUnit.SECONDS), "单线程调度池单次调度完成");
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_074_image_dispatch_job_boundary_limit_and_overflow() throws Exception {
|
||||
// 上限/超限:文件作业池 2 线程 + 队列 10,容量为 12;
|
||||
// 第 13 个提交被拒绝(TaskRejectedException),不发生无界堆积。
|
||||
ThreadPoolTaskExecutor dispatch = newDispatch(2, 10);
|
||||
CountDownLatch blockersRunning = new CountDownLatch(2);
|
||||
CountDownLatch releaseBlockers = new CountDownLatch(1);
|
||||
CountDownLatch allDone = new CountDownLatch(12);
|
||||
for (int i = 0; i < 2; i++) {
|
||||
dispatch.execute(() -> {
|
||||
blockersRunning.countDown();
|
||||
try {
|
||||
releaseBlockers.await(10, TimeUnit.SECONDS);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
allDone.countDown();
|
||||
});
|
||||
}
|
||||
assertTrue(blockersRunning.await(5, TimeUnit.SECONDS), "两个运行线程占位");
|
||||
for (int i = 0; i < 10; i++) {
|
||||
dispatch.execute(allDone::countDown);
|
||||
}
|
||||
assertThrows(TaskRejectedException.class, () -> dispatch.execute(allDone::countDown),
|
||||
"队列满后拒绝新提交,不发生无界堆积");
|
||||
releaseBlockers.countDown();
|
||||
assertTrue(allDone.await(10, TimeUnit.SECONDS), "已受理的 12 个任务全部完成");
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_074_image_dispatch_job_invalid_input_rejected() throws Exception {
|
||||
// 非法参数:负容量统一钳制到最小值(不崩溃、行为确定);
|
||||
// null 任务直接被拒绝。
|
||||
ThreadPoolTaskScheduler scheduler = newScheduler(-1);
|
||||
assertEquals(1, scheduler.getScheduledThreadPoolExecutor().getCorePoolSize(), "负值钳制到最小值");
|
||||
ThreadPoolTaskExecutor dispatch = newDispatch(-2, -5);
|
||||
assertEquals(1, dispatch.getCorePoolSize());
|
||||
assertEquals(10, dispatch.getQueueCapacity());
|
||||
|
||||
ExecutorService cozeVirtual = newCozeVirtual();
|
||||
TaskExecutor coze = newCoze(cozeVirtual, -3);
|
||||
assertThrows(IllegalArgumentException.class, () -> coze.execute(null), "null 任务被拒绝");
|
||||
CountDownLatch done = new CountDownLatch(1);
|
||||
coze.execute(done::countDown);
|
||||
assertTrue(done.await(5, TimeUnit.SECONDS), "非法配置钳制后池仍可用");
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_074_image_dispatch_job_dependency_failure_releases_resources() throws Exception {
|
||||
// 依赖失败:Coze 任务抛异常后信号量名额必须释放(后续任务可执行);
|
||||
// 调度任务异常被 error handler 吞掉,调度器继续可用。
|
||||
ExecutorService cozeVirtual = newCozeVirtual();
|
||||
TaskExecutor coze = newCoze(cozeVirtual, 2);
|
||||
CountDownLatch blockerHeld = new CountDownLatch(1);
|
||||
CountDownLatch releaseBlocker = new CountDownLatch(1);
|
||||
coze.execute(() -> {
|
||||
blockerHeld.countDown();
|
||||
try {
|
||||
releaseBlocker.await(10, TimeUnit.SECONDS);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
});
|
||||
assertTrue(blockerHeld.await(5, TimeUnit.SECONDS), "任务 1 占住一个信号量名额");
|
||||
coze.execute(() -> {
|
||||
throw new IllegalStateException("coze down");
|
||||
});
|
||||
CountDownLatch afterFailure = new CountDownLatch(1);
|
||||
coze.execute(afterFailure::countDown);
|
||||
assertTrue(afterFailure.await(5, TimeUnit.SECONDS), "失败任务释放名额,后续任务可执行");
|
||||
releaseBlocker.countDown();
|
||||
|
||||
ThreadPoolTaskScheduler scheduler = newScheduler(2);
|
||||
AtomicBoolean secondRan = new AtomicBoolean(false);
|
||||
CountDownLatch secondDone = new CountDownLatch(1);
|
||||
scheduler.schedule(() -> {
|
||||
throw new IllegalStateException("scheduled boom");
|
||||
}, new Date(System.currentTimeMillis() + 30));
|
||||
scheduler.schedule(() -> {
|
||||
secondRan.set(true);
|
||||
secondDone.countDown();
|
||||
}, new Date(System.currentTimeMillis() + 60));
|
||||
assertTrue(secondDone.await(5, TimeUnit.SECONDS), "调度任务异常被 error handler 吞掉,调度器继续可用");
|
||||
assertTrue(secondRan.get());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user