Spring Batch分区作业异常处理:按类型写入不同错误文件
Absolutely! You can absolutely implement this requirement using either Spring Batch Listeners or a custom error-handling Writer, and even categorize failed records into different files based on exception types. Let me break down how to approach this step by step, with a focus on your partitioned, multi-threaded job scenario:
1. Use SkipListeners (Batch-Native Fault Tolerance)
Spring Batch's skip mechanism is tailor-made for handling failed records without aborting the entire job. Pair it with a SkipListener to capture skipped records and route them to different error files based on exception type.
Step 1: Configure Skip Rules in Your Partitioned Step
First, enable skipping for the exceptions you want to handle, and attach your listener:
@Bean public Step partitionedProcessingStep(StepBuilderFactory stepBuilderFactory, ItemReader<InputRecord> itemReader, ItemProcessor<InputRecord, OutputRecord> itemProcessor, ItemWriter<OutputRecord> itemWriter, CategorizedSkipListener categorizedSkipListener) { return stepBuilderFactory.get("partitionedProcessingStep") .<InputRecord, OutputRecord>chunk(100) .reader(itemReader) .processor(itemProcessor) .writer(itemWriter) .faultTolerant() // Define which exceptions to skip (adjust based on your needs) .skip(ValidationException.class) .skip(IOException.class) .skip(Exception.class) // Catch-all for other exceptions .skipLimit(Integer.MAX_VALUE) // Skip as many as needed (adjust if you want a cap) .listener(categorizedSkipListener) .build(); }
Step 2: Create a Categorized Skip Listener
This listener will catch skipped records, wrap them with error details, and route them to the correct error file based on the exception type:
@Component @Scope("step") // Critical! Ensures each partition gets its own listener instance (thread-safe) public class CategorizedSkipListener implements SkipListener<InputRecord, OutputRecord> { private static final Logger log = LoggerFactory.getLogger(CategorizedSkipListener.class); private final Map<Class<? extends Exception>, ItemWriter<ErrorRecord>> errorWriterMap; // Inject pre-configured writers for each exception type public CategorizedSkipListener(Map<Class<? extends Exception>, ItemWriter<ErrorRecord>> errorWriterMap) { this.errorWriterMap = errorWriterMap; } @Override public void onSkipInProcess(InputRecord failedRecord, Throwable exception) { wrapAndWriteErrorRecord(failedRecord, exception); } @Override public void onSkipInWrite(OutputRecord failedRecord, Throwable exception) { // If the error happens during writing, adjust the wrapped record as needed wrapAndWriteErrorRecord(convertOutputToInput(failedRecord), exception); } @Override public void onSkipInRead(Throwable exception) { // Handle read errors if needed (e.g., log invalid lines in the CSV) log.error("Failed to read record: {}", exception.getMessage()); } private void wrapAndWriteErrorRecord(InputRecord originalRecord, Throwable exception) { ErrorRecord errorRecord = new ErrorRecord( originalRecord, exception.getMessage(), exception.getClass().getSimpleName() ); // Get the writer for this exception type, or fall back to a default ItemWriter<ErrorRecord> targetWriter = errorWriterMap.getOrDefault( exception.getClass(), errorWriterMap.get(Exception.class) ); try { targetWriter.write(Collections.singletonList(errorRecord)); } catch (Exception writeError) { log.error("Failed to write error record to file: {}", originalRecord, writeError); } } // Helper method if you need to convert output records back to input for error logging private InputRecord convertOutputToInput(OutputRecord outputRecord) { // Implement based on your data models return new InputRecord(); } }
Step 3: Configure Error File Writers
Create separate FlatFileItemWriter beans for each exception type, and map them to their corresponding exceptions. For partitioned jobs, use dynamic filenames to avoid thread conflicts:
@Bean @StepScope public FlatFileItemWriter<ErrorRecord> validationErrorWriter(@Value("#{stepExecutionContext['partitionId']}") String partitionId) { return buildErrorWriter(String.format("validation_errors_partition_%s.csv", partitionId)); } @Bean @StepScope public FlatFileItemWriter<ErrorRecord> ioErrorWriter(@Value("#{stepExecutionContext['partitionId']}") String partitionId) { return buildErrorWriter(String.format("io_errors_partition_%s.csv", partitionId)); } @Bean @StepScope public FlatFileItemWriter<ErrorRecord> defaultErrorWriter(@Value("#{stepExecutionContext['partitionId']}") String partitionId) { return buildErrorWriter(String.format("general_errors_partition_%s.csv", partitionId)); } // Helper method to avoid duplicate writer configuration private FlatFileItemWriter<ErrorRecord> buildErrorWriter(String fileName) { FlatFileItemWriter<ErrorRecord> writer = new FlatFileItemWriter<>(); writer.setResource(new FileSystemResource(fileName)); writer.setAppendAllowed(true); // Append to existing files if needed writer.setLineAggregator(new DelimitedLineAggregator<>() {{ setDelimiter(","); setFieldExtractor(new BeanWrapperFieldExtractor<>() {{ setNames(new String[]{"originalRecord", "errorMessage", "exceptionType"}); }}); }}); return writer; } @Bean public Map<Class<? extends Exception>, ItemWriter<ErrorRecord>> errorWriterMap( FlatFileItemWriter<ErrorRecord> validationErrorWriter, FlatFileItemWriter<ErrorRecord> ioErrorWriter, FlatFileItemWriter<ErrorRecord> defaultErrorWriter) { Map<Class<? extends Exception>, ItemWriter<ErrorRecord>> map = new HashMap<>(); map.put(ValidationException.class, validationErrorWriter); map.put(IOException.class, ioErrorWriter); map.put(Exception.class, defaultErrorWriter); return map; }
2. Alternative: Custom ItemProcessListener + Error Writer
If you don't want to use the skip mechanism, you can use an ItemProcessListener to catch exceptions during processing, then route the failed record to your categorized error writers directly. This works similarly to the SkipListener approach, but you'll need to handle skipping the failed record manually (e.g., return null from the processor to skip it from the writer).
3. Key Considerations for Partitioned Jobs
- Thread Safety: Use
@StepScopefor all listeners and writers involved in error handling. This ensures each partition gets its own instance, preventing race conditions when writing to files. - Dynamic Filenames: Include the partition ID in error filenames (as shown above) to avoid multiple threads writing to the same file simultaneously.
- Error File Cleanup: If needed, add a step to merge partition-specific error files into a single file after the job completes.
- Exception Granularity: Be specific with your exception types (avoid catching
Exception.classunless necessary) to ensure records are routed to the correct error files.
内容的提问来源于stack exchange,提问作者bdetweiler

