Skip to content

Latest commit

 

History

2 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 

Repository files navigation

分布式离线任务处理器

一个基于 Spring Boot 的离线任务处理示例项目,主要演示两类任务处理能力:

  • 定时任务调度:从数据库读取待执行任务,根据 cron 表达式更新下次执行时间,并调用对应业务处理器。
  • Fork/Join 任务处理:将一个主任务拆分为多个子任务并发执行,所有子任务完成后再做结果聚合。

项目适合作为定时任务、批处理任务、异步拆分任务、失败重试和任务状态流转的学习样例。

功能特性

  • 基于 MyBatis-Plus 访问 MySQL 任务表。
  • 基于 @ScheduleTaskHandler 注册业务定时任务处理器。
  • 支持 XXL-JOB 触发定时任务扫描入口。
  • 支持主任务拆分为子任务,并通过线程池并发处理。
  • 支持子任务执行完成后发布 Spring 事件,自动检查主任务是否可以进入后置聚合阶段。
  • 支持主任务和子任务失败重试配置。
  • 使用乐观版本号控制任务状态更新,降低重复处理风险。

技术栈

  • Java
  • Spring Boot 2.7.9
  • MyBatis-Plus 3.5.7
  • MySQL Connector/J 5.1.49
  • Quartz CronExpression
  • XXL-JOB Core 2.4.1
  • Hutool
  • Lombok

项目结构

src/main/java/org/practice
├── Application.java                         # Spring Boot 启动入口
├── controller                               # Web/调试入口
├── entity                                   # 任务表实体
├── forkjoin                                 # Fork/Join 主任务、子任务处理逻辑
│   └── impl/TestForkJoinTaskHandler.java    # Fork/Join 示例处理器
├── mapper                                   # MyBatis-Plus Mapper
├── schedule                                 # 定时任务扫描与处理器注册
├── service                                  # 业务 Service
└── util/RetryUtil.java                      # 重试工具

核心流程

定时任务流程

  1. ScheduleTaskManager 查询 next_execute_time 早于当前时间的任务配置。
  2. 命中任务后,根据任务的 cron 表达式计算并更新下一次执行时间。
  3. 通过 task_handler 找到对应的 @ScheduleTaskHandler 方法。
  4. 将数据库中的 JSON 参数反序列化后传给业务 handler 执行。

XXL-JOB 入口:

@XxlJob("queryAndHandleScheduleTask")
public void queryAndHandleTask()

本地调试入口:

GET /scheduleTaskManager/queryNeedHandleTask
GET /scheduleTaskManager/queryAndHandleTask

Fork/Join 任务流程

  1. 调用 ForkJoinTaskManager.submitTask 提交主任务。
  2. 根据 ForkJoinTaskHandler.generateSubTaskParamList 生成子任务参数。
  3. 保存主任务和子任务,主任务进入处理中,子任务进入待处理。
  4. 定时扫描待处理子任务,并发执行 executeSubTask
  5. 每个子任务结束后发布 ForkJoinSubTaskFinishedEvent
  6. 监听器检查同一主任务下所有子任务是否结束。
  7. 所有子任务完成后执行 doAfterAllSubTaskFinished 聚合结果。

示例处理器:

@Service
public class TestForkJoinTaskHandler implements ForkJoinTaskHandler {
    @Override
    public String name() {
        return "testForkJoinTaskHandler";
    }
}

本地运行

1. 准备环境

  • JDK 8 或更高版本
  • Maven 3.6+
  • MySQL 5.7+ 或兼容版本
  • 可选:XXL-JOB Admin

2. 配置数据库连接

默认配置位于 src/main/resources/application.yml

spring:
  datasource:
    url: ${DB_URL:jdbc:mysql://127.0.0.1:3306/fork_join_task_db?useUnicode=true&characterEncoding=UTF-8}
    username: ${DB_USERNAME:root}
    password: ${DB_PASSWORD:}

建议通过环境变量覆盖:

export DB_URL='jdbc:mysql://127.0.0.1:3306/fork_join_task_db?useUnicode=true&characterEncoding=UTF-8'
export DB_USERNAME='root'
export DB_PASSWORD='your_password'

Windows PowerShell:

$env:DB_URL='jdbc:mysql://127.0.0.1:3306/fork_join_task_db?useUnicode=true&characterEncoding=UTF-8'
$env:DB_USERNAME='root'
$env:DB_PASSWORD='your_password'

3. 启动应用

mvn spring-boot:run

应用默认端口:

http://localhost:8081

如何新增一个定时任务处理器

在 Spring Bean 中新增带 @ScheduleTaskHandler 的方法:

@Component
public class DemoScheduleHandler {
    @ScheduleTaskHandler("demoHandler")
    public void handle(DemoParam param) {
        // 业务逻辑
    }
}

数据库任务配置中的 task_handler 填写 demoHandlerparam 填写可反序列化为 DemoParam 的 JSON 字符串。

如何新增一个 Fork/Join 处理器

实现 ForkJoinTaskHandler

@Service
public class DemoForkJoinTaskHandler implements ForkJoinTaskHandler {
    @Override
    public String name() {
        return "demoForkJoinTaskHandler";
    }

    @Override
    public List<JSONObject> generateSubTaskParamList(JSONObject param) {
        // 将主任务参数拆分为多个子任务参数
    }

    @Override
    public JSONObject executeSubTask(JSONObject param) {
        // 执行单个子任务
    }

    @Override
    public JSONObject doAfterAllSubTaskFinished(List<ForkJoinSubTask> subTaskList, JSONObject mainTaskExecuteParam) {
        // 聚合所有子任务结果
    }
}

提交任务示例可参考 TestScheduleHandler

forkJoinTaskManager.submitTask("testForkJoinTaskHandler", JSONUtil.parseObj(param), null);

配置说明

配置项 说明
server.port 应用启动端口,默认 8081
spring.datasource.url MySQL 连接地址,可由 DB_URL 覆盖
spring.datasource.username MySQL 用户名,可由 DB_USERNAME 覆盖
spring.datasource.password MySQL 密码,可由 DB_PASSWORD 覆盖
xxl.job.admin.addresses XXL-JOB Admin 地址
xxl.job.executor.appname XXL-JOB 执行器名称
xxl.job.executor.logpath XXL-JOB 日志路径

注意事项

  • 当前仓库未包含建表 SQL,运行前需要根据 entitymapper 自行创建对应 MySQL 表。
  • 示例代码中的 controller 多数是预留入口,核心逻辑集中在 scheduleforkjoin 包。
  • 本项目偏向任务处理模型演示,生产使用前建议补充权限控制、任务幂等、监控告警、表结构迁移脚本和完整测试。

About

No description, website, or topics provided.

Resources

Stars

1 star

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages