Spring Batch异常处理:将跳过的无效记录以ResponseEntity返回
问题描述
我是Spring Boot新手,正在做一个Spring Batch小项目积累经验。业务场景是:有两个.csv文件,分别存储员工和公司经理数据,需要读取文件并将每条记录存入数据库。已实现通过控制器端点上传MultipartFile并启动Job,但遇到以下问题:
- 需要处理多种验证:用JSR 380验证实体,同时检查业务规则(比如员工直属经理必须和员工同部门,否则抛出异常)
- 无效或不符合逻辑的错误记录需跳过不入库,但要把这些记录存入Map/List,最终通过ResponseEntity返回给客户端,让客户端知晓需修正的行
我知道要用到Listeners,但不知道如何将异常记录存入集合并返回成ResponseEntity。以下是我的相关代码:
相关代码
EmployeeBatchConfig.java
@Configuration @EnableBatchProcessing @AllArgsConstructor public class EmployeeBatchConfig { private JobBuilderFactory jobBuilderFactory; private StepBuilderFactory stepBuilderFactory; private EmployeeRepository employeeRepository; private EmployeeItemWriter employeeItemWriter; @Bean @StepScope public FlatFileItemReader<EmployeeDto> itemReader(@Value("# {jobParameters[fullPathFileName]}") final String pathFile) { FlatFileItemReader<EmployeeDto> flatFileItemReader = new FlatFileItemReader<>(); flatFileItemReader.setResource(new FileSystemResource(new File(pathFile))); flatFileItemReader.setName("CSV-Reader"); flatFileItemReader.setLinesToSkip(1); flatFileItemReader.setLineMapper(lineMapper()); return flatFileItemReader; } private LineMapper<EtudiantDto> lineMapper() { DefaultLineMapper<EtudiantDto> lineMapper = new DefaultLineMapper<> (); DelimitedLineTokenizer lineTokenizer = new DelimitedLineTokenizer(); lineTokenizer.setDelimiter(","); lineTokenizer.setStrict(false); lineTokenizer.setNames("Username", "lastName", "firstName", "departement", "supervisor"); BeanWrapperFieldSetMapper<EmployeeDto> fieldSetMapper = new BeanWrapperFieldSetMapper<>(); fieldSetMapper.setTargetType(EmployeeDto.class); lineMapper.setLineTokenizer(lineTokenizer); lineMapper.setFieldSetMapper(fieldSetMapper); return lineMapper; } @Bean public EmployeeProcessor processor() { return new EmployeeProcessor(); /*Create a bean processor to skip invalid rows*/ } @Bean public RepositoryItemWriter<Employee> writer() { RepositoryItemWriter<Employee> writer = new RepositoryItemWriter<>(); writer.setRepository(employeeRepository); writer.setMethodName("save"); return writer; } @Bean public Step step1(FlatFileItemReader<EmployeeDto> itemReader) { return stepBuilderFactory.get("slaveStep").<EmployeeDto, Employee>chunk(5) .reader(itemReader) .processor(processor()) .writer(employeeItemWriter) .faultTolerant() .listener(skipListener()) .skip(SkipException.class) .skipLimit(10) .skipPolicy(skipPolicy()) .build(); } @Bean @Qualifier("executeJobEmployee") public Job runJob(FlatFileItemReader<Employee> itemReader) { return jobBuilderFactory .get("importEmployee") .flow(step1(itemReader)) .end() .build(); } @Bean public SkipPolicy skipPolicy(){ return new ExceptionSkipPolicy(); } @Bean public SkipListener<EmployeeDto, Employee> skipListener(){ return new StepSkipListener(); } /*@Bean public ExecutionContext executionContext(){ return new ExecutionContext(); }*/ }
EmployeeProcessor.java
public class EmployeeProcessor implements ItemProcessor<EmployeeDto, Employee>{ @Autowired private SupervisorService managerService; @Override public Employee process(@Valid EmployeeDto item) throws Exception, SkipException { ManagerDto manager = SupervisorService.findSupervisorById(item.getSupervisor()); //retrieve the manager of the employee and compare departement if(!(manager.getDepartement().equals(item.getDepartement()))) { throw new SkipException("Manager Invalid", item); //return null; } return ObjectMapperUtils.map(item, Employee.class); } }
MySkipPolicy.java
public class MySkipPolicy implements SkipPolicy { @Override public boolean shouldSkip(Throwable throwable, int i) throws SkipLimitExceededException { return true; } }
StepSkipListenerPolicy.java
public class StepSkipListener implements SkipListener<EmployeeDto, Number> { @Override // item reader public void onSkipInRead(Throwable throwable) { System.out.println("In OnSkipReader"); } @Override // item writer public void onSkipInWrite(Number item, Throwable throwable) { System.out.println("Nooooooooo "); } //@SneakyThrows @Override // item processor public void onSkipInProcess(@Valid EmployeeDto employee, Throwable throwable){ System.out.println("Process... "); /* I guess this is where I should work, but how do I deal with the exception occur? How do I know which exception I would get ? */ } }
SkipException.java
public class SkipException extends Exception { private Map<String, EmployeeDto> errors = new HashMap<>(); public SkipException(String errorMessage, EmployeeDto employee) { super(); this.errors.put(errorMessage, employee); } public Map<String, EmployeeDto> getErrors() { return this.errors; } }
JobController.java
@RestController @RequestMapping("/upload") public class JobController { @Autowired private JobLauncher jobLauncher; @Autowired @Qualifier("executeJobEmployee") private Job job; private final String EMPLOYEE_FOLDER = "C:/Users/Project/Employee/"; @PostMapping("/employee") public ResponseEntity<Object> importEmployee(@RequestParam("file") MultipartFile multipartFile) throws JobInterruptedException, SkipException, IllegalStateException, IOException, FlatFileParseException{ try { String fileName = multipartFile.getOriginalFilename(); File fileToImport= new File(EMPLOYEE_FOLDER + fileName); multipartFile.transferTo(fileToImport); JobParameters jobParameters = new JobParametersBuilder() .addString("fullPathFileName", EMPLOYEE_FOLDER + fileName) .addLong("startAt", System.currentTimeMillis()) .toJobParameters(); JobExecution jobExecution = this.jobLauncher.run(job, jobParameters); ExecutionContext executionContext = jobExecution.getExecutionContext(); System.out.println("My Skiped items : " + executionContext.toString()); } catch (ConstraintViolationException | FlatFileParseException | JobRestartException | JobInstanceAlreadyCompleteException | JobParametersInvalidException | JobExecutionAlreadyRunningException e) { e.printStackTrace(); return new ResponseEntity<>(e.getMessage(), HttpStatus.BAD_REQUEST); } return new ResponseEntity<>("Employee inserted succesfully", HttpStatus.OK); } }
解决方案
1. 改造SkipListener收集错误记录
让SkipListener负责收集所有跳过的记录及错误信息,通过ThreadLocal避免多线程冲突,最后存入JobExecution的ExecutionContext供Controller获取:
public class StepSkipListener implements SkipListener<EmployeeDto, Employee> { private final ThreadLocal<List<Map<String, Object>>> skippedRecords = ThreadLocal.withInitial(ArrayList::new); @Override public void onSkipInRead(Throwable throwable) { String errorMsg = "读取失败: " + throwable.getMessage(); Integer lineNum = throwable instanceof FlatFileParseException ? ((FlatFileParseException) throwable).getLineNumber() : null; skippedRecords.get().add(Map.of("error", errorMsg, "lineNumber", lineNum)); } @Override public void onSkipInWrite(Employee item, Throwable throwable) { String errorMsg = "写入失败: " + throwable.getMessage(); skippedRecords.get().add(Map.of("error", errorMsg, "employee", item)); } @Override public void onSkipInProcess(EmployeeDto employee, Throwable throwable) { String errorMsg; if (throwable instanceof SkipException) { errorMsg = ((SkipException) throwable).getMessage(); } else if (throwable instanceof ConstraintViolationException) { ConstraintViolationException violationEx = (ConstraintViolationException) throwable; errorMsg = violationEx.getConstraintViolations().stream() .map(v -> v.getPropertyPath() + ": " + v.getMessage()) .collect(Collectors.joining("; ")); } else { errorMsg = "处理失败: " + throwable.getMessage(); } skippedRecords.get().add(Map.of("error", errorMsg, "employeeData", employee)); } public List<Map<String, Object>> getSkippedRecords() { return skippedRecords.get(); } public void clear() { skippedRecords.remove(); } }
2. 配置Listener关联到Job和Step
修改EmployeeBatchConfig,注册JobExecutionListener用于在Job结束时将错误记录存入ExecutionContext,同时修正lineMapper的类型错误:
@Configuration @EnableBatchProcessing @AllArgsConstructor public class EmployeeBatchConfig { // ... 原有注入Bean @Bean public StepSkipListener skipListener(){ return new StepSkipListener(); } @Bean public JobExecutionListener jobExecutionListener(StepSkipListener skipListener) { return new JobExecutionListener() { @Override public void beforeJob(JobExecution jobExecution) { skipListener.clear(); } @Override public void afterJob(JobExecution jobExecution) { List<Map<String, Object>> skippedRecords = skipListener.getSkippedRecords(); if (!skippedRecords.isEmpty()) { jobExecution.getExecutionContext().put("skippedRecords", skippedRecords); } } }; } @Bean public Step step1(FlatFileItemReader<EmployeeDto> itemReader, StepSkipListener skipListener) { return stepBuilderFactory.get("slaveStep").<EmployeeDto, Employee>chunk(5) .reader(itemReader) .processor(processor()) .writer(writer()) .faultTolerant() .listener(skipListener) .skip(SkipException.class) .skip(ConstraintViolationException.class) .skipLimit(10) .skipPolicy(new MySkipPolicy()) .build(); } @Bean @Qualifier("executeJobEmployee") public Job runJob(Step step1, JobExecutionListener jobExecutionListener) { return jobBuilderFactory .get("importEmployee") .flow(step1) .end() .listener(jobExecutionListener) .build(); } private LineMapper<EmployeeDto> lineMapper() { DefaultLineMapper<EmployeeDto> lineMapper = new DefaultLineMapper<>(); DelimitedLineTokenizer lineTokenizer = new DelimitedLineTokenizer(); lineTokenizer.setDelimiter(","); lineTokenizer.setStrict(false); lineTokenizer.setNames("Username", "lastName", "firstName", "departement", "supervisor"); BeanWrapperFieldSetMapper<EmployeeDto> fieldSetMapper = new BeanWrapperFieldSetMapper<>(); fieldSetMapper.setTargetType(EmployeeDto.class); lineMapper.setLineTokenizer(lineTokenizer); lineMapper.setFieldSetMapper(fieldSetMapper); return lineMapper; } }
3. 修正Processor的验证逻辑
手动触发JSR380验证,确保实体校验生效:
@Component public class EmployeeProcessor implements ItemProcessor<EmployeeDto, Employee>{ @Autowired private SupervisorService managerService; @Autowired private Validator validator; @Override public Employee process(EmployeeDto item) throws Exception { Set<ConstraintViolation<EmployeeDto>> violations = validator.validate(item); if (!violations.isEmpty()) { throw new ConstraintViolationException(violations); } ManagerDto manager = managerService.findSupervisorById(item.getSupervisor()); if(!manager.getDepartement().equals(item.getDepartement())) { throw new SkipException("经理与员工不同部门", item); } return ObjectMapperUtils.map(item, Employee.class); } }
4. 调整SkipPolicy缩小跳过范围
避免跳过所有异常,只处理业务和验证相关异常:
public class MySkipPolicy implements SkipPolicy { @Override public boolean shouldSkip(Throwable throwable, int skipCount) throws SkipLimitExceededException { return throwable instanceof SkipException || throwable instanceof ConstraintViolationException; } }
5. 在Controller返回错误记录
从JobExecution中取出跳过的记录,组装成响应返回给客户端:
@RestController @RequestMapping("/upload") public class JobController { @Autowired private JobLauncher jobLauncher; @Autowired @Qualifier("executeJobEmployee") private Job job; private final String EMPLOYEE_FOLDER = "C:/Users/Project/Employee/"; @PostMapping("/employee") public ResponseEntity<Object> importEmployee(@RequestParam("file") MultipartFile multipartFile) { try { String fileName = multipartFile.getOriginalFilename(); File fileToImport= new File(EMPLOYEE_FOLDER + fileName); multipartFile.transferTo(fileToImport); JobParameters jobParameters = new JobParametersBuilder() .addString("fullPathFileName", EMPLOYEE_FOLDER + fileName) .addLong("startAt", System.currentTimeMillis()) .toJobParameters(); JobExecution jobExecution = jobLauncher.run(job, jobParameters); ExecutionContext executionContext = jobExecution.getExecutionContext(); List<Map<String, Object>> skippedRecords = (List<Map<String, Object>>) executionContext.get("skippedRecords"); if (skippedRecords != null && !skippedRecords.isEmpty()) { Map<String, Object> response = new HashMap<>(); response.put("message", "部分数据导入成功,存在错误记录"); response.put("skippedRecords", skippedRecords); return new ResponseEntity<>(response, HttpStatus.PARTIAL_CONTENT); } else { return new ResponseEntity<>("所有员工数据导入成功", HttpStatus.OK); } } catch (Exception e) { e.printStackTrace(); return new ResponseEntity<>("导入失败: " + e.getMessage(), HttpStatus.BAD_REQUEST); } } }
内容的提问来源于stack exchange,提问作者Suresh Vanalya
相关产品推荐
相关产品推荐

