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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 01:57:39