如何在使用Hazelcast Durable ExecutorService提交Callable/Runnable任务时获取开始、调度和结束时间戳?
Hazelcast Durable ExecutorService 获取任务时间戳的实现方案
Hazelcast Durable ExecutorService 本身没有直接提供获取任务提交、调度、执行时间戳的API,你可以通过以下两种常用方式实现需求:
方法一:包装任务,手动埋点记录时间
通过自定义包装类包裹原任务,在任务执行的关键节点记录时间,适合需要在任务内部感知时间的场景。
Callable任务包装示例
import java.util.concurrent.Callable; import java.util.HashMap; import java.util.Map; public class TimedCallableWrapper<V> implements Callable<V> { private final Callable<V> delegate; private final Map<String, Long> timestamps = new HashMap<>(); public TimedCallableWrapper(Callable<V> delegate) { this.delegate = delegate; // 记录任务提交时间(初始化时即提交时刻) timestamps.put("submittedTime", System.currentTimeMillis()); } @Override public V call() throws Exception { // 记录任务实际开始执行时间 timestamps.put("startedTime", System.currentTimeMillis()); try { return delegate.call(); } finally { // 记录任务执行结束时间 timestamps.put("completedTime", System.currentTimeMillis()); } } // 对外暴露时间戳获取方法 public Map<String, Long> getTimestamps() { return timestamps; } }
使用方式
// 用包装类包裹原Callable任务 TimedCallableWrapper<String> timedTask = new TimedCallableWrapper<>(yourOriginalCallable); // 提交到Durable ExecutorService DurableExecutorFuture<String> future = durableExecutor.submit(timedTask); // 任务完成后获取时间戳 Map<String, Long> timestamps = timedTask.getTimestamps(); System.out.println("提交时间: " + timestamps.get("submittedTime")); System.out.println("开始执行时间: " + timestamps.get("startedTime")); System.out.println("执行结束时间: " + timestamps.get("completedTime"));
注意:这种方式无法直接获取Hazelcast集群内部的任务调度时间,若需要调度时间,推荐使用方法二。
方法二:通过TaskListener监听任务生命周期
利用Hazelcast提供的TaskListener接口,监听任务的提交、启动、完成事件,在回调中记录对应时间戳,能精准捕获任务在集群中的调度节点时间。
自定义TaskListener实现
import com.hazelcast.core.TaskListener; import java.util.concurrent.ConcurrentHashMap; import java.util.Map; public class TimedTaskListener implements TaskListener<Object> { private final Map<String, Map<String, Long>> taskTimeMap = new ConcurrentHashMap<>(); @Override public void onSubmitted(String taskId) { Map<String, Long> timeRecord = new ConcurrentHashMap<>(); timeRecord.put("submittedTime", System.currentTimeMillis()); taskTimeMap.put(taskId, timeRecord); } @Override public void onStarted(String taskId) { // onStarted触发时,任务已被调度到节点并开始执行,此时间可作为调度+开始时间 taskTimeMap.get(taskId).put("scheduledStartedTime", System.currentTimeMillis()); } @Override public void onCompleted(String taskId, Object result) { taskTimeMap.get(taskId).put("completedTime", System.currentTimeMillis()); } @Override public void onFailed(String taskId, Throwable throwable) { // 任务失败时同样记录结束时间 taskTimeMap.get(taskId).put("completedTime", System.currentTimeMillis()); } // 根据任务ID获取对应时间戳 public Map<String, Long> getTaskTimestamps(String taskId) { return taskTimeMap.get(taskId); } }
使用方式
// 初始化监听器 TimedTaskListener taskListener = new TimedTaskListener(); // 提交任务并绑定监听器 DurableExecutorFuture<String> future = durableExecutor.submit(yourOriginalCallable, taskListener); // 获取任务唯一ID String taskId = future.getTaskId(); // 任务完成后获取时间戳 Map<String, Long> timestamps = taskListener.getTaskTimestamps(taskId); System.out.println("提交时间: " + timestamps.get("submittedTime")); System.out.println("调度并开始执行时间: " + timestamps.get("scheduledStartedTime")); System.out.println("执行结束时间: " + timestamps.get("completedTime"));
关键注意事项
- 集群环境下,时间戳基于任务执行节点的系统时间,需确保集群所有节点时间同步(如使用NTP服务),否则时间戳会存在偏差。
- 对于持久化的Durable任务,若任务因节点故障重启,时间戳会重新记录,需根据业务需求做特殊处理。
内容的提问来源于stack exchange,提问作者Manish Sharma
相关产品推荐
相关产品推荐

