验证码: 看不清楚,换一张 查询 注册会员,免验证
  • {{ basic.site_slogan }}
  • 打开微信扫一扫,
    您还可以在这里找到我们哟

    关注我们

如何扩展ExecutorService

阅读:399 来源:乙速云 作者:代码code

如何扩展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> tasks) throws InterruptedException {
        return delegate.invokeAll(tasks);
    }

    @Override
    public  List> invokeAll(Collection> tasks, long timeout, TimeUnit unit) throws InterruptedException {
        return delegate.invokeAll(tasks, timeout, unit);
    }

    @Override
    public  T invokeAny(Collection> tasks) throws InterruptedException, ExecutionException {
        return delegate.invokeAny(tasks);
    }

    @Override
    public  T invokeAny(Collection> 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、用户信息)

如果你愿意,我可以:

  • ✅ 给你一个生产级可观测线程池
  • ✅ 写一个支持优雅关闭 + 拒绝策略扩展的线程池
  • ✅ 对比 ExecutorService vs ForkJoinPool

你可以直接说你的使用场景

分享到:
*特别声明:以上内容来自于网络收集,著作权属原作者所有,如有侵权,请联系我们: hlamps#outlook.com (#换成@)。
相关文章
{{ v.title }}
{{ v.description||(cleanHtml(v.content)).substr(0,100)+'···' }}
你可能感兴趣
推荐阅读 更多>