You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Batch调度任务编码实现求助:含触发频率、重试、告警等配置

Spring Batch 批处理作业完整实现方案

以下是覆盖所有需求的Spring Batch作业实现,包含作业定义、调度、重试、告警全流程配置:

一、作业基础定义(作业ID/描述&步骤)

  • 作业ID:user-data-import-job
  • 作业描述:批量读取CSV用户数据,清洗后写入MySQL数据库,包含3个核心步骤:
    1. 读取步骤:从指定目录读取CSV文件
    2. 处理步骤:校验用户手机号、邮箱格式,补全缺失字段
    3. 写入步骤:将清洗后的数据插入用户表

作业配置代码

import org.springframework.batch.core.Job;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.configuration.annotation.EnableBatchProcessing;
import org.springframework.batch.core.configuration.annotation.JobBuilderFactory;
import org.springframework.batch.core.configuration.annotation.StepBuilderFactory;
import org.springframework.batch.item.database.JdbcBatchItemWriter;
import org.springframework.batch.item.file.FlatFileItemReader;
import org.springframework.batch.item.file.mapping.BeanWrapperFieldSetMapper;
import org.springframework.batch.item.file.mapping.DefaultLineMapper;
import org.springframework.batch.item.file.transform.DelimitedLineTokenizer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.io.FileSystemResource;
import org.springframework.batch.item.database.BeanPropertyItemSqlParameterSourceProvider;
import org.springframework.batch.item.ItemProcessor;

import javax.sql.DataSource;

@Configuration
@EnableBatchProcessing
public class BatchJobConfiguration {

    private final JobBuilderFactory jobBuilderFactory;
    private final StepBuilderFactory stepBuilderFactory;
    private final DataSource dataSource;

    public BatchJobConfiguration(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory, DataSource dataSource) {
        this.jobBuilderFactory = jobBuilderFactory;
        this.stepBuilderFactory = stepBuilderFactory;
        this.dataSource = dataSource;
    }

    // 1. 读取CSV文件的Reader
    @Bean
    public FlatFileItemReader<User> userItemReader() {
        FlatFileItemReader<User> reader = new FlatFileItemReader<>();
        reader.setResource(new FileSystemResource("/data/users.csv"));
        reader.setLineMapper(new DefaultLineMapper<>() {{
            setLineTokenizer(new DelimitedLineTokenizer() {{
                setNames("id", "name", "phone", "email");
            }});
            setFieldSetMapper(new BeanWrapperFieldSetMapper<>() {{
                setTargetType(User.class);
            }});
        }});
        return reader;
    }

    // 2. 数据处理器
    @Bean
    public UserDataProcessor userDataProcessor() {
        return new UserDataProcessor();
    }

    // 3. 写入数据库的Writer
    @Bean
    public JdbcBatchItemWriter<User> userItemWriter() {
        JdbcBatchItemWriter<User> writer = new JdbcBatchItemWriter<>();
        writer.setDataSource(dataSource);
        writer.setSql("INSERT INTO users (id, name, phone, email) VALUES (:id, :name, :phone, :email)");
        writer.setItemSqlParameterSourceProvider(new BeanPropertyItemSqlParameterSourceProvider<>());
        return writer;
    }

    // 定义步骤
    @Bean
    public Step userImportStep() {
        return stepBuilderFactory.get("user-import-step")
                .<User, User>chunk(100)
                .reader(userItemReader())
                .processor(userDataProcessor())
                .writer(userItemWriter())
                // 异常重试配置(对应需求4)
                .faultTolerant()
                .retryLimit(3) // 异常发生后重试3次
                .retry(Exception.class)
                .build();
    }

    // 定义作业
    @Bean
    public Job userDataImportJob() {
        return jobBuilderFactory.get("user-data-import-job")
                .description("批量导入用户数据至MySQL,包含读取、清洗、写入三个步骤")
                .start(userImportStep())
                .build();
    }
}

// 数据实体类
class User {
    private Long id;
    private String name;
    private String phone;
    private String email;
    // getter/setter省略
}

// 数据处理器类
class UserDataProcessor implements ItemProcessor<User, User> {
    @Override
    public User process(User user) throws Exception {
        // 清洗逻辑:格式化手机号、校验邮箱
        user.setPhone(user.getPhone().replaceAll("[^0-9]", ""));
        if (!user.getEmail().matches("^[A-Za-z0-9+_.-]+@[A-Za-z0-9.-]+$")) {
            throw new IllegalArgumentException("邮箱格式错误:" + user.getEmail());
        }
        return user;
    }
}

二、多频率调度配置

使用Quartz实现全类型触发调度,以下是不同频率对应的Cron表达式及配置代码:

触发频率对应Cron表达式

  • 分钟级(每5分钟):0 */5 * * * ?
  • 小时级(每2小时):0 0 */2 * * ?
  • 每日一次(凌晨2点):0 0 2 * * ?
  • 指定时间(每月1号凌晨3点):0 0 3 1 * ?
  • 每周(每周日凌晨1点):0 0 1 ? * SUN
  • 每两周(第1、3个周日凌晨1点):0 0 1 ? * SUN#1,SUN#3
  • 每月(每月最后一天凌晨2点):0 0 2 L * ?
  • 每季度(每季度第一天凌晨3点):0 0 3 1 1,4,7,10 ?
  • 循环(每10秒一次):*/10 * * * * ?

调度配置代码

import org.quartz.*;
import org.springframework.batch.core.Job;
import org.springframework.batch.core.JobParameters;
import org.springframework.batch.core.launch.JobLauncher;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.retry.support.RetryTemplate;
import org.springframework.retry.policy.SimpleRetryPolicy;
import org.springframework.retry.backoff.FixedBackOffPolicy;

import java.io.File;

@Configuration
public class QuartzSchedulerConfiguration {

    private final JobLauncher jobLauncher;
    private final Job userDataImportJob;

    public QuartzSchedulerConfiguration(JobLauncher jobLauncher, Job userDataImportJob) {
        this.jobLauncher = jobLauncher;
        this.userDataImportJob = userDataImportJob;
    }

    // 定义Quartz JobDetail
    @Bean
    public JobDetail batchJobDetail() {
        return JobBuilder.newJob(BatchQuartzJob.class)
                .withIdentity("user-data-import-quartz-job")
                .storeDurably()
                .build();
    }

    // 定义触发器(以每日凌晨2点为例,可替换为其他Cron表达式)
    @Bean
    public Trigger batchJobTrigger() {
        CronScheduleBuilder scheduleBuilder = CronScheduleBuilder.cronSchedule("0 0 2 * * ?");
        return TriggerBuilder.newTrigger()
                .forJob(batchJobDetail())
                .withIdentity("user-data-import-trigger")
                .withSchedule(scheduleBuilder)
                .build();
    }

    // Quartz执行Spring Batch作业的实现类
    public static class BatchQuartzJob implements org.quartz.Job {
        private final JobLauncher jobLauncher;
        private final Job userDataImportJob;

        public BatchQuartzJob(JobLauncher jobLauncher, Job userDataImportJob) {
            this.jobLauncher = jobLauncher;
            this.userDataImportJob = userDataImportJob;
        }

        @Override
        public void execute(JobExecutionContext context) throws JobExecutionException {
            try {
                // 检查依赖条件(如CSV文件是否存在,对应需求3)
                checkDependency();
                jobLauncher.run(userDataImportJob, new JobParameters());
            } catch (Exception e) {
                // 依赖未满足或执行异常时触发告警
                AlertService.sendExceptionAlert("user-data-import-job", e.getMessage(), "user-import-step");
                throw new JobExecutionException(e);
            }
        }

        // 依赖检查&重试逻辑
        private void checkDependency() throws Exception {
            RetryTemplate retryTemplate = new RetryTemplate();
            retryTemplate.setRetryPolicy(new SimpleRetryPolicy(3)); // 依赖未满足重试3次
            retryTemplate.setBackOffPolicy(new FixedBackOffPolicy() {{
                setBackOffPeriod(5000); // 每次重试间隔5秒
            }});
            retryTemplate.execute(context -> {
                boolean fileExists = new File("/data/users.csv").exists();
                if (!fileExists) {
                    throw new IllegalStateException("依赖文件不存在");
                }
                return null;
            });
        }
    }
}

三、告警配置实现

1. 监控告警(作业未触发/超时)

  • 沟通层级:
    • 一级:dev-team@company.com、138xxxx1234
    • 二级:tech-lead@company.com、139xxxx4567
  • 邮件模板:

【监控告警】作业${jobId}触发异常
告警时间:${triggerTime}
告警类型:${alertType}(未触发/超时)
详细描述:${description}
请及时排查处理。

  • SMS模板:

监控告警:作业${jobId}${alertType},时间${triggerTime},请立即排查。

2. 异常告警(作业执行失败)

  • 沟通层级:
    • 一级:dev-oncall@company.com、137xxxx7890
    • 二级:cto@company.com、136xxxx0123
  • 邮件模板:

【异常告警】作业${jobId}执行失败
失败时间:${failTime}
异常信息:${exceptionMsg}
失败步骤:${failedStep}
请立即处理。

  • SMS模板:

异常告警:作业${jobId}执行失败,步骤${failedStep},请立即处理。

告警服务代码

import org.springframework.mail.SimpleMailMessage;
import org.springframework.mail.javamail.JavaMailSender;
import org.springframework.stereotype.Service;

import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;

@Service
public class AlertService {

    private final JavaMailSender mailSender;
    // 假设注入SMS服务客户端
    // private final SmsClient smsClient;

    public AlertService(JavaMailSender mailSender) {
        this.mailSender = mailSender;
    }

    // 发送监控告警
    public void sendMonitorAlert(String jobId, String alertType, String description) {
        String time = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
        // 发送一级邮件
        sendMail("dev-team@company.com", "【监控告警】作业" + jobId + "异常",
                String.format("【监控告警】作业`%s`触发异常\n告警时间:%s\n告警类型:%s\n详细描述:%s\n请及时排查处理。", jobId, time, alertType, description));
        // 发送一级SMS
        // sendSms("138xxxx1234", String.format("监控告警:作业%s%s,时间%s,请立即排查。", jobId, alertType, time));

        // 若30分钟未处理,发送二级告警(可结合定时任务实现)
    }

    // 发送异常告警
    public void sendExceptionAlert(String jobId, String exceptionMsg, String failedStep) {
        String time = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
        // 发送一级邮件
        sendMail("dev-oncall@company.com", "【异常告警】作业" + jobId + "执行失败",
                String.format("【异常告警】作业`%s`执行失败\n失败时间:%s\n异常信息:%s\n失败步骤:%s\n请立即处理。", jobId, time, exceptionMsg, failedStep));
        // 发送一级SMS
        // sendSms("137xxxx7890", String.format("异常告警:作业%s执行失败,步骤%s,请立即处理。", jobId, failedStep));

        // 若10分钟未处理,发送二级告警
    }

    private void sendMail(String to, String subject, String content) {
        SimpleMailMessage message = new SimpleMailMessage();
        message.setTo(to);
        message.setSubject(subject);
        message.setText(content);
        mailSender.send(message);
    }

    // private void sendSms(String phone, String content) {
    //     smsClient.send(phone, content);
    // }
}

内容的提问来源于stack exchange,提问作者hemantn28

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.08 02:31:15