如何在Quarkus中实现由事件触发的带参周期性任务?
Java & Quarkus 实现带参数的动态周期性任务
一、Java 标准库方案:ScheduledExecutorService + 任务追踪
ScheduledExecutorService 本身支持动态提交带参数的周期性任务,核心是通过线程安全的映射关联用户ID与任务实例,实现创建启动、删除停止的生命周期管理。
实现步骤:
- 用
ConcurrentHashMap<Long, ScheduledFuture<?>>维护用户ID和对应定时任务的映射,确保线程安全。 - 用户创建成功后,提交带用户ID参数的周期性任务到线程池,将返回的
ScheduledFuture存入映射。 - 用户删除时,从映射中取出任务实例,调用
cancel(true)终止任务并移除映射条目。
代码示例:
import java.util.concurrent.*; public class UserTaskManager { private final ConcurrentHashMap<Long, ScheduledFuture<?>> userTaskMap = new ConcurrentHashMap<>(); private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(5); // 创建用户后启动追踪任务 public void startUserTrackingTask(Long userId) { if (userTaskMap.containsKey(userId)) { return; } Runnable trackingTask = () -> trackUserStatus(userId); // 初始延迟1秒,每30秒执行一次 ScheduledFuture<?> future = scheduler.scheduleAtFixedRate(trackingTask, 1, 30, TimeUnit.SECONDS); userTaskMap.put(userId, future); } // 删除用户时停止追踪任务 public void stopUserTrackingTask(Long userId) { ScheduledFuture<?> future = userTaskMap.remove(userId); if (future != null && !future.isCancelled()) { future.cancel(true); } } // 实际状态追踪逻辑 private void trackUserStatus(Long userId) { System.out.println("Tracking user status for ID: " + userId); // 此处编写数据库查询、状态校验等业务代码 } // 应用关闭时清理线程池 public void shutdown() { scheduler.shutdown(); try { if (!scheduler.awaitTermination(10, TimeUnit.SECONDS)) { scheduler.shutdownNow(); } } catch (InterruptedException e) { scheduler.shutdownNow(); } } }
二、Quarkus 专属方案
Quarkus提供了更贴合其生态的动态定时任务支持,以下两种方案适配不同场景:
1. Quarkus Scheduler 动态任务
通过Quarkus内置的Scheduler接口,可动态创建带参数的定时任务,同样需要映射管理任务生命周期。
代码示例:
import io.quarkus.scheduler.Scheduler; import io.quarkus.scheduler.ScheduledTask; import jakarta.enterprise.context.ApplicationScoped; import java.util.concurrent.ConcurrentHashMap; @ApplicationScoped public class QuarkusUserTaskManager { private final ConcurrentHashMap<Long, ScheduledTask> userTaskMap = new ConcurrentHashMap<>(); private final Scheduler scheduler; // 构造注入Scheduler public QuarkusUserTaskManager(Scheduler scheduler) { this.scheduler = scheduler; } public void startUserTracking(Long userId) { if (userTaskMap.containsKey(userId)) { return; } // 动态创建任务,指定唯一标识、执行间隔和业务逻辑 ScheduledTask task = scheduler.newTask() .withIdentity("track-user-" + userId) .withInterval(30) // 每30秒执行一次 .execute(() -> trackUserStatus(userId)) .schedule(); userTaskMap.put(userId, task); } public void stopUserTracking(Long userId) { ScheduledTask task = userTaskMap.remove(userId); if (task != null) { task.cancel(); } } private void trackUserStatus(Long userId) { System.out.println("Quarkus tracking user: " + userId); // 业务逻辑实现 } }
2. Reactive 场景:Vertx 定时器 + 事件总线
如果是Quarkus Reactive应用,可结合Vertx定时器和事件总线,通过事件触发任务的启停。
代码示例:
import io.quarkus.vertx.ConsumeEvent; import io.vertx.core.Vertx; import jakarta.enterprise.context.ApplicationScoped; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; @ApplicationScoped public class ReactiveUserTaskManager { private final ConcurrentHashMap<Long, Long> userTimerIds = new ConcurrentHashMap<>(); private final Vertx vertx; public ReactiveUserTaskManager(Vertx vertx) { this.vertx = vertx; } // 监听用户创建事件(方法A创建用户后发送此事件) @ConsumeEvent("user.created") public void handleUserCreated(Long userId) { if (userTimerIds.containsKey(userId)) { return; } // 启动周期性定时器,返回timerId用于后续取消 long timerId = vertx.setPeriodic(TimeUnit.SECONDS.toMillis(30), id -> { trackUserStatus(userId); }); userTimerIds.put(userId, timerId); } // 监听用户删除事件 @ConsumeEvent("user.deleted") public void handleUserDeleted(Long userId) { Long timerId = userTimerIds.remove(userId); if (timerId != null) { vertx.cancelTimer(timerId); } } private void trackUserStatus(Long userId) { System.out.println("Reactive tracking user: " + userId); // Reactive风格的业务逻辑(如Panache Reactive数据库查询) } }
关于@Scheduled注解的局限
Quarkus的@Scheduled是静态声明式注解,任务在应用启动时就完成初始化,无法动态传递参数或根据事件启停,仅适用于固定的、无参数的定时任务场景。
内容的提问来源于stack exchange,提问作者seeusoon
相关产品推荐
相关产品推荐

