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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 05:30:53