Skip to content
第 70 / 250 章后端⏱ 10 分钟阅读

第 70 章:定时任务与异步任务

学习目标

  • 掌握 Spring @Scheduled@Async 的使用
  • 学会分布式调度 XXL-Job 与 PowerJob
  • 理解任务幂等、动态调度、监控告警

一、定时任务的场景

二、Spring @Scheduled(单机版)

java
@SpringBootApplication
@EnableScheduling                                  // ① 启用调度
public class TaskflowApplication { }
java
@Component
@RequiredArgsConstructor
@Slf4j
public class ScheduledTasks {

    private final OrderMapper orderMapper;
    private final LocalMessageMapper localMessageMapper;

    // ② fixedRate:上次开始时间算
    @Scheduled(fixedRate = 5000)
    public void task1() {
        log.info("每 5 秒执行一次");
    }

    // ③ fixedDelay:上次结束时间算(推荐,避免堆积)
    @Scheduled(fixedDelay = 5000)
    public void task2() {
        log.info("上次结束后 5 秒执行");
    }

    // ④ cron 表达式(最灵活)
    @Scheduled(cron = "0 0 2 * * ?")               // 每天凌晨 2 点
    public void cleanLog() {
        log.info("清理过期日志");
        logMapper.deleteExpired();
    }

    // ⑤ cron 表达式:每分钟第 0 秒
    @Scheduled(cron = "0 * * * * ?")
    public void sendPendingMessages() {
        List<LocalMessage> msgs = localMessageMapper.selectPending(100);
        msgs.forEach(msg -> rabbitTemplate.send(msg.getTopic(), msg.getPayload()));
    }
}

Cron 表达式

秒   分   时   日   月   周
0    0    2    *    *    ?       每天凌晨 2 点
0    */5  *    *    *    ?       每 5 分钟
0    0    9-18 *    *    ?       9 点到 18 点整点
0    30   9    1    *    ?       每月 1 号 9:30
0    0    0    *    *    MON-FRI 工作日 0 点
字段范围特殊字符
0-59* , - /
0-59* , - /
0-23* , - /
1-31* , - / ? L W
1-12* , - /
0-7 (0 和 7 都代表周日)* , - / ? L #

⚠️ ? 只能用在「日」或「周」,因为两者互斥(月日和周几不能同时指定)。

@Scheduled 的局限性

java
// ❌ 1. 单点故障:一台机器跑挂了,任务就没了
// ❌ 2. 无法水平扩展:多台机器会重复执行(除非手动加锁)
// ❌ 3. 任务不能动态调整:改 cron 要重启服务
// ❌ 4. 没有重试、没有监控、没有告警

生产环境一定要用分布式调度框架:XXL-Job、PowerJob、Elastic-Job。

三、@Async 异步任务

java
@SpringBootApplication
@EnableAsync                                       // ① 启用异步
public class TaskflowApplication { }

@Configuration
public class AsyncConfig implements AsyncConfigurer {

    @Override
    public Executor getAsyncExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(8);               // ② 核心线程
        executor.setMaxPoolSize(16);              // 最大线程
        executor.setQueueCapacity(200);
        executor.setKeepAliveSeconds(60);
        executor.setThreadNamePrefix("async-");
        executor.setRejectedExecutionHandler(
                new ThreadPoolExecutor.CallerRunsPolicy());
        // ③ 等待所有任务完成再关闭
        executor.setWaitForTasksToCompleteOnShutdown(true);
        executor.setAwaitTerminationSeconds(60);
        executor.initialize();
        return executor;
    }
}
java
@Service
@RequiredArgsConstructor
public class UserService {

    private final EmailService emailService;

    // ① 注册时异步发欢迎邮件
    public void register(UserDTO dto) {
        userMapper.insert(user);
        emailService.sendWelcome(user.getEmail());    // ② 同步调用,耗时
    }
}

@Service
public class EmailService {

    @Async                                          // ③ 异步执行
    public void sendWelcome(String email) {
        // 模拟发邮件
        try { Thread.sleep(3000); } catch (Exception e) {}
        log.info("已发送欢迎邮件至 {}", email);
    }
}

@Async 的坑

java
// ❌ 坑 1:同类自调用
@Service
public class UserService {

    public void register() {
        this.sendWelcome();                          // ❌ 不走代理,异步失效
    }

    @Async
    public void sendWelcome() { ... }
}

// ❌ 坑 2:返回值的任务无法直接获取结果
@Async
public String compute() { return "result"; }        // ❌ 返回 null
String result = userService.compute();

// ✅ 用 CompletableFuture
@Async
public CompletableFuture<String> compute() {
    return CompletableFuture.completedFuture("result");
}
CompletableFuture<String> future = userService.compute();
String result = future.get(5, TimeUnit.SECONDS);

// ❌ 坑 3:异常无法被全局异常处理器捕获
// @Async 方法抛异常会进入异步线程的异常处理
// ✅ 用 AsyncUncaughtExceptionHandler 处理

四、XXL-Job 分布式调度

架构

部署调度中心

bash
# 下载 xxl-job-admin 源码
# 修改 application.properties:数据库地址
# 启动调度中心
# 访问 http://localhost:8080/xxl-job-admin 默认账号 admin/123456

执行器集成

xml
<dependency>
    <groupId>com.xuxueli</groupId>
    <artifactId>xxl-job-core</artifactId>
    <version>2.4.0</version>
</dependency>
yaml
xxl:
  job:
    admin:
      addresses: http://localhost:8080/xxl-job-admin
    executor:
      appname: taskflow-executor
      port: 9999
      logpath: /data/applogs/xxl-job/jobhandler
      logretentiondays: 30
    accessToken: default_token
java
@Configuration
public class XxlJobConfig {

    @Value("${xxl.job.admin.addresses}")
    private String adminAddresses;

    @Value("${xxl.job.accessToken}")
    private String accessToken;

    @Bean
    public XxlJobSpringExecutor xxlJobExecutor() {
        XxlJobSpringExecutor executor = new XxlJobSpringExecutor();
        executor.setAdminAddresses(adminAddresses);
        executor.setAccessToken(accessToken);
        executor.setAppname("taskflow-executor");
        executor.setPort(9999);
        executor.setLogPath("/data/applogs/xxl-job/jobhandler");
        executor.setLogRetentionDays(30);
        return executor;
    }
}

编写任务

java
@Component
@Slf4j
public class OrderJobHandler {

    @XxlJob("orderCloseJob")
    public ReturnT<String> closeTimeoutOrders(String param) {
        XxlJobHelper.log("开始关闭超时订单");

        // ① 解析参数(JSON 或字符串)
        int timeoutMinutes = 30;
        if (param != null && !param.isEmpty()) {
            timeoutMinutes = Integer.parseInt(param);
        }

        // ② 分片参数(分片广播)
        int shardIndex = XxlJob.getShardIndex();
        int shardTotal = XxlJob.getShardTotal();

        // ③ 执行业务
        int count = orderService.closeTimeoutOrders(timeoutMinutes, shardIndex, shardTotal);
        XxlJobHelper.log("本次处理 {} 条", count);

        return ReturnT.SUCCESS;
    }

    @XxlJob("dailyReportJob")
    public ReturnT<String> generateDailyReport(String param) {
        try {
            reportService.generate();
            return ReturnT.SUCCESS;
        } catch (Exception e) {
            XxlJobHelper.log("生成日报失败", e);
            return ReturnT.FAILED;                  // 失败会重试
        }
    }
}

调度中心配置

路由策略含义
第一个固定选第一台机器
最后一个固定选最后一台
轮询轮流选
随机随机选
一致性 HASH同一个任务参数总是到同一台
最不频繁使用选最闲的
故障转移选第一个,失败转下一个
忙碌转移选空闲的
分片广播所有机器都执行(并行)

分片广播实战

java
@XxlJob("shardingJob")
public ReturnT<String> shardingJob(String param) {
    int shardIndex = XxlJob.getShardIndex();     // 当前分片索引(0,1,2...)
    int shardTotal = XxlJob.getShardTotal();     // 总分片数

    // 处理第 shardIndex 段
    int pageSize = 1000;
    int start = shardIndex * pageSize;
    int end = start + pageSize;

    List<Order> orders = orderMapper.selectRange(start, end);
    for (Order order : orders) {
        process(order);
    }

    return ReturnT.SUCCESS;
}
// 100 万订单,10 个分片,每台机器处理 10 万,10 倍加速

五、PowerJob 现代化调度(推荐)

PowerJob(瓴犀)对比 XXL-Job:

  • 支持 MapReduce 分布式计算
  • 支持工作流(DAG 任务依赖)
  • 内置日志、监控、告警
  • 支持秒级调度
xml
<dependency>
    <groupId>tech.powerjob</groupId>
    <artifactId>powerjob-spring-boot-starter</artifactId>
    <version>5.1.0</version>
</dependency>
yaml
powerjob:
  worker:
    app-name: taskflow-worker
    enable-test-mode: false
    port: 27777
    transport: netty
    registry-address: localhost:7700                  # PowerJob Server 地址
    store-strategy: disk
java
@Component
@Slf4j
public class OrderJobs implements BasicProcessor {

    @Override
    public ProcessResult process(TaskContext ctx) throws Exception {
        // ① 任务参数
        String jobParams = ctx.getJobParams();
        int timeoutMinutes = 30;

        // ② 业务处理
        int count = orderService.closeTimeoutOrders(timeoutMinutes,
                ctx.getShardingId(), ctx.getShardingTotalCount());

        // ③ 返回结果(用于上层决策)
        return new ProcessResult(true, "处理 " + count + " 条订单");
    }
}

六、幂等性保证(关键!)

java
@XxlJob("orderCloseJob")
public ReturnT<String> orderCloseJob(String param) {
    // ① 业务幂等 key(同一参数不重复执行)
    String bizKey = "job:closeOrder:" + LocalDate.now() + ":" + param;

    // ② Redis SETNX 抢锁(5 分钟有效期)
    Boolean first = redis.opsForValue().setIfAbsent(bizKey, "1", 5, TimeUnit.MINUTES);
    if (Boolean.FALSE.equals(first)) {
        log.warn("任务正在执行中,跳过本次触发");
        return ReturnT.SUCCESS;
    }

    try {
        // ③ 执行业务
        orderService.closeTimeoutOrders(30);
        return ReturnT.SUCCESS;
    } finally {
        redis.delete(bizKey);
    }
}

七、任务监控与告警

yaml
# XXL-Job 调度中心配置:邮件告警
xxl.job:
  admin:
    addresses: http://localhost:8080/xxl-job-admin

# PowerJob 告警:钉钉、飞书、企业微信、邮件
powerjob:
  worker:
    alarm:
      enable: true
      webhook: https://oapi.dingtalk.com/robot/send?access_token=xxx

八、选型对比

框架学习成本特性推荐
Spring @Scheduled极低单机简单任务临时脚本
XXL-Job成熟、社区大绝大多数场景
PowerJobDAG 工作流、MapReduce复杂任务编排
Elastic-Job弹性、分片互联网大厂
Quartz老牌遗留系统

九、本章小结

要点关键
@Scheduled单机 cron / fixedRate / fixedDelay
@Async异步任务,注意同类调用失效
分布式调度XXL-Job(最常用)/ PowerJob(现代化)
任务幂等Redis SETNX + 业务幂等双保险
分片广播10 台机器分片处理 10 倍加速
告警邮件 / 钉钉 / 飞书
选型简单用 Spring,分布式用 XXL-Job

动手练习

练习 1:基础题

实现一个每天凌晨 2 点清理 30 天前日志的 @Scheduled 任务,并配置线程池参数。

练习 2:进阶题

集成 XXL-Job 调度中心,实现一个「每分钟扫描超时订单并自动关闭」的任务。要点:

  • 调度中心后台配置 Cron
  • 任务幂等(防重复)
  • 失败重试
  • 钉钉告警

练习 3:思考题

你的系统有 1000 万张订单,需要每天凌晨做一次全量对账。如何用 XXL-Job 的分片广播加速?


下一章第 71 章:文件上传与对象存储

本站基于 VitePress 构建 · 由 Codebook 团队维护