如何扩展ExecutorService
在 Java 中,ExecutorService 本身是一个接口,不能直接“扩展”它的功能,但可以通过实现接口、包装(装饰)、继承已有实现类、或结合其他并发工具来“扩展”它。下面从常见方式与实战场景给你一个系统说明。
一、最推荐方式:装饰器模式(Wrapper)
这是最常用、最安全的扩展方式,不破坏线程池原有行为。
示例:给 ExecutorService 增加“任务提交日志”
public class LoggingExecutorService implements ExecutorService {
private final ExecutorService delegate;
public LoggingExecutorService(ExecutorService delegate) {
this.delegate = delegate;
}
@Override
public Future submit(Callable task) {
System.out.println("Submit callable: " + task);
return delegate.submit(task);
}
@Override
public Future> submit(Runnable task) {
System.out.println("Submit runnable: " + task);
return delegate.submit(task);
}
// 其余方法全部委托给 delegate
@Override
public void shutdown() {
delegate.shutdown();
}
@Override
public List shutdownNow() {
return delegate.shutdownNow();
}
@Override
public boolean isShutdown() {
return delegate.isShutdown();
}
@Override
public boolean isTerminated() {
return delegate.isTerminated();
}
@Override
public boolean awaitTermination(long timeout, TimeUnit unit) throws InterruptedException {
return delegate.awaitTermination(timeout, unit);
}
@Override
public List> invokeAll(Collection extends Callable> tasks) throws InterruptedException {
return delegate.invokeAll(tasks);
}
@Override
public List> invokeAll(Collection extends Callable> tasks, long timeout, TimeUnit unit) throws InterruptedException {
return delegate.invokeAll(tasks, timeout, unit);
}
@Override
public T invokeAny(Collection extends Callable> tasks) throws InterruptedException, ExecutionException {
return delegate.invokeAny(tasks);
}
@Override
public T invokeAny(Collection extends Callable> tasks, long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
return delegate.invokeAny(tasks, timeout, unit);
}
@Override
public void execute(Runnable command) {
delegate.execute(command);
}
}
✅ 优点
- 不破坏原有线程池
- 可叠加多个功能(日志、监控、限流)
- 对调用方透明
二、继承 ThreadPoolExecutor(适合高级定制)
如果你用的是 ThreadPoolExecutor,可以直接继承它。
常见扩展点
| 方法 | 用途 |
|---|---|
beforeExecute |
任务执行前 |
afterExecute |
任务执行后 |
terminated |
线程池关闭后 |
示例:统计任务执行时间
public class TimingThreadPool extends ThreadPoolExecutor {
public TimingThreadPool(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue workQueue) {
super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue);
}
@Override
protected void beforeExecute(Thread t, Runnable r) {
super.beforeExecute(t, r);
System.out.println("Task start: " + r);
}
@Override
protected void afterExecute(Runnable r, Throwable t) {
super.afterExecute(r, t);
System.out.println("Task end: " + r);
}
}
✅ 适合:
- 监控
- 统计
- 异常统一处理
⚠️ 注意:
- 只能用于
ThreadPoolExecutor - 不能用于
ForkJoinPool
三、扩展任务(Runnable / Callable)
如果你只是想“增强任务”,而不是线程池本身:
public class LoggingTask implements Runnable {
private final Runnable task;
public LoggingTask(Runnable task) {
this.task = task;
}
@Override
public void run() {
long start = System.currentTimeMillis();
try {
task.run();
} finally {
System.out.println("Cost: " + (System.currentTimeMillis() - start));
}
}
}
使用:
executor.submit(new LoggingTask(() -> System.out.println("hello")));
✅ 最轻量
❌ 侵入提交代码
四、结合 CompletableFuture(现代方式)
如果你使用的是 Java 8+,可以考虑不扩展 ExecutorService,而是增强使用方式:
CompletableFuture.supplyAsync(() -> {
return doWork();
}, executor)
.thenApply(...)
.exceptionally(...)
✅ 非阻塞
✅ 链式扩展
✅ 更灵活
五、扩展场景总结
| 需求 | 推荐方式 |
|---|---|
| 增加日志 / 监控 | 装饰器 |
| 控制任务执行行为 | 继承 ThreadPoolExecutor |
| 统计 / 超时 / 重试 | 包装任务 |
| 异步编排 | CompletableFuture |
| 限流 / 熔断 | 外层封装 |
六、典型扩展案例
- 任务提交限流(Semaphore)
- 提交失败自动重试
- 线程池指标上报(QPS、队列长度)
- 上下文传递(TraceId、用户信息)
如果你愿意,我可以:
- ✅ 给你一个生产级可观测线程池
- ✅ 写一个支持优雅关闭 + 拒绝策略扩展的线程池
- ✅ 对比
ExecutorServicevsForkJoinPool
你可以直接说你的使用场景。