如何关联CompletableFuture任务与线程池,监控提交到执行的空闲时长?
问题
我需要收集异步任务从提交到执行的等待时长,用来监控基于fixedThreadPool实现的ExecutorService和提交的任务。当前用CompletableFuture组合异步任务,要实现两个目标:
- 给
CompletableFuture.supplyAsync()传入的Supplier添加标识,比如任务名称 - 监控线程池中
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
相关产品推荐
相关产品推荐

