Spring Cloud Task案例

wen java案例 3

本文目录导读:

Spring Cloud Task案例

  1. 项目结构
  2. Maven依赖 (pom.xml)
  3. 主应用类
  4. 基础任务实现
  5. 数据处理任务(复杂任务示例)
  6. 定时任务
  7. 服务层
  8. 报表服务
  9. 配置类
  10. 模型类
  11. 配置文件 (application.yml)
  12. 控制器
  13. 具体命令执行示例
  14. 测试类
  15. Dockerfile (可选)
  16. 使用说明
  17. 特性说明

我将为您提供一个完整的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 查看各种监控指标。

特性说明

  1. 任务生命周期管理:Spring Cloud Task自动管理任务状态
  2. 数据持久化:任务状态存储在数据库中
  3. 参数传递:支持通过命令行或REST API传参
  4. 定时任务:支持Cron表达式定时执行
  5. 监控告警:通过Actuator提供监控端点
  6. 优雅关闭:支持任务取消和清理

这个案例包含了Spring Cloud Task的主要功能,可以作为实际项目的参考模板。

抱歉,评论功能暂时关闭!