You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何关联CompletableFuture任务与线程池,监控提交到执行的空闲时长?

问题

我需要收集异步任务从提交到执行的等待时长,用来监控基于fixedThreadPool实现的ExecutorService和提交的任务。当前用CompletableFuture组合异步任务,要实现两个目标:

  1. 给CompletableFuture.supplyAsync()传入的Supplier添加标识,比如任务名称
  2. 监控线程池中Runnable从提交到执行的等待时间,通过扩展ThreadPoolExecutor,在beforeExecute()和execute()里加日志

现在的核心难点是关联任务标识和线程池的监控数据:线程池处理的是Runnable,但CompletableFuture内部会把Supplier包装成Runnable,没法直接从这个Runnable拿到任务标识;而且CompletableFuture组合后的任务是内部调用execute(),没法手动给每个Runnable加标识。

示例代码:

List<Foo> results = createResults();
results.forEach(r -> CompletableFuture.completedFuture(r)
  .thenCompose(result -> addSomething(result, something)
  .thenCombine(addSomethingElse(result, somethingElse), (r1, r2) -> result)
  .thenCompose(r -> doSomething(result).thenCompose(this::setSomething))
  .thenApply(v -> result))));

// 后续调用join等待完成
listOfFutures.join()

其中内部方法示例:

private CompletableFuture<Foo> setSomething(Foo foo) {
    return CompletableFuture.supplyAsync(() -> {
        foo.description = "Setting something";
        return foo;
    }, myExecutorService);
}
解决方案

核心思路是给Supplier包装一层带标识的包装类,同时让线程池能识别这个包装类对应的Runnable,从而关联任务标识和监控数据。

1. 实现带标识的Supplier包装类

定义一个带任务ID/名称的Supplier包装类,同时给生成的Runnable注入标识,让线程池能识别:

import java.util.concurrent.Supplier;

public class IdentifiedSupplier<T> implements Supplier<T> {
    private final String taskId;
    private final Supplier<T> delegate;

    public IdentifiedSupplier(String taskId, Supplier<T> delegate) {
        this.taskId = taskId;
        this.delegate = delegate;
    }

    @Override
    public T get() {
        return delegate.get();
    }

    // 获取任务标识
    public String getTaskId() {
        return taskId;
    }

    // 把带标识的Supplier包装成线程池可识别的Runnable
    public static <T> Runnable wrapAsRunnable(IdentifiedSupplier<T> supplier) {
        return new Runnable() {
            @Override
            public void run() {
                supplier.get();
            }

            // 暴露任务标识的方法,供线程池提取
            public String getTaskId() {
                return supplier.getTaskId();
            }
        };
    }
}

2. 封装自定义CompletableFuture工具方法

不用重写CompletableFuture,直接封装工具方法,自动给Supplier加标识,提交带标识的Runnable到线程池:

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;
import java.util.concurrent.Supplier;

public class TrackableCompletableFuture {
    public static <T> CompletableFuture<T> supplyAsync(String taskId, Supplier<T> supplier, Executor executor) {
        IdentifiedSupplier<T> identifiedSupplier = new IdentifiedSupplier<>(taskId, supplier);
        // 提交我们包装后的Runnable,确保线程池能拿到标识
        CompletableFuture<T> future = new CompletableFuture<>();
        executor.execute(() -> {
            try {
                T result = identifiedSupplier.get();
                future.complete(result);
            } catch (Throwable t) {
                future.completeExceptionally(t);
            }
        });
        return future;
    }
}

3. 扩展ThreadPoolExecutor实现监控与关联

自定义线程池,在execute()记录任务提交时间,beforeExecute()计算等待时长并关联任务标识:

import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;

public class MonitoredThreadPoolExecutor extends ThreadPoolExecutor {
    // 存储任务标识与提交时间的映射
    private final ConcurrentHashMap<String, Long> taskSubmitTimeMap = new ConcurrentHashMap<>();
    // 自动生成任务ID的生成器(用于未自定义标识的任务)
    private final AtomicLong taskIdGenerator = new AtomicLong(0);

    public MonitoredThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, BlockingQueue<Runnable> workQueue) {
        super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue);
    }

    @Override
    public void execute(Runnable command) {
        String taskId = null;
        // 尝试从Runnable中提取任务标识
        try {
            taskId = (String) command.getClass().getMethod("getTaskId").invoke(command);
        } catch (Exception e) {
            // 无自定义标识时自动生成
            taskId = "auto-task-" + taskIdGenerator.incrementAndGet();
        }
        // 记录任务提交时间
        taskSubmitTimeMap.put(taskId, System.currentTimeMillis());
        super.execute(command);
    }

    @Override
    protected void beforeExecute(Thread t, Runnable r) {
        super.beforeExecute(t, r);
        String taskId = null;
        // 再次提取任务标识
        try {
            taskId = (String) r.getClass().getMethod("getTaskId").invoke(r);
        } catch (Exception e) {
            taskId = "unknown-task";
        }
        // 计算并输出等待时长
        Long submitTime = taskSubmitTimeMap.remove(taskId);
        if (submitTime != null) {
            long waitTime = System.currentTimeMillis() - submitTime;
            // 这里可以替换为上报监控系统的逻辑
            System.out.printf("任务[%s] | 等待时长: %d ms | 执行线程: %s%n", taskId, waitTime, t.getName());
        }
    }
}

4. 改造业务代码使用自定义工具

把原来的CompletableFuture.supplyAsync()替换为自定义工具方法,传入任务标识:

private CompletableFuture<Foo> setSomething(Foo foo) {
    // 用业务唯一标识作为任务ID,比如foo的ID
    String taskId = "set-something-" + foo.getId();
    return TrackableCompletableFuture.supplyAsync(taskId, () -> {
        foo.description = "Setting something";
        return foo;
    }, myExecutorService);
}

注意事项

  • 对于CompletableFuture链式调用(如thenCompose、thenApply),默认会用ForkJoinPool.commonPool(),如果需要监控这些后续任务,要使用带线程池参数的异步方法(如thenApplyAsync(..., myExecutorService))
  • taskSubmitTimeMap在beforeExecute中会移除已执行任务的记录,不会造成内存泄漏
  • 可以扩展IdentifiedSupplier添加更多元数据(如业务类型、优先级),满足更复杂的监控需求

内容的提问来源于stack exchange,提问作者LisekKL

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 14:15:49