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

如何在Quarkus中实现由事件触发的带参周期性任务?

Java & Quarkus 实现带参数的动态周期性任务

一、Java 标准库方案:ScheduledExecutorService + 任务追踪

ScheduledExecutorService 本身支持动态提交带参数的周期性任务,核心是通过线程安全的映射关联用户ID与任务实例,实现创建启动、删除停止的生命周期管理。

实现步骤:

  1. 用ConcurrentHashMap<Long, ScheduledFuture<?>>维护用户ID和对应定时任务的映射,确保线程安全。
  2. 用户创建成功后,提交带用户ID参数的周期性任务到线程池,将返回的ScheduledFuture存入映射。
  3. 用户删除时,从映射中取出任务实例,调用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 02:20:28