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

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 @StepScope for 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.class unless necessary) to ensure records are routed to the correct error files.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:13:01