任务分配系统听起来是个很大的概念,但拆开看其实很简单:有一批任务需要处理,有若干个工作线程可以干活,中间需要一个调度器把任务合理地分发给各个线程。生产环境里常用的MQ、Elastic-Job、XXL-Job本质上都是这个模型的复杂版。今天我们就用Java自带的线程池和阻塞队列,从零实现一个简易任务分配系统,把里面的设计思路和容易踩的坑都讲清楚。

一、整体架构设计:三层模型怎么划分
在动手写代码之前,先想清楚系统的分层。一个简易任务分配系统通常分为三层:任务提交层、调度层和执行层。任务提交层负责接收外部请求,把任务封装成统一的对象放入队列;调度层是核心,维护一个阻塞队列和若干工作线程,决定任务什么时候被取出、分给谁;执行层就是各个工作线程本身,它们循环从队列中取任务并执行。
这个划分的好处是职责清晰。如果以后想扩展成分布式版本,只需要把阻塞队列换成Redis或消息队列,任务提交层和执行层几乎不用动。很多初学者写这种系统时喜欢把所有逻辑塞进一个类,结果任务类型一变就要大改,这就是没有做好抽象的代价。
定义任务时,建议抽象出统一接口,包含任务ID、任务类型、执行方法和超时时间这几个基本属性。下面是任务模型的基础定义:
public interface Task {
String getTaskId();
String getTaskType();
void execute() throws Exception;
long getTimeoutMillis();
}
public class SimpleTask implements Task {
private final String taskId;
private final Runnable action;
public SimpleTask(String taskId, Runnable action) {
this.taskId = taskId;
this.action = action;
}
@Override
public String getTaskId() { return taskId; }
@Override
public String getTaskType() { return "default"; }
@Override
public void execute() { action.run(); }
@Override
public long getTimeoutMillis() { return 30_000; }
}这个抽象看似简单,但它是整个系统可扩展性的基础。后续要加优先级任务、延时任务,都只需要新增实现类,不用改调度器代码。
二、任务分发核心:BlockingQueue的使用与线程池配置
调度层的核心是BlockingQueue。工作线程调用take()方法从队列取任务,队列空了就自动阻塞挂起,不消耗CPU;有新任务进来时被唤醒继续工作。这种生产者消费者模式天然解决了任务堆积和线程空转的问题。
选择队列实现时有讲究。LinkedBlockingQueue可以设置容量也可以不设,不设容量时默认是近乎无界的,任务堆积可能导致内存溢出,生产环境建议显式设置容量;ArrayBlockingQueue容量固定,公平模式可选;如果要支持优先级调度,就用PriorityBlockingQueue,但要注意任务必须实现Comparable接口。
下面是一个支持优先级的任务调度器实现:
public class TaskScheduler {
private final PriorityBlockingQueue<PriorityTask> queue =
new PriorityBlockingQueue<>(100);
private final ExecutorService workerPool;
private final AtomicInteger submittedCount = new AtomicInteger(0);
private volatile boolean running = true;
public TaskScheduler(int workerCount) {
this.workerPool = Executors.newFixedThreadPool(workerCount);
for (int i = 0; i < workerCount; i++) {
workerPool.submit(this::workerLoop);
}
}
public void submit(PriorityTask task) {
if (!running) {
throw new IllegalStateException("调度器已关闭");
}
queue.put(task);
submittedCount.incrementAndGet();
}
private void workerLoop() {
while (running || !queue.isEmpty()) {
try {
Task task = queue.poll(1, TimeUnit.SECONDS);
if (task != null) {
executeSafely(task);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
}
private void executeSafely(Task task) {
try {
task.execute();
} catch (Exception e) {
System.err.println("任务执行失败: " + task.getTaskId()
+ ", 原因: " + e.getMessage());
}
}
public void shutdown() {
running = false;
workerPool.shutdown();
}
}注意这里没有直接用ExecutorService的submit去执行任务,而是自己维护了工作循环。这样做的好处是我们可以在循环里插入监控逻辑、重试逻辑,对任务的执行过程有完全的掌控力。如果只是简单场景,直接用ThreadPoolExecutor也是完全可以的,代码量更少。
工作线程数不是越多越好。CPU密集型任务建议设为CPU核心数加一,IO密集型任务可以设为CPU核心数的两倍左右。这个经验值不一定精确,但比拍脑袋定一个50要靠谱得多。
三、任务状态管理与线程安全:最容易踩的坑
任务分发出去之后,状态管理是个大问题。一个任务通常有几种状态:待执行、执行中、成功、失败、超时。如果用普通的int字段加if判断来改状态,在多线程环境下几乎必然出问题——两个工作线程同时判断任务处于待执行状态,然后同时去执行,任务就被重复处理了。
解决这类问题有两把钥匙。第一把是AtomicReference配合CAS操作,状态变更必须原子完成,谁CAS成功谁才有执行权:
public enum TaskStatus { PENDING, RUNNING, SUCCESS, FAILED }
public class TaskStatusManager {
private final AtomicReference<TaskStatus> status =
new AtomicReference<>(TaskStatus.PENDING);
public boolean tryStart() {
// 只有PENDING状态才能改为RUNNING,CAS失败说明被别的线程抢了
return status.compareAndSet(TaskStatus.PENDING, TaskStatus.RUNNING);
}
public void finish(boolean success) {
status.set(success ? TaskStatus.SUCCESS : TaskStatus.FAILED);
}
public TaskStatus get() {
return status.get();
}
}第二把钥匙是状态机约束。合法的状态流转只能是PENDING到RUNNING,RUNNING到SUCCESS或FAILED,其他任何流转都应该被拒绝。把流转规则收敛到一处代码,比散落在各处的if判断安全得多。
还有一个常见坑是工作线程里的异常吞噬。如果任务执行抛出了RuntimeException而没有任何捕获,使用ExecutorService.submit时异常会被存在Future里悄悄吞掉,你只会看到任务莫名消失。所以上面的executeSafely方法里做了统一捕获,实际项目中还应该把失败任务记录下来,甚至放回队列做有限次重试。重试时注意带上重试次数上限,否则一个持续失败的任务会把系统拖入死循环。
最后说说优雅停机。系统关闭时不能直接杀线程,否则正在执行的任务会中断在半路。正确做法是设置停止标志、停止接收新任务、等待队列中剩余任务执行完(或超时后强制退出),ThreadPoolExecutor的shutdown和awaitTermination组合就是干这个的。虽然是简易系统,这些细节做扎实了,离生产可用的距离就没那么远了。
Java任务分配系统多线程任务调度任务分配项目实战修改时间:2026-09-09 05:58:33