Spring Batch调度任务编码实现求助:含触发频率、重试、告警等配置
Spring Batch 批处理作业完整实现方案
以下是覆盖所有需求的Spring Batch作业实现,包含作业定义、调度、重试、告警全流程配置:
一、作业基础定义(作业ID/描述&步骤)
- 作业ID:
user-data-import-job - 作业描述:批量读取CSV用户数据,清洗后写入MySQL数据库,包含3个核心步骤:
- 读取步骤:从指定目录读取CSV文件
- 处理步骤:校验用户手机号、邮箱格式,补全缺失字段
- 写入步骤:将清洗后的数据插入用户表
作业配置代码
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
相关产品推荐
相关产品推荐

