本文目录导读:

- 项目结构
- Maven依赖 (pom.xml)
- 主应用类
- 基础任务实现
- 数据处理任务(复杂任务示例)
- 定时任务
- 服务层
- 报表服务
- 配置类
- 模型类
- 配置文件 (application.yml)
- 控制器
- 具体命令执行示例
- 测试类
- Dockerfile (可选)
- 使用说明
- 特性说明
我将为您提供一个完整的Spring Cloud Task案例,包含批处理任务、任务调度和任务管理。
项目结构
spring-cloud-task-demo
├── pom.xml
├── src/main/java/com/example/task
│ ├── TaskDemoApplication.java
│ ├── config
│ │ ├── TaskConfig.java
│ │ └── BatchConfig.java
│ ├── task
│ │ ├── SimpleTask.java
│ │ ├── DataProcessingTask.java
│ │ └── ScheduledTasks.java
│ ├── service
│ │ ├── OrderService.java
│ │ └── ReportService.java
│ ├── listener
│ │ └── TaskExecutionListener.java
│ └── model
│ ├── Order.java
│ └── TaskResult.java
└── src/main/resources
└── application.yml
Maven依赖 (pom.xml)
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>2.7.14</version>
<relativePath/>
</parent>
<groupId>com.example</groupId>
<artifactId>spring-cloud-task-demo</artifactId>
<version>1.0.0-SNAPSHOT</version>
<packaging>jar</packaging>
<properties>
<java.version>1.8</java.version>
<spring-cloud.version>2021.0.8</spring-cloud.version>
</properties>
<dependencies>
<!-- Spring Cloud Task -->
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-task</artifactId>
</dependency>
<!-- Spring Batch -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-batch</artifactId>
</dependency>
<!-- Web -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- JPA -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-jpa</artifactId>
</dependency>
<!-- MySQL -->
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>8.0.28</version>
</dependency>
<!-- H2 Database for testing -->
<dependency>
<groupId>com.h2database</groupId>
<artifactId>h2</artifactId>
<scope>runtime</scope>
</dependency>
<!-- Lombok -->
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<!-- Actuator for monitoring -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<!-- Test -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-dependencies</artifactId>
<version>${spring-cloud.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
主应用类
package com.example.task;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.cloud.task.configuration.EnableTask;
import org.springframework.scheduling.annotation.EnableScheduling;
@SpringBootApplication
@EnableTask
@EnableScheduling
public class TaskDemoApplication {
public static void main(String[] args) {
SpringApplication.run(TaskDemoApplication.class, args);
}
}
基础任务实现
package com.example.task.task;
import org.springframework.cloud.task.listener.TaskExecutionListener;
import org.springframework.cloud.task.repository.TaskExecution;
import org.springframework.stereotype.Component;
@Component
public class SimpleTask implements TaskExecutionListener {
@Override
public void onTaskStartup(TaskExecution taskExecution) {
System.out.println("=== Task Started === " + taskExecution.getTaskName());
}
@Override
public void onTaskEnd(TaskExecution taskExecution) {
System.out.println("=== Task Ended === " + taskExecution.getTaskName() +
" Exit Code: " + taskExecution.getExitCode());
}
@Override
public void onTaskFailed(TaskExecution taskExecution, Throwable throwable) {
System.err.println("=== Task Failed === " + taskExecution.getTaskName());
throwable.printStackTrace();
}
public void executeTask(String[] args) {
System.out.println("执行简单任务,参数:" + String.join(", ", args));
// 模拟任务处理
try {
System.out.println("正在处理简单任务...");
Thread.sleep(2000);
System.out.println("简单任务处理完成");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("任务执行失败", e);
}
}
}
数据处理任务(复杂任务示例)
package com.example.task.task;
import com.example.task.service.OrderService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.time.LocalDateTime;
import java.util.concurrent.atomic.AtomicInteger;
@Slf4j
@Component
public class DataProcessingTask {
@Autowired
private OrderService orderService;
private final AtomicInteger progress = new AtomicInteger(0);
private volatile boolean isCancelled = false;
private volatile boolean isRunning = false;
public TaskResult processOrders(String[] args) {
isRunning = true;
isCancelled = false;
progress.set(0);
TaskResult result = new TaskResult();
result.setStartTime(LocalDateTime.now());
try {
log.info("开始处理订单数据,参数: {}", String.join(", ", args));
// 获取需要处理的订单
int totalOrders = orderService.getOrderCount();
result.setTotalCount(totalOrders);
int batchSize = 100;
int processedCount = 0;
// 分批处理订单
while (processedCount < totalOrders && !isCancelled) {
int limit = Math.min(batchSize, totalOrders - processedCount);
orderService.processOrders(processedCount, limit);
processedCount += limit;
progress.set((processedCount * 100) / totalOrders);
log.info("数据处理进度: {}% ({}/{})", progress.get(), processedCount, totalOrders);
// 模拟处理时间
Thread.sleep(1000);
}
if (isCancelled) {
result.setStatus("CANCELLED");
result.setMessage("任务被取消");
log.warn("订单处理任务被取消");
} else {
result.setStatus("SUCCESS");
result.setMessage("订单处理完成");
log.info("订单处理任务完成");
}
result.setProcessedCount(processedCount);
} catch (Exception e) {
log.error("订单处理失败", e);
result.setStatus("FAILED");
result.setMessage(e.getMessage());
} finally {
isRunning = false;
result.setEndTime(LocalDateTime.now());
}
return result;
}
public void cancelProcessing() {
isCancelled = true;
log.info("收到取消任务信号");
}
public int getProgress() {
return progress.get();
}
public boolean isRunning() {
return isRunning;
}
}
定时任务
package com.example.task.task;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.text.SimpleDateFormat;
import java.util.Date;
@Slf4j
@Component
public class ScheduledTasks {
@Autowired
private ReportService reportService;
private static final SimpleDateFormat dateFormat =
new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
// 每分钟执行一次
@Scheduled(cron = "0 * * * * ?")
public void generateDailyReport() {
log.info("生成每日报告任务开始 - {}", dateFormat.format(new Date()));
reportService.generateReport();
log.info("生成每日报告任务结束 - {}", dateFormat.format(new Date()));
}
// 每小时执行一次
@Scheduled(cron = "0 0 * * * ?")
public void cleanOldData() {
log.info("清理旧数据任务开始 - {}", dateFormat.format(new Date()));
// 清理操作
log.info("清理旧数据任务结束 - {}", dateFormat.format(new Date()));
}
// 每天中午12点执行
@Scheduled(cron = "0 0 12 * * ?")
public void sendNotifications() {
log.info("开始发送通知 - {}", dateFormat.format(new Date()));
// 发送通知逻辑
log.info("通知发送完成 - {}", dateFormat.format(new Date()));
}
}
服务层
package com.example.task.service;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import javax.annotation.PostConstruct;
import java.util.ArrayList;
import java.util.List;
import java.util.Random;
import java.util.concurrent.ConcurrentHashMap;
@Slf4j
@Service
public class OrderService {
private List<Order> orders = new ArrayList<>();
private ConcurrentHashMap<Long, Order> processedOrders = new ConcurrentHashMap<>();
private Random random = new Random();
@PostConstruct
public void init() {
// 模拟生成1000条订单数据
for (int i = 0; i < 1000; i++) {
Order order = new Order();
order.setId((long) i);
order.setOrderNo("ORD" + String.format("%06d", i));
order.setAmount(random.nextDouble() * 1000);
order.setStatus("NEW");
order.setCustomerName("Customer" + i);
orders.add(order);
}
log.info("初始化订单数据完成,共 {} 条", orders.size());
}
public void processOrders(int start, int limit) {
for (int i = start; i < start + limit; i++) {
if (i < orders.size()) {
Order order = orders.get(i);
// 模拟处理订单
order.setStatus("PROCESSED");
order.setProcessedTime(new java.util.Date());
processedOrders.put(order.getId(), order);
// 模拟处理耗时
try {
Thread.sleep(10);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
}
public int getOrderCount() {
return orders.size();
}
public List<Order> getProcessedOrders() {
return new ArrayList<>(processedOrders.values());
}
public double getTotalProcessedAmount() {
return processedOrders.values().stream()
.mapToDouble(Order::getAmount)
.sum();
}
}
报表服务
package com.example.task.service;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.io.File;
import java.io.FileWriter;
import java.io.IOException;
import java.text.SimpleDateFormat;
import java.util.Date;
@Slf4j
@Service
public class ReportService {
@Autowired
private OrderService orderService;
public void generateReport() {
SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd");
String dateStr = sdf.format(new Date());
String fileName = "report_" + dateStr + ".txt";
try {
File file = new File("reports/" + fileName);
file.getParentFile().mkdirs();
try (FileWriter writer = new FileWriter(file)) {
writer.write("=== 订单处理报表 ===\n");
writer.write("生成时间: " + new Date() + "\n");
writer.write("总订单数: " + orderService.getOrderCount() + "\n");
writer.write("已处理订单数: " + orderService.getProcessedOrders().size() + "\n");
writer.write("总金额: " + String.format("%.2f",
orderService.getTotalProcessedAmount()) + " 元\n");
}
log.info("报表生成成功: {}", file.getAbsolutePath());
} catch (IOException e) {
log.error("报表生成失败", e);
}
}
}
配置类
package com.example.task.config;
import com.example.task.listener.TaskExecutionListener;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class TaskConfig {
@Bean
public TaskExecutionListener taskExecutionListener() {
return new TaskExecutionListener();
}
@Bean
public org.springframework.batch.core.launch.JobLauncher jobLauncher() {
return new org.springframework.batch.core.launch.support.SimpleJobLauncher();
}
}
模型类
package com.example.task.model;
import lombok.Data;
import java.util.Date;
@Data
public class Order {
private Long id;
private String orderNo;
private Double amount;
private String status;
private String customerName;
private Date processedTime;
}
package com.example.task.model;
import lombok.Data;
import java.time.LocalDateTime;
@Data
public class TaskResult {
private String status;
private String message;
private LocalDateTime startTime;
private LocalDateTime endTime;
private int totalCount;
private int processedCount;
public long getDuration() {
if (startTime != null && endTime != null) {
return java.time.Duration.between(startTime, endTime).getSeconds();
}
return 0;
}
}
配置文件 (application.yml)
spring:
application:
name: spring-cloud-task-demo
datasource:
url: jdbc:h2:mem:testdb
driver-class-name: org.h2.Driver
username: sa
password:
jpa:
hibernate:
ddl-auto: create-drop
show-sql: true
cloud:
task:
table-prefix: TASK_
initialize-enabled: true
close-context-enabled: true
h2:
console:
enabled: true
path: /h2-console
batch:
job:
enabled: false
server:
port: 8080
management:
endpoints:
web:
exposure:
include: "*"
logging:
level:
root: INFO
com.example.task: DEBUG
控制器
package com.example.task.controller;
import com.example.task.model.TaskResult;
import com.example.task.task.DataProcessingTask;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
@Slf4j
@RestController
@RequestMapping("/api/tasks")
public class TaskController {
@Autowired
private DataProcessingTask dataProcessingTask;
@PostMapping("/start")
public TaskResult startTask(@RequestBody(required = false) String[] args) {
log.info("启动数据处理任务");
return dataProcessingTask.processOrders(args != null ? args : new String[]{});
}
@PostMapping("/cancel")
public String cancelTask() {
dataProcessingTask.cancelProcessing();
return "取消请求已发送";
}
@GetMapping("/progress")
public int getProgress() {
return dataProcessingTask.getProgress();
}
@GetMapping("/status")
public String getStatus() {
return dataProcessingTask.isRunning() ? "RUNNING" : "IDLE";
}
}
具体命令执行示例
package com.example.task.command;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.CommandLineRunner;
import org.springframework.stereotype.Component;
import com.example.task.model.TaskResult;
import com.example.task.task.DataProcessingTask;
import com.example.task.task.SimpleTask;
@Slf4j
@Component
public class TaskCommandRunner implements CommandLineRunner {
@Autowired
private SimpleTask simpleTask;
@Autowired
private DataProcessingTask dataProcessingTask;
@Override
public void run(String... args) throws Exception {
log.info("任务执行程序启动");
// 检查是否有命令行参数
if (args.length > 0 && args[0].equals("run-demo")) {
log.info("执行演示任务");
// 执行简单任务
simpleTask.executeTask(args);
// 执行数据处理任务
TaskResult result = dataProcessingTask.processOrders(args);
log.info("任务执行结果: {}", result.getStatus());
}
}
}
测试类
package com.example.task;
import org.junit.jupiter.api.Test;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.beans.factory.annotation.Autowired;
import com.example.task.service.OrderService;
@SpringBootTest
class TaskDemoApplicationTests {
@Autowired
private OrderService orderService;
@Test
void contextLoads() {
}
@Test
void testOrderService() {
System.out.println("测试订单服务");
assert orderService.getOrderCount() == 1000;
}
}
Dockerfile (可选)
FROM openjdk:8-jdk-alpine
VOLUME /tmp
ARG JAR_FILE=target/*.jar
COPY ${JAR_FILE} app.jar
ENTRYPOINT ["java","-jar","/app.jar"]
使用说明
运行任务
# 运行Spring Boot应用 mvn spring-boot:run # 或者打包运行 mvn clean package java -jar target/spring-cloud-task-demo-1.0.0-SNAPSHOT.jar # 带参数运行 java -jar spring-cloud-task-demo.jar run-demo --parameter1=value1
通过API触发任务
# 启动任务 curl -X POST http://localhost:8080/api/tasks/start -H "Content-Type: application/json" -d '["param1","param2"]' # 取消任务 curl -X POST http://localhost:8080/api/tasks/cancel # 查询进度 curl http://localhost:8080/api/tasks/progress # 查询状态 curl http://localhost:8080/api/tasks/status
监控Dashboard
访问 http://localhost:8080/actuator 查看各种监控指标。
特性说明
- 任务生命周期管理:Spring Cloud Task自动管理任务状态
- 数据持久化:任务状态存储在数据库中
- 参数传递:支持通过命令行或REST API传参
- 定时任务:支持Cron表达式定时执行
- 监控告警:通过Actuator提供监控端点
- 优雅关闭:支持任务取消和清理
这个案例包含了Spring Cloud Task的主要功能,可以作为实际项目的参考模板。