パーティショニングとステップの並列実行
マルチスレッドステップ、パーティショニング、リモートチャンク処理戦略でスループットを向上させます。
「パーティショニングとステップの並列実行」はCoddyKit上の無料Spring Boot 4 Complete Guideレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応の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.
よくある質問
「パーティショニングとステップの並列実行」レッスンは無料ですか?
はい。「パーティショニングとステップの並列実行」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Spring Boot 4 Complete Guideコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Spring Boot 4 Complete Guideコースには全4レッスンが含まれています。
「パーティショニングとステップの並列実行」で何を学びますか?
マルチスレッドステップ、パーティショニング、リモートチャンク処理戦略でスループットを向上させます。 ブラウザで直接実行するハンズオンコードでSpring Boot 4 Complete Guideを演習し、24時間対応の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フィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- ジョブ、ステップ、JobRepositoryモデル
- チャンク指向のReader-Processor-Writerフロー
- フォールトトレランス、スキップ、リトライポリシー
- パーティショニングとステップの並列実行