创建 Spring Batch 批处理服务

创建 Spring Batch 批处理服务

本指南带你创建一个基础批处理解决方案。

将要构建什么

构建一个服务,从 CSV 表格导入数据,用自定义代码转换,再把最终结果存入数据库。

需要什么

• 约15分钟。

• 熟悉的文本编辑器或 IDE。

• Java 17 或更新版本。

• Gradle 7.5+ 或 Maven 3.5+。

• 也可直接把代码导入 Spring Tool Suite(STS)、IntelliJ IDEA 或 VSCode。

如何完成指南

与其他 Spring 入门指南一样,你可以从零开始逐步完成,也可以跳过已熟悉的基础设置;两条路径最终都得到可工作的代码。

从零开始时,继续阅读“使用 Spring Initializr”。

若要跳过基础步骤:

• 下载并解压源代码,或执行 git clone https://github.com/spring-guides/gs-batch-processing.git。

• 进入 gs-batch-processing/initial。

• 跳到“业务数据”一节。

完成后,可与 gs-batch-processing/complete 中的代码比较。

使用 Spring Initializr

可以打开按本指南预先配置的 Spring Initializr 项目,点击 Generate 下载 ZIP;也可以按以下步骤手动初始化。

手动初始化步骤:

• 打开 Spring Initializr,它会引入所需依赖并完成大部分设置。

• 选择 Gradle 或 Maven,并选择编程语言。本指南假设选择 Java。

• 点击 Dependencies,选择 Spring Batch JDBC 和 HyperSQL Database。

• 点击 Generate。

• 下载生成的 ZIP,其中包含按选项配置好的应用。

> 若 IDE 集成了 Spring Initializr,可以直接在 IDE 中完成这些步骤。

> 也可以 fork GitHub 项目,再用 IDE 或其他编辑器打开。

业务数据

通常,客户或业务分析师会提供表格。这个简单示例使用虚构数据,文件为 src/main/resources/sample-data.csv:

Jill,Doe
Joe,Doe
Justin,Doe
Jane,Doe
John,Doe

每行包含名字与姓氏,以逗号分隔。这种常见格式无需定制即可由 Spring 处理。

接着编写 SQL 脚本,创建存放数据的表。文件为 src/main/resources/schema-all.sql:

DROP TABLE people IF EXISTS;

CREATE TABLE people  (
    person_id BIGINT IDENTITY NOT NULL PRIMARY KEY,
    first_name VARCHAR(20),
    last_name VARCHAR(20)
);

> Spring Boot 启动时自动运行 schema-@@platform@@.sql,其中 -all 是适用于所有平台的默认值。

创建业务类

了解输入和输出格式后,可以用代码表示一行数据。以下是 src/main/java/com/example/batchprocessing/Person.java:

package com.example.batchprocessing;

public record Person(String firstName, String lastName) {

}

可通过构造函数,传入名字和姓氏来实例化 Person record。

创建中间处理器

批处理常见模式是读取数据、转换,再输出到其他地方。这里编写一个简单转换器,将姓名改为大写。以下是 src/main/java/com/example/batchprocessing/PersonItemProcessor.java:

package com.example.batchprocessing;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import org.springframework.batch.infrastructure.item.ItemProcessor;
public class PersonItemProcessor implements ItemProcessor<Person, Person> {

  private static final Logger log = LoggerFactory.getLogger(PersonItemProcessor.class);

  @Override
  public Person process(final Person person) {
    final String firstName = person.firstName().toUpperCase();
    final String lastName = person.lastName().toUpperCase();

    final Person transformedPerson = new Person(firstName, lastName);
    log.info("Converting ({}) into ({})", person, transformedPerson);

    return transformedPerson;
  }

}

PersonItemProcessor 实现 Spring Batch 的 ItemProcessor 接口,便于接入稍后定义的批处理作业。按该接口约定,它接收 Person,转换后返回姓名大写的 Person。

> 输入类型与输出类型不必相同。读取某个数据源后,应用的数据流有时需要另一种数据类型。

组装批处理作业

接下来组装真正的作业。Spring Batch 提供许多工具类,减少定制代码,让你专注于业务逻辑。

先创建 Spring @Configuration 类,例如 src/main/java/com/example/batchprocessing/BatchConfiguration.java。本例使用内存数据库,程序结束后数据也会消失。向 BatchConfiguration 添加以下 bean,定义 reader、processor 和 writer:

@Bean
public FlatFileItemReader<Person> reader() {
  return new FlatFileItemReaderBuilder<Person>()
    .name("personItemReader")
    .resource(new ClassPathResource("sample-data.csv"))
    .delimited()
    .names("firstName", "lastName")
    .targetType(Person.class)
    .build();
}

@Bean
public PersonItemProcessor processor() {
  return new PersonItemProcessor();
}

@Bean
public JdbcBatchItemWriter<Person> writer(DataSource dataSource) {
  return new JdbcBatchItemWriterBuilder<Person>()
    .sql("INSERT INTO people (first_name, last_name) VALUES (:firstName, :lastName)")
    .dataSource(dataSource)
    .beanMapped()
    .build();
}

第一段代码定义输入、处理器和输出。

• reader() 创建 ItemReader,查找 sample-data.csv,逐行解析出构造 Person 所需的信息。

• processor() 创建此前定义的 PersonItemProcessor,用于把数据改为大写。

• writer(DataSource) 创建面向 JDBC 的 ItemWriter,自动取得 Spring Boot 创建的 DataSource;包含插入单个 Person 的 SQL,参数由 Java record 的组件提供。

BatchConfiguration.java 的下一段定义实际作业:

@Bean
public Job importUserJob(JobRepository jobRepository, Step step1, JobCompletionNotificationListener listener) {
  return new JobBuilder(jobRepository)
    .listener(listener)
    .start(step1)
    .build();
}

@Bean
public Step step1(JobRepository jobRepository, DataSourceTransactionManager transactionManager,
          FlatFileItemReader<Person> reader, PersonItemProcessor processor, JdbcBatchItemWriter<Person> writer) {
  return new StepBuilder(jobRepository)
    .<Person, Person>chunk(3)
          .transactionManager(transactionManager)
    .reader(reader)
    .processor(processor)
    .writer(writer)
    .build();
}

第一个方法定义作业,第二个方法定义单个步骤。作业由步骤构成,每个步骤可包含 reader、processor 和 writer。

然后按顺序列出各步骤,本例只有一个步骤;结束构建时,Java API 会生成配置完整的作业。

步骤定义规定每次写入多少数据,本例最多一次写三条。随后使用注入的 bean 配置 reader、processor 和 writer。

> chunk() 前缀为 <Person,Person>,因为这是泛型方法,表示每个处理块的输入与输出类型,与 ItemReader<Person> 和 ItemWriter<Person> 对应。

最后需要在作业结束时收到通知。以下为 src/main/java/com/example/batchprocessing/JobCompletionNotificationListener.java:

package com.example.batchprocessing;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import org.springframework.batch.core.BatchStatus;
import org.springframework.batch.core.job.JobExecution;
import org.springframework.batch.core.listener.JobExecutionListener;
import org.springframework.jdbc.core.DataClassRowMapper;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;

@Component
public class JobCompletionNotificationListener implements JobExecutionListener {

  private static final Logger log = LoggerFactory.getLogger(JobCompletionNotificationListener.class);

  private final JdbcTemplate jdbcTemplate;

  public JobCompletionNotificationListener(JdbcTemplate jdbcTemplate) {
    this.jdbcTemplate = jdbcTemplate;
  }

  @Override
  public void afterJob(JobExecution jobExecution) {
    if (jobExecution.getStatus() == BatchStatus.COMPLETED) {
      log.info("!!! JOB FINISHED! Time to verify the results");

      jdbcTemplate
          .query("SELECT first_name, last_name FROM people", new DataClassRowMapper<>(Person.class))
          .forEach(person -> log.info("Found <{}> in the database.", person));
    }
  }
}

JobCompletionNotificationListener 等待作业进入 BatchStatus.COMPLETED,然后用 JdbcTemplate 检查结果。

让应用可执行

批处理可以嵌入 Web 应用或 WAR;这里采用更简单的独立应用方式,将所有内容打包到单个可执行 JAR,用普通 Java main() 方法启动。

Spring Initializr 已创建应用类,本例无需进一步修改。以下是 src/main/java/com/example/batchprocessing/BatchProcessingApplication.java:

package com.example.batchprocessing;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;

@SpringBootApplication
public class BatchProcessingApplication {

  public static void main(String[] args) {
    System.exit(SpringApplication.exit(SpringApplication.run(BatchProcessingApplication.class, args)));
  }
}

@SpringBootApplication 是便利注解,组合了以下功能:

• @Configuration:把该类标记为应用上下文的 bean 定义来源。

• @EnableAutoConfiguration:让 Spring Boot 根据类路径、其他 bean 和属性设置添加 bean。例如,类路径存在 spring-webmvc 时,会把应用视为 Web 应用,并启用 DispatcherServlet 等关键行为。

• @ComponentScan:让 Spring 搜索 com/example 包中的其他组件、配置和服务,从而找到控制器。

main() 使用 SpringApplication.run() 启动应用。整个示例没有一行 XML,也没有 web.xml,完全使用 Java,不必手工配置基础设施。

注意:SpringApplication.exit() 和 System.exit() 确保作业结束后 JVM 退出。详情参见 Spring Boot 应用退出。

为了演示,示例注入 JdbcTemplate,查询数据库并打印批处理插入的人员姓名。

> 应用没有使用 @EnableBatchProcessing。以前可以用该注解启用 Spring Boot 对 Spring Batch 的自动配置;现在,定义带该注解的 bean 或扩展 DefaultBatchConfiguration,会让自动配置退让,由应用完全控制 Spring Batch 配置。

构建可执行 JAR

可通过 Gradle 或 Maven 从命令行运行,也可以构建包含所有依赖、类和资源的可执行 JAR。这样便于在整个开发周期中跨环境交付、管理版本和部署服务。

使用 Gradle 时,运行 ./gradlew bootRun;也可以执行 ./gradlew build 构建 JAR,再这样运行:

java -jar build/libs/gs-batch-processing-0.0.1-SNAPSHOT.jar

使用 Maven 时,运行 ./mvnw spring-boot:run;也可以执行 ./mvnw clean package 构建,再这样运行:

java -jar target/gs-batch-processing-0.0.1-SNAPSHOT.jar

作业为每个转换的人输出一行日志。结束后还会显示数据库查询结果。原文预期输出如下:

Converting (Person[firstName=Jill, lastName=Doe]) into (Person[firstName=JILL, lastName=DOE])
Converting (Person[firstName=Joe, lastName=Doe]) into (Person[firstName=JOE, lastName=DOE])
Converting (Person[firstName=Justin, lastName=Doe]) into (Person[firstName=JUSTIN, lastName=DOE])
Converting (Person[firstName=Jane, lastName=Doe]) into (Person[firstName=JANE, lastName=DOE])
Converting (Person[firstName=John, lastName=Doe]) into (Person[firstName=JOHN, lastName=DOE])
Found <Person[firstName=JILL, lastName=DOE]> in the database.
Found <Person[firstName=JOE, lastName=DOE]> in the database.
Found <Person[firstName=JUSTIN, lastName=DOE]> in the database.
Found <Person[firstName=JANE, lastName=DOE]> in the database.
Found <Person[firstName=JOHN, lastName=DOE]> in the database.

小结

至此,你已构建一个从表格读取数据、进行处理并写入数据库的批处理作业。

延伸阅读

以下指南也可能有帮助:

• 使用 Spring Boot 构建应用

• 使用 GemFire 访问数据

• 使用 JPA 访问数据

• 使用 MongoDB 访问数据

• 使用 MySQL 访问数据

想编写新指南或参与现有指南,可参见 贡献指南。


来源:Creating a Batch Service。© 2005–2026 Broadcom 及贡献者。代码采用 Apache 2.0;原文写作采用署名、禁止演绎的 Creative Commons 许可证(CC BY-ND)。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容