Spring Batch分区与跳过策略:百万级CSV转Oracle并行方案咨询
Great question! Let's break this down step by step for your Spring Batch scenario involving million-row CSV data and Oracle database writes.
Q1: Multithreading (task-executor) vs. Partitioner – Which is Better?
First, let’s clarify the core differences between these two parallel approaches, since your use case deals with large-scale CSV-to-database processing:
Multithreading with
task-executor:- This runs a single Step's
ItemReader/Processor/Writeracross multiple threads. - The critical catch: Most CSV readers (like
FlatFileItemReader) are not thread-safe by default. Sharing one reader instance across threads will cause duplicate rows, missing data, or corrupted reads because the reader’s position pointer is shared. - You could configure a reader per thread (using
scope="step"), but this gets inefficient for million-row files—each thread would need to re-scan the file to find its starting point, wasting resources. - Best suited for small-to-medium datasets or scenarios where your reader is naturally thread-safe (e.g., reading from a partitioned database table).
- This runs a single Step's
Partitioner:
- This splits your input into distinct, non-overlapping partitions (e.g., split the CSV into chunks of N rows, or by logical keys if your data has them). Each partition is processed by its own independent Step instance.
- For CSV files, you can build a custom
Partitionerthat calculates start/end row positions for each partition. Each partition’sFlatFileItemReaderis configured to read only its assigned range, eliminating thread safety issues entirely. - When writing to Oracle, each partition can use optimized batch inserts (via
JdbcBatchItemWriter), which plays nicely with Oracle’s bulk insert capabilities. You also avoid write contention since each partition handles its own dataset. - This is the far better choice for your million-row scenario: It’s more scalable, avoids thread safety pitfalls, and aligns with batch processing best practices for large datasets.
Q2: Implementing Skip Policy with Partitioner to Skip & Log Failed Inserts
When using Partitioner, each partition’s Step can independently handle skips. Here’s how to set this up effectively:
1. Custom Skip Policy
Create a SkipPolicy to define which exceptions should trigger a skip (e.g., unique constraint violations, data type mismatches):
import org.springframework.batch.core.step.skip.SkipLimitExceededException; import org.springframework.batch.core.step.skip.SkipPolicy; import java.sql.SQLIntegrityConstraintViolationException; public class OracleInsertSkipPolicy implements SkipPolicy { private final int skipLimit; public OracleInsertSkipPolicy(int skipLimit) { this.skipLimit = skipLimit; } @Override public boolean shouldSkip(Throwable throwable, int skipCount) throws SkipLimitExceededException { // Skip constraint violations (adjust exceptions to match your failure cases) if (throwable instanceof SQLIntegrityConstraintViolationException && skipCount < skipLimit) { return true; } // Re-throw other exceptions to fail the step return false; } }
2. Skip Listener for Logging Skipped Records
Add a SkipListener to log details of skipped items (critical for debugging and data reconciliation):
import org.springframework.batch.core.SkipListener; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class CsvSkipListener<T> implements SkipListener<T, T> { private static final Logger logger = LoggerFactory.getLogger(CsvSkipListener.class); @Override public void onSkipInRead(Throwable throwable) { logger.error("Failed to read CSV row: {}", throwable.getMessage()); } @Override public void onSkipInWrite(T item, Throwable throwable) { logger.error("Skipped inserting item {} due to: {}", item.toString(), throwable.getMessage()); } @Override public void onSkipInProcess(T item, Throwable throwable) { logger.error("Skipped processing item {} due to: {}", item.toString(), throwable.getMessage()); } }
3. Custom CSV Partitioner
Implement a Partitioner to split your CSV into manageable chunks:
import org.springframework.batch.core.partition.support.Partitioner; import org.springframework.batch.item.ExecutionContext; import java.io.File; import java.io.IOException; import java.nio.file.Files; import java.util.HashMap; import java.util.Map; public class CsvPartitioner implements Partitioner { private String inputFile; private int rowsPerPartition; @Override public Map<String, ExecutionContext> partition(int gridSize) { Map<String, ExecutionContext> partitions = new HashMap<>(); try { long totalRows = Files.lines(new File(inputFile).toPath()).count() - 1; // Subtract header row int startRow = 1; // Skip header while (startRow <= totalRows) { int endRow = Math.min(startRow + rowsPerPartition - 1, (int) totalRows); ExecutionContext context = new ExecutionContext(); context.putInt("startRow", startRow); context.putInt("endRow", endRow); context.putString("inputFile", inputFile); partitions.put("partition-" + startRow, context); startRow += rowsPerPartition; } } catch (IOException e) { throw new RuntimeException("Failed to count CSV rows", e); } return partitions; } // Setters for configuration public void setInputFile(String inputFile) { this.inputFile = inputFile; } public void setRowsPerPartition(int rowsPerPartition) { this.rowsPerPartition = rowsPerPartition; } }
4. Batch Configuration
Wire everything together in your @Configuration class:
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.core.partition.PartitionHandler; import org.springframework.batch.core.partition.support.TaskExecutorPartitionHandler; 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.scheduling.concurrent.ThreadPoolTaskExecutor; import javax.sql.DataSource; @Configuration @EnableBatchProcessing public class BatchConfig { private final JobBuilderFactory jobBuilderFactory; private final StepBuilderFactory stepBuilderFactory; private final DataSource dataSource; public BatchConfig(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory, DataSource dataSource) { this.jobBuilderFactory = jobBuilderFactory; this.stepBuilderFactory = stepBuilderFactory; this.dataSource = dataSource; } @Bean public FlatFileItemReader<YourDataModel> csvItemReader(ExecutionContext context) { FlatFileItemReader<YourDataModel> reader = new FlatFileItemReader<>(); reader.setResource(new FileSystemResource(context.getString("inputFile"))); reader.setLinesToSkip(context.getInt("startRow")); reader.setMaxItemCount(context.getInt("endRow") - context.getInt("startRow") + 1); DefaultLineMapper<YourDataModel> lineMapper = new DefaultLineMapper<>(); DelimitedLineTokenizer tokenizer = new DelimitedLineTokenizer(); tokenizer.setNames("column1", "column2", "column3"); // Match your CSV columns lineMapper.setLineTokenizer(tokenizer); BeanWrapperFieldSetMapper<YourDataModel> fieldSetMapper = new BeanWrapperFieldSetMapper<>(); fieldSetMapper.setTargetType(YourDataModel.class); lineMapper.setFieldSetMapper(fieldSetMapper); reader.setLineMapper(lineMapper); return reader; } @Bean public JdbcBatchItemWriter<YourDataModel> oracleItemWriter() { JdbcBatchItemWriter<YourDataModel> writer = new JdbcBatchItemWriter<>(); writer.setDataSource(dataSource); writer.setSql("INSERT INTO your_table (col1, col2, col3) VALUES (:column1, :column2, :column3)"); writer.setItemSqlParameterSourceProvider(new BeanPropertyItemSqlParameterSourceProvider<>()); writer.setBatchSize(1000); // Adjust based on Oracle's optimal batch size return writer; } @Bean public Step slaveStep() { return stepBuilderFactory.get("slaveStep") .<YourDataModel, YourDataModel>chunk(1000) .reader(csvItemReader(null)) // Populated by partitioner at runtime .writer(oracleItemWriter()) .faultTolerant() .skipPolicy(new OracleInsertSkipPolicy(100)) // Allow up to 100 skips .listener(new CsvSkipListener<>()) .build(); } @Bean public PartitionHandler partitionHandler() { TaskExecutorPartitionHandler handler = new TaskExecutorPartitionHandler(); handler.setStep(slaveStep()); handler.setTaskExecutor(taskExecutor()); handler.setGridSize(4); // Number of parallel partitions (adjust to your server capacity) return handler; } @Bean public ThreadPoolTaskExecutor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(4); executor.setMaxPoolSize(4); executor.setThreadNamePrefix("batch-thread-"); return executor; } @Bean public Step masterStep() { return stepBuilderFactory.get("masterStep") .partitioner(slaveStep().getName(), csvPartitioner()) .partitionHandler(partitionHandler()) .build(); } @Bean public CsvPartitioner csvPartitioner() { CsvPartitioner partitioner = new CsvPartitioner(); partitioner.setInputFile("/path/to/your/large.csv"); partitioner.setRowsPerPartition(250000); // Split 1M rows into 4 partitions return partitioner; } @Bean public Job csvToOracleJob() { return jobBuilderFactory.get("csvToOracleJob") .start(masterStep()) .build(); } }
Key Notes:
- Partition Size: Adjust
rowsPerPartitionandgridSizebased on your server’s CPU/memory and Oracle’s capacity. 4 partitions of 250k rows each is a reasonable starting point for 1M rows. - Batch Size: Tune
batchSizeinJdbcBatchItemWriterto optimize Oracle bulk inserts (1000-5000 rows per batch typically works well). - Skip Limit: Set
skipLimitto prevent infinite skips if there’s a systemic issue with your data. - Logging: The
CsvSkipListenerensures you have a clear record of every skipped item for post-job reconciliation.
内容的提问来源于stack exchange,提问作者hafs

