在 Java 中,ExecutorService 本身是一个接口,不能直接“扩展”它的功能,但可以通过实现接口、包装(装饰)、继承已有实现类、或结合其他并发工具来“扩展”它。下面从常见方式与实战场景给你一个系统说明。
这是最常用、最安全的扩展方式,不破坏线程池原有行为。
public class LoggingExecutorService implements ExecutorService {
private final ExecutorService delegate;
public LoggingExecutorService(ExecutorService delegate) {
this.delegate = delegate;
}
@Override
public <T> Future<T> submit(Callable<T> 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<Runnable> 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 <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) throws InterruptedException {
return delegate.invokeAll(tasks);
}
@Override
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException {
return delegate.invokeAll(tasks, timeout, unit);
}
@Override
public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException {
return delegate.invokeAny(tasks);
}
@Override
public <T> T invokeAny(Collection<? extends Callable<T>> 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,可以直接继承它。
| 方法 | 用途 |
|---|---|
beforeExecute |
任务执行前 |
afterExecute |
任务执行后 |
terminated |
线程池关闭后 |
public class TimingThreadPool extends ThreadPoolExecutor {
public TimingThreadPool(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> 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);
}
}
✅ 适合:
⚠️ 注意:
ThreadPoolExecutorForkJoinPool如果你只是想“增强任务”,而不是线程池本身:
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")));
✅ 最轻量
❌ 侵入提交代码
如果你使用的是 Java 8+,可以考虑不扩展 ExecutorService,而是增强使用方式:
CompletableFuture.supplyAsync(() -> {
return doWork();
}, executor)
.thenApply(...)
.exceptionally(...)
✅ 非阻塞
✅ 链式扩展
✅ 更灵活
| 需求 | 推荐方式 |
|---|---|
| 增加日志 / 监控 | 装饰器 |
| 控制任务执行行为 | 继承 ThreadPoolExecutor |
| 统计 / 超时 / 重试 | 包装任务 |
| 异步编排 | CompletableFuture |
| 限流 / 熔断 | 外层封装 |
如果你愿意,我可以:
ExecutorService vs ForkJoinPool你可以直接说你的使用场景。
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。