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

Spring Batch分区与跳过策略:百万级CSV转Oracle并行方案咨询

Spring Batch: Parallel Processing for Million-Row CSV to Oracle + Skip Policy with Partitioner

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/Writer across 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).
  • 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 Partitioner that calculates start/end row positions for each partition. Each partition’s FlatFileItemReader is 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 rowsPerPartition and gridSize based 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 batchSize in JdbcBatchItemWriter to optimize Oracle bulk inserts (1000-5000 rows per batch typically works well).
  • Skip Limit: Set skipLimit to prevent infinite skips if there’s a systemic issue with your data.
  • Logging: The CsvSkipListener ensures you have a clear record of every skipped item for post-job reconciliation.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 08:08:32