Spring Boot 4 Complete Guide · Урок

Разбиение и параллельное выполнение шагов

Повышайте пропускную способность с помощью многопоточных шагов, разбиения и стратегий удалённой обработки блоков.

Урок 4 из 413 шагов

«Разбиение и параллельное выполнение шагов» — бесплатный урок Spring Boot 4 Complete Guide на CoddyKit. Это урок 4 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения 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.

Можно начать бесплатно

Изучай Java с ИИ-репетитором — бесплатно

Пиши и запускай код прямо в браузере, получай мгновенную помощь от ИИ-репетитора 24/7 и продолжи учиться на сайте или в приложении.

Курсы
21
Уроки
84

Часто задаваемые вопросы

Урок «Разбиение и параллельное выполнение шагов» бесплатный?

Да — полный текст урока «Разбиение и параллельное выполнение шагов» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс Spring Boot 4 Complete Guide, подпишись на CoddyKit PRO. Курс Spring Boot 4 Complete Guide содержит 4 уроков всего.

Чему я научусь в уроке «Разбиение и параллельное выполнение шагов»?

Повышайте пропускную способность с помощью многопоточных шагов, разбиения и стратегий удалённой обработки блоков. Ты практикуешь Spring Boot 4 Complete Guide с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать Spring Boot 4 Complete Guide?

Предыдущий опыт не требуется. Spring Boot 4 Complete Guide на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 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