能否结合Spring Async与Java8 Parallel Stream处理CSV文件?
Solution for Parallel CSV Processing with Scheduler
Core Workflow
We'll implement parallelism at two levels:
- Inter-file parallelism: Use Spring's
@Asyncto process multiple CSV files simultaneously. - Intra-file parallelism: Split records into chunks of 1000 when a file's record count exceeds this threshold, then process chunks in parallel.
The full pipeline:
- Fetch FTP file paths from the database
- Download and parse each file in separate threads
- Split large record sets into manageable chunks
- Validate and persist each chunk in parallel
- Save validated data to the target database table
Key Implementation Details
- Async File Processing: Annotate the file handler method with
@Asyncto offload each file's processing to a separate thread (ensure@EnableAsyncis enabled in your Spring config). - Chunking Logic: Split records into sublists of 1000 entries to avoid overwhelming the database and leverage parallel processing for large files.
- Parallel Chunk Processing: Use
parallelStreamto process record chunks concurrently, balancing performance and resource usage.
Complete Code Example
import org.springframework.scheduling.annotation.Async; import org.springframework.stereotype.Service; import java.io.InputStream; import java.util.List; import java.util.stream.Collectors; @Service public class CSVProcessorService { private final CSVFileRecordsRepository filesRepo; private final CSVParser csvParser; private final DataRepository dataRepo; // Constructor injection for dependencies public CSVProcessorService(CSVFileRecordsRepository filesRepo, CSVParser csvParser, DataRepository dataRepo) { this.filesRepo = filesRepo; this.csvParser = csvParser; this.dataRepo = dataRepo; } // Scheduler entry point to trigger processing public void triggerCSVProcessing() { List<CSVFileRecords> files = filesRepo.findAll(); files.forEach(this::processSingleFile); } @Async public void processSingleFile(CSVFileRecords file) { // Download CSV from FTP InputStream fileStream = fetchFtpFile(file.getFtpUrl()); // Parse CSV into list of Data objects List<Data> allRecords = csvParser.csvToBean(fileStream); // Split records into chunks of 1000 (adjustable) int chunkSize = 1000; List<List<Data>> recordChunks = splitRecordsIntoChunks(allRecords, chunkSize); // Process each chunk in parallel recordChunks.parallelStream().forEach(this::processRecordChunk); } private void processRecordChunk(List<Data> chunk) { // Validate records in the chunk List<Data> validRecords = chunk.stream() .filter(this::validateDataRecord) .collect(Collectors.toList()); // Batch save valid records to database dataRepo.saveAll(validRecords); } private boolean validateDataRecord(Data data) { // Implement custom validation logic (e.g., non-null fields, format checks) return data != null && data.getRequiredField() != null && !data.getRequiredField().isBlank(); } private List<List<Data>> splitRecordsIntoChunks(List<Data> records, int chunkSize) { return records.stream() .collect(Collectors.groupingBy(index -> index / chunkSize)) .values() .stream() .collect(Collectors.toList()); } private InputStream fetchFtpFile(String ftpUrl) { // Implement FTP download logic (handle authentication, connections) return null; // Replace with actual implementation } }
Critical Notes
- Configure a custom thread pool for
@Asyncto control the maximum number of parallel file processing threads (prevents resource exhaustion). - Add exception handling in
processSingleFileandprocessRecordChunkto isolate failures (one bad file/chunk won't break the entire batch). - Adjust
chunkSizebased on your database's batch insert limits and server performance.
内容的提问来源于stack exchange,提问作者G10
相关产品推荐
相关产品推荐

