第 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_tokenjava
@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: diskjava
@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 | 中 | 成熟、社区大 | 绝大多数场景 |
| PowerJob | 中 | DAG 工作流、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 章:文件上传与对象存储 →