0Pricing
Spring Boot 4 Complete Guide · 课时

分区与并行步骤执行

通过多线程步骤、分区和远程分块策略提升吞吐量。

分区与并行步骤执行 是 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 — creates ExecutionContext maps, one per partition
  • PartitionHandler — 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-integration dependency
  • 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 SimpleAsyncTaskExecutor and wrap stateful readers in SynchronizedItemStreamReader
  • Parallel steps via split() flows run independent steps concurrently within one job
  • Local partitioning divides a dataset into ID/date ranges; a TaskExecutorPartitionHandler runs 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 反馈 — 无需本地设置。

此课程中的所有课时

  1. 作业、步骤与 JobRepository 模型
  2. 面向块的读取器—处理器—写入器流程
  3. 容错、跳过与重试策略
  4. 分区与并行步骤执行
← 返回 Spring Boot 4 Complete Guide