Spring Batch任务架构咨询与实现示例:MySQL转CSV传GCS
Spring Batch 架构设计与实现指导
核心问题解答
1. CSV 格式转换:用 Writer,而非 Processor
Spring Batch 组件的职责边界明确:
- Processor:负责数据的业务转换、过滤、校验(比如字段映射、数据清洗),不承担持久化/格式序列化工作。
- Writer:专门负责将处理后的数据写入目标介质,Spring Batch 提供的
FlatFileItemWriter是为文本/CSV 输出设计的标准组件,内置分隔符、引号、表头配置,完全适配 CSV 生成需求,无需用 Processor 做格式转换。
2. 是否需要分两步?两种方案可选
方案一:单步完成(推荐,高效)
直接从 MySQL 读取数据,用 FlatFileItemWriter 写入本地临时 CSV,再通过 Step 监听回调完成 GCS 上传。这种方式避免二次读取文件的 IO 开销,适合大多数常规场景。
方案二:分两步执行(适用于大文件/中间校验场景)
第一步:MySQL → CSV 文件;第二步:读取本地 CSV → 上传 GCS。适合需要对生成的 CSV 做校验、或文件体积过大内存无法承载的情况,Spring Batch 的分区/分片机制也能更好适配大文件并行处理。
实现示例
依赖配置(Maven)
<dependencies> <!-- Spring Batch 核心 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-batch</artifactId> </dependency> <!-- MySQL 驱动 --> <dependency> <groupId>com.mysql</groupId> <artifactId>mysql-connector-j</artifactId> </dependency> <!-- Spring Cloud GCP GCS 支持 --> <dependency> <groupId>com.google.cloud</groupId> <artifactId>spring-cloud-gcp-starter-storage</artifactId> </dependency> </dependencies>
方案一:单步实现(MySQL → CSV → GCS)
1. 实体类(以用户表为例)
public class User { private Long id; private String username; private String email; // 构造器、getter/setter 省略 }
2. MySQL 数据读取器
@Bean public JdbcCursorItemReader<User> userReader(DataSource dataSource) { return new JdbcCursorItemReaderBuilder<User>() .dataSource(dataSource) .sql("SELECT id, username, email FROM users") .rowMapper((rs, rowNum) -> new User( rs.getLong("id"), rs.getString("username"), rs.getString("email") )) .name("userReader") .build(); }
3. CSV 文件写入器
@Bean public FlatFileItemWriter<User> csvWriter() { String tempFilePath = System.getProperty("java.io.tmpdir") + "/users.csv"; return new FlatFileItemWriterBuilder<User>() .resource(new FileSystemResource(tempFilePath)) .name("csvWriter") .delimited() .delimiter(",") .names("id", "username", "email") .headerCallback(writer -> writer.write("id,username,email")) .build(); }
4. GCS 上传监听器
@Component public class GcsUploadListener implements StepExecutionListener { private final Storage storage; private final String gcsBucketName = "your-bucket-name"; // 替换为你的 GCS 桶名 public GcsUploadListener(Storage storage) { this.storage = storage; } @Override public ExitStatus afterStep(StepExecution stepExecution) { String tempFilePath = System.getProperty("java.io.tmpdir") + "/users.csv"; File csvFile = new File(tempFilePath); if (csvFile.exists()) { try { BlobId blobId = BlobId.of(gcsBucketName, "exports/users.csv"); BlobInfo blobInfo = BlobInfo.newBuilder(blobId).build(); storage.create(blobInfo, Files.readAllBytes(csvFile.toPath())); csvFile.delete(); return ExitStatus.COMPLETED; } catch (IOException e) { return ExitStatus.FAILED.addExitDescription(e.getMessage()); } } return ExitStatus.FAILED.addExitDescription("CSV 文件未生成"); } }
5. 批处理 Job 配置
@Configuration @EnableBatchProcessing public class BatchConfig { private final JobBuilderFactory jobBuilderFactory; private final StepBuilderFactory stepBuilderFactory; public BatchConfig(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory) { this.jobBuilderFactory = jobBuilderFactory; this.stepBuilderFactory = stepBuilderFactory; } @Bean public Step exportAndUploadStep(JdbcCursorItemReader<User> userReader, FlatFileItemWriter<User> csvWriter, GcsUploadListener gcsUploadListener) { return stepBuilderFactory.get("exportAndUploadStep") .<User, User>chunk(1000) .reader(userReader) .writer(csvWriter) .listener(gcsUploadListener) .build(); } @Bean public Job exportUserJob(Step exportAndUploadStep) { return jobBuilderFactory.get("exportUserJob") .start(exportAndUploadStep) .build(); } }
方案二:分两步实现
1. 第一步:MySQL → CSV(与方案一的 reader、writer 一致)
2. 第二步:读取 CSV → 上传 GCS
// CSV 文件读取器 @Bean public FlatFileItemReader<List<String>> csvFileReader() { String tempFilePath = System.getProperty("java.io.tmpdir") + "/users.csv"; return new FlatFileItemReaderBuilder<List<String>>() .resource(new FileSystemResource(tempFilePath)) .name("csvFileReader") .lineMapper(new DefaultLineMapper<List<String>>() {{ setLineTokenizer(new DelimitedLineTokenizer() {{ setDelimiter(","); setNames("id", "username", "email"); }}); setFieldSetMapper(fieldSet -> Arrays.asList( fieldSet.readString("id"), fieldSet.readString("username"), fieldSet.readString("email") )); }}) .skippedLinesCallback(line -> {}) // 跳过表头 .build(); } // GCS 上传写入器 @Bean public ItemWriter<List<String>> gcsWriter(Storage storage) { String bucketName = "your-bucket-name"; String gcsFilePath = "exports/users.csv"; BlobId blobId = BlobId.of(bucketName, gcsFilePath); BlobInfo blobInfo = BlobInfo.newBuilder(blobId).build(); WriteChannel writeChannel = storage.writer(blobInfo); return items -> { try (BufferedWriter writer = new BufferedWriter(Channels.newWriter(writeChannel, StandardCharsets.UTF_8.name()))) { writer.write("id,username,email\n"); for (List<String> item : items) { writer.write(String.join(",", item) + "\n"); } } }; } // 两步 Job 配置 @Bean public Step exportToCsvStep(JdbcCursorItemReader<User> userReader, FlatFileItemWriter<User> csvWriter) { return stepBuilderFactory.get("exportToCsvStep") .<User, User>chunk(1000) .reader(userReader) .writer(csvWriter) .build(); } @Bean public Step uploadToGcsStep(FlatFileItemReader<List<String>> csvFileReader, ItemWriter<List<String>> gcsWriter) { return stepBuilderFactory.get("uploadToGcsStep") .<List<String>, List<String>>chunk(1000) .reader(csvFileReader) .writer(gcsWriter) .build(); } @Bean public Job exportUserTwoStepJob(Step exportToCsvStep, Step uploadToGcsStep) { return jobBuilderFactory.get("exportUserTwoStepJob") .start(exportToCsvStep) .next(uploadToGcsStep) .build(); }
注意事项
- 临时文件管理:上传完成后务必删除本地临时文件,避免磁盘占用。
- 大文件处理:大体积 CSV 建议使用 GCS 分块上传,或结合 Spring Batch 分区机制并行处理。
- 异常处理:添加重试、跳过策略,比如读取 MySQL 失败自动重试,写入时跳过无效数据。
内容的提问来源于stack exchange,提问作者Pedro Buttenbender
相关产品推荐
相关产品推荐

