分区与并行步骤执行
通过多线程步骤、分区和远程分块策略提升吞吐量。
分区与并行步骤执行 是 CoddyKit 上的免费 Spring Boot 4 Complete Guide 课时。 这是第 4 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Spring Boot 4 Complete Guide 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Spring Boot 4 Complete Guide 课程共包含 4 节课。
本课时的部分内容尚未翻译,以英文显示。
Why Scale Spring Batch?
Single-threaded Spring Batch jobs process one chunk at a time — fine for small datasets, but too slow for millions of records. When batch throughput becomes a bottleneck, Spring Batch offers four scaling strategies:
- Multi-threaded Step — parallel threads within a single JVM step
- Parallel Steps — independent steps run concurrently in a flow
- Partitioning — divide data into partitions, each processed by a worker step
- Remote Chunking — offload chunk processing to remote workers over messaging middleware
Each strategy has different complexity/throughput trade-offs. This lesson covers all four, starting with the simplest.
Multi-Threaded Steps with TaskExecutor
The easiest way to add parallelism is to inject a TaskExecutor into your Step. Spring Batch will execute chunks concurrently on a thread pool. Important: the ItemReader must be thread-safe (e.g. stateless, or use SynchronizedItemStreamReader).
Configure a multi-threaded step like this:
import org.springframework.batch.core.Step;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.batch.item.file.FlatFileItemReader;
import org.springframework.batch.item.file.builder.FlatFileItemReaderBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.io.ClassPathResource;
import org.springframework.core.task.SimpleAsyncTaskExecutor;
import org.springframework.transaction.PlatformTransactionManager;
@Configuration
public class MultiThreadedStepConfig {
@Bean
public Step multiThreadedStep(JobRepository jobRepository,
PlatformTransactionManager txManager,
FlatFileItemReader<String> reader) {
return new StepBuilder("multiThreadedStep", jobRepository)
.<String, String>chunk(100, txManager)
.reader(reader)
.writer(items -> items.forEach(System.out::println))
.taskExecutor(new SimpleAsyncTaskExecutor())
.throttleLimit(4) // max concurrent threads
.build();
}
}Thread-Safe Readers with SynchronizedItemStreamReader
Standard FlatFileItemReader is not thread-safe because it maintains internal state (current line position). Wrapping it in SynchronizedItemStreamReader serialises read() calls, making it safe for multi-threaded steps without altering processing or writing parallelism.
import org.springframework.batch.item.file.FlatFileItemReader;
import org.springframework.batch.item.file.builder.FlatFileItemReaderBuilder;
import org.springframework.batch.item.support.SynchronizedItemStreamReader;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.io.ClassPathResource;
@Configuration
public class SafeReaderConfig {
@Bean
public SynchronizedItemStreamReader<String> synchronizedReader() {
FlatFileItemReader<String> delegate = new FlatFileItemReaderBuilder<String>()
.name("lineReader")
.resource(new ClassPathResource("data/input.csv"))
.lineMapper((line, lineNumber) -> line)
.build();
SynchronizedItemStreamReader<String> reader = new SynchronizedItemStreamReader<>();
reader.setDelegate(delegate);
return reader;
}
}Parallel Steps with Split Flows
When you have independent steps (e.g. loading products and loading customers simultaneously), use a split() flow so both steps run concurrently. Spring Batch's FlowBuilder supports this natively.
Steps inside a split share no data — they must operate on separate resources to avoid contention.
import org.springframework.batch.core.Job;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.job.builder.FlowBuilder;
import org.springframework.batch.core.job.builder.JobBuilder;
import org.springframework.batch.core.job.flow.Flow;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.task.SimpleAsyncTaskExecutor;
@Configuration
public class ParallelStepsConfig {
@Bean
public Flow productFlow(Step loadProductsStep) {
return new FlowBuilder<Flow>("productFlow")
.start(loadProductsStep)
.build();
}
@Bean
public Flow customerFlow(Step loadCustomersStep) {
return new FlowBuilder<Flow>("customerFlow")
.start(loadCustomersStep)
.build();
}
@Bean
public Job parallelJob(JobRepository jobRepository,
Flow productFlow,
Flow customerFlow) {
return new JobBuilder("parallelJob", jobRepository)
.start(productFlow)
.split(new SimpleAsyncTaskExecutor())
.add(customerFlow)
.end()
.build();
}
}Introduction to Partitioning
Partitioning divides a dataset into non-overlapping partitions, each processed independently by a worker step. A manager step (formerly called master) creates the partitions and delegates them.
Partitioner— createsExecutionContextmaps, one per partitionPartitionHandler— decides how worker steps are launched (local or remote)TaskExecutorPartitionHandler— runs workers locally in a thread pool
Each worker step receives its own ExecutionContext with partition-specific parameters (e.g. row range, file path).
Implementing a Custom Partitioner
A Partitioner returns a Map<String, ExecutionContext> where each entry represents one partition. The map keys become the partition names visible in the Job Repository.
This example partitions a table by ID range — each worker handles a slice of rows:
import org.springframework.batch.core.partition.support.Partitioner;
import org.springframework.batch.item.ExecutionContext;
import java.util.HashMap;
import java.util.Map;
public class RangePartitioner implements Partitioner {
private final long totalRows;
private final int gridSize;
public RangePartitioner(long totalRows, int gridSize) {
this.totalRows = totalRows;
this.gridSize = gridSize;
}
@Override
public Map<String, ExecutionContext> partition(int gridSize) {
long partitionSize = totalRows / gridSize;
Map<String, ExecutionContext> partitions = new HashMap<>();
for (int i = 0; i < gridSize; i++) {
ExecutionContext ctx = new ExecutionContext();
long minId = i * partitionSize + 1;
long maxId = (i == gridSize - 1) ? totalRows : (i + 1) * partitionSize;
ctx.putLong("minId", minId);
ctx.putLong("maxId", maxId);
ctx.putString("name", "partition" + i);
partitions.put("partition" + i, ctx);
}
return partitions;
}
}Wiring the Partitioned Step
Once you have a Partitioner, wire it into a manager step using StepBuilder.partitioner(). The TaskExecutorPartitionHandler runs each worker step in a thread pool within the same JVM.
The worker step reads using @StepScope beans that pull minId/maxId from the partition's ExecutionContext.
import org.springframework.batch.core.Step;
import org.springframework.batch.core.partition.support.TaskExecutorPartitionHandler;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.task.SimpleAsyncTaskExecutor;
import org.springframework.transaction.PlatformTransactionManager;
@Configuration
public class PartitionedStepConfig {
@Bean
public Step managerStep(JobRepository jobRepository,
Step workerStep,
RangePartitioner partitioner) {
TaskExecutorPartitionHandler handler = new TaskExecutorPartitionHandler();
handler.setStep(workerStep);
handler.setTaskExecutor(new SimpleAsyncTaskExecutor());
handler.setGridSize(8); // 8 parallel workers
return new StepBuilder("managerStep", jobRepository)
.partitioner("workerStep", partitioner)
.partitionHandler(handler)
.build();
}
}Step-Scoped Worker Beans
Worker step beans must be declared with @StepScope so Spring creates a new instance per partition, injecting each partition's ExecutionContext values via @Value("#{stepExecutionContext['minId']}").
This pattern ensures each thread reads a completely independent row range with no shared state:
import org.springframework.batch.core.Step;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.scope.context.StepSynchronizationManager;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.batch.item.database.JdbcPagingItemReader;
import org.springframework.batch.item.database.Order;
import org.springframework.batch.item.database.support.SqlPagingQueryProviderFactoryBean;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Scope;
import org.springframework.context.annotation.ScopedProxyMode;
import javax.sql.DataSource;
import java.util.Map;
@Configuration
public class WorkerStepConfig {
@Bean
@Scope(value = "step", proxyMode = ScopedProxyMode.TARGET_CLASS)
public JdbcPagingItemReader<Order> workerReader(
DataSource dataSource,
@Value("#{stepExecutionContext['minId']}") Long minId,
@Value("#{stepExecutionContext['maxId']}") Long maxId) throws Exception {
JdbcPagingItemReader<Order> reader = new JdbcPagingItemReader<>();
reader.setDataSource(dataSource);
reader.setPageSize(100);
reader.setRowMapper((rs, i) -> new Order(rs.getLong("id"), rs.getString("status")));
reader.setSelectClause("SELECT id, status");
reader.setFromClause("FROM orders");
reader.setWhereClause("WHERE id BETWEEN " + minId + " AND " + maxId);
reader.setSortKeys(Map.of("id", org.springframework.batch.item.database.Order.ASCENDING));
reader.afterPropertiesSet();
return reader;
}
}Remote Partitioning with Spring Integration
Remote Partitioning moves worker steps to separate JVM processes (or Kubernetes pods). The manager sends partition StepExecutionRequest messages over a message broker (RabbitMQ, Kafka, etc.) and workers respond with results.
- Add
spring-batch-integrationdependency - Manager uses
MessageChannelPartitionHandler - Workers listen on an input channel with
StepExecutionRequestHandler
This scales horizontally — add more worker pods to increase throughput without redeploying the manager.
Remote Chunking Architecture
Remote Chunking differs from partitioning: the manager reads all data and sends individual chunks to remote workers for processing and writing. Workers don't access the datasource directly.
- Lower latency per item (manager controls read order)
- Network becomes the bottleneck at high throughput
- Guaranteed delivery requires a durable broker (no data loss on worker crash)
Use remote chunking when processing/writing is the CPU bottleneck, not reading. Use remote partitioning when reading is also slow.
// Remote Chunking manager configuration (spring-batch-integration)
import org.springframework.batch.integration.chunk.RemoteChunkingManagerStepBuilderFactory;
import org.springframework.batch.core.Step;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.channel.QueueChannel;
@Configuration
public class RemoteChunkingConfig {
private final RemoteChunkingManagerStepBuilderFactory managerStepBuilderFactory;
public RemoteChunkingConfig(RemoteChunkingManagerStepBuilderFactory factory) {
this.managerStepBuilderFactory = factory;
}
@Bean
public DirectChannel requests() { return new DirectChannel(); }
@Bean
public QueueChannel replies() { return new QueueChannel(); }
@Bean
public Step remoteChunkingManagerStep() {
return managerStepBuilderFactory
.get("remoteChunkingManager")
.<String, String>chunk(200)
.reader(flatFileReader()) // manager reads
.outputChannel(requests()) // sends chunks to workers
.inputChannel(replies()) // receives ack from workers
.build();
}
private org.springframework.batch.item.ItemReader<String> flatFileReader() {
// returns a configured FlatFileItemReader
return null; // replace with actual reader bean
}
}Choosing the Right Strategy
Picking the wrong strategy leads to either under-utilization or unnecessary complexity. Use this decision guide:
- Multi-threaded step — data fits in one source, reader can be made thread-safe, simplest option
- Parallel steps — independent data loads that don't share state
- Local partitioning — large single-source dataset, single JVM, ID/date ranges easy to define
- Remote partitioning — data is too large for one JVM, horizontal scaling needed, workers can reach the data source
- Remote chunking — processing/writing is the bottleneck, workers are stateless, a message broker is already in your stack
Always prefer local strategies first — they are easier to monitor, debug, and restart after failure.
Knowledge Check: Partitioning vs Remote Chunking
Test your understanding of when to apply each Spring Batch scaling strategy.
Lesson Recap: Partitioning and Parallel Execution
In this lesson you learned how Spring Batch scales throughput beyond single-threaded processing:
- Multi-threaded steps add parallelism with minimal config — use
SimpleAsyncTaskExecutorand wrap stateful readers inSynchronizedItemStreamReader - Parallel steps via
split()flows run independent steps concurrently within one job - Local partitioning divides a dataset into ID/date ranges; a
TaskExecutorPartitionHandlerruns worker steps in a thread pool; worker beans use@StepScope+@Value("#{stepExecutionContext[...]}")" - Remote partitioning distributes workers across JVMs/pods over a message broker — best when reading is the bottleneck and workers can reach the data source
- Remote chunking ships pre-read chunks to stateless remote workers — best when processing/writing is the bottleneck
Always start with the simplest strategy that meets your throughput requirements, and prefer local strategies to keep observability and failure recovery straightforward.
常见问题解答
「分区与并行步骤执行」课时是免费的吗?
是的 — 「分区与并行步骤执行」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Spring Boot 4 Complete Guide 课程的其余内容,请升级到 CoddyKit PRO。 Spring Boot 4 Complete Guide 课程共包含 4 节课。
「分区与并行步骤执行」这节课中我会学到什么?
通过多线程步骤、分区和远程分块策略提升吞吐量。 你通过在浏览器中直接运行的动手代码来练习 Spring Boot 4 Complete Guide,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 Spring Boot 4 Complete Guide 需要有经验吗?
无需任何先前经验。CoddyKit 上的 Spring Boot 4 Complete Guide 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 4 节课,共 4 节。
「分区与并行步骤执行」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 Spring Boot 4 Complete Guide 课中编写并运行代码吗?
能。每节 Spring Boot 4 Complete Guide 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。