Spring Boot 4 – komplett guide · Lektion

Partitionering och parallell körning av steg

Öka genomströmningen med flertrådade steg, partitionering och strategier för fjärrbaserad chunk-bearbetning.

Lektion 4 av 413 steg

Partitionering och parallell körning av steg är en gratis lektion i Spring Boot 4 – komplett guide på CoddyKit. Detta är lektion 4 av 4. Du kan läsa vilka 3 lektioner som helst i den här lärvägen kostnadsfritt i sin helhet – därefter låser CoddyKit PRO upp alla lektioner, plus praktisk övning med en inbyggd kodredigerare och en AI-lärare dygnet runt. Den ingår i lärvägen för Spring Boot 4 – komplett guide, och Era framsteg synkroniseras mellan webben och CoddyKit-appen. Kursen i Spring Boot 4 – komplett guide innehåller totalt 4 lektioner.

Varför skala Spring Batch?

Spring Batch-jobb med en tråd bearbetar en chunk i taget — det fungerar bra för små datamängder, men är för långsamt för miljontals poster. När batchgenomströmningen blir en flaskhals erbjuder Spring Batch fyra skalningsstrategier:

  • Steg med flera trådar — parallella trådar inom ett steg i en enda JVM
  • Parallella steg — oberoende steg körs samtidigt i ett flöde
  • Partitionering — dela upp data i partitioner som var och en bearbetas av ett arbetarsteg
  • Remote Chunking — lägg ut chunk-bearbetningen på fjärrarbetare via meddelandemellanvara

Varje strategi innebär olika avvägningar mellan komplexitet och genomströmning. I den här lektionen behandlas alla fyra, med början i den enklaste.

Steg med flera trådar och TaskExecutor

Det enklaste sättet att lägga till parallellism är att injicera en TaskExecutor i Ert Step. Spring Batch kör då chunker samtidigt i en trådpool. Viktigt: ItemReader måste vara trådsäker (till exempel tillståndslös eller använda SynchronizedItemStreamReader).

Konfigurera ett steg med flera trådar så här:

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();
    }
}

Trådsäkra läsare med SynchronizedItemStreamReader

Standardklassen FlatFileItemReader är inte trådsäker eftersom den lagrar internt tillstånd (den aktuella radens position). Genom att omsluta den med SynchronizedItemStreamReader serialiseras anropen till read(), vilket gör den säker för steg med flera trådar utan att parallelliteten i bearbetning eller skrivning påverkas.

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;
    }
}

Parallella steg med delade flöden

När Ni har oberoende steg (till exempel samtidig inläsning av produkter och kunder) använder Ni ett split()-flöde så att båda stegen körs samtidigt. Spring Batchs FlowBuilder har inbyggt stöd för detta.

Steg i en split delar inga data — de måste arbeta med separata resurser för att undvika konkurrens.

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();
    }
}

Introduktion till partitionering

Partitionering delar upp en datamängd i icke-överlappande partitioner, som var och en bearbetas oberoende. Ett hanterarsteg (tidigare kallat master) skapar partitionerna och delegerar dem.

  • Partitioner — skapar ExecutionContext-mappar, en per partition
  • PartitionHandler — avgör hur arbetarstegen startas (lokalt eller på distans)
  • TaskExecutorPartitionHandler — kör arbetarna lokalt i en trådpool

Varje arbetarsteg får sin egen ExecutionContext med partitionsspecifika parametrar (till exempel radintervall eller filsökväg).

Implementera en anpassad Partitioner

En Partitioner returnerar en Map<String, ExecutionContext> där varje post representerar en partition. Map-nycklarna blir de partitionsnamn som visas i jobbförrådet.

Det här exemplet partitionerar en tabell efter ID-intervall — varje arbetare hanterar ett utsnitt av raderna:

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;
    }
}

Koppla in det partitionerade steget

När Ni har en Partitioner kopplar Ni in den i ett hanterarsteg med StepBuilder.partitioner(). TaskExecutorPartitionHandler kör varje arbetar steg i en trådpool i samma JVM.

Arbetarsteget läser med @StepScope-bönor som hämtar minId/maxId från partitionens 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();
    }
}

Stegomfattade arbetar-bönor

Arbetarstegens bönor måste deklareras med @StepScope så att Spring skapar en ny instans per partition och injicerar varje partitions värden från ExecutionContext via @Value("#{stepExecutionContext['minId']}").

Detta mönster säkerställer att varje tråd läser ett helt separat radintervall utan delat tillstånd:

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;
    }
}

Fjärrpartitionering med Spring Integration

Fjärrpartitionering flyttar arbetarsteget till separata JVM-processer (eller Kubernetes-poddar). Hanteraren skickar meddelanden av typen StepExecutionRequest för partitionerna via en meddelandemäklare (RabbitMQ, Kafka med flera), och arbetarna svarar med resultat.

  • Lägg till beroendet spring-batch-integration
  • Hanteraren använder MessageChannelPartitionHandler
  • Arbetarna lyssnar på en indatakanal med StepExecutionRequestHandler

Detta skalar horisontellt — lägg till fler arbetar-poddar för att öka genomströmningen utan att distribuera om hanteraren.

Arkitektur för Remote Chunking

Remote Chunking skiljer sig från partitionering: hanteraren läser alla data och skickar enskilda chunker till fjärrarbetare för bearbetning och skrivning. Arbetarna har ingen direkt åtkomst till datakällan.

  • Lägre fördröjning per objekt (hanteraren styr läsordningen)
  • Nätverket blir flaskhalsen vid hög genomströmning
  • Garanterad leverans kräver en beständig meddelandemäklare (inga dataförluster om en arbetare kraschar)

Använd remote chunking när bearbetning eller skrivning är CPU-flaskhalsen, inte läsningen. Använd fjärrpartitionering när även läsningen är långsam.

// 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
    }
}

Välj rätt strategi

Om Ni väljer fel strategi leder det antingen till outnyttjade resurser eller onödig komplexitet. Använd denna beslutsvägledning:

  • Steg med flera trådar — data ryms i en enda källa, läsaren kan göras trådsäker och detta är det enklaste alternativet
  • Parallella steg — oberoende datainläsningar som inte delar tillstånd
  • Lokal partitionering — stor datamängd från en enda källa, en JVM, och ID- eller datumintervall som är enkla att definiera
  • Fjärrpartitionering — data är för stor för en JVM, horisontell skalning behövs och arbetarna kan nå datakällan
  • Remote Chunking — bearbetning eller skrivning är flaskhalsen, arbetarna är tillståndslösa och en meddelandemäklare finns redan i Er stack

Föredra alltid lokala strategier först — de är enklare att övervaka, felsöka och starta om efter fel.

Kunskapstest: partitionering jämfört med Remote Chunking

Testa Er förståelse av när Ni ska använda de olika skalningsstrategierna i Spring Batch.

Lektionssammanfattning: partitionering och parallell körning

I den här lektionen lärde Ni Er hur Spring Batch skalar genomströmningen utöver bearbetning med en enda tråd:

  • Steg med flera trådar lägger till parallellism med minimal konfiguration — använd SimpleAsyncTaskExecutor och omslut tillståndsbevarande läsare med SynchronizedItemStreamReader
  • Parallella steg via split()-flöden kör oberoende steg samtidigt inom samma jobb
  • Lokal partitionering delar upp en datamängd i ID- eller datumintervall; en TaskExecutorPartitionHandler kör arbetarsteg i en trådpool; arbetar-bönor använder @StepScope + @Value("#{stepExecutionContext[...]}")"
  • Fjärrpartitionering distribuerar arbetare över JVM:er eller poddar via en meddelandemäklare — bäst när läsningen är flaskhalsen och arbetarna kan nå datakällan
  • Remote Chunking skickar förlästa chunker till tillståndslösa fjärrarbetare — bäst när bearbetning eller skrivning är flaskhalsen

Börja alltid med den enklaste strategin som uppfyller Era krav på genomströmning och föredra lokala strategier för att hålla övervakning och felåterställning okomplicerade.

Gratis att börja

Lär dig Java med en AI-lärare – gratis

Skriv och kör riktig kod i webbläsaren, få omedelbar hjälp av en AI-lärare dygnet runt och fortsätt där du slutade – på webben eller i appen.

Kurser
21
Lektioner
84

Vanliga frågor

Är lektionen ”Partitionering och parallell körning av steg” gratis?

Ja – du kan läsa vilka 3 lektioner som helst i lärvägen Spring Boot 4 – komplett guide, inklusive ”Partitionering och parallell körning av steg”, kostnadsfritt i sin helhet här på webben. Därefter låser CoddyKit PRO upp alla lektioner, plus interaktiv övning med en inbyggd kodredigerare och en AI-lärare dygnet runt. Kursen i Spring Boot 4 – komplett guide innehåller totalt 4 lektioner.

Vad lär jag mig i ”Partitionering och parallell körning av steg”?

Öka genomströmningen med flertrådade steg, partitionering och strategier för fjärrbaserad chunk-bearbetning. Ni övar på Spring Boot 4 – komplett guide med praktisk kod som körs direkt i webbläsaren, medan en AI-handledare som är tillgänglig dygnet runt svarar på Era frågor under lektionen.

Behöver jag någon erfarenhet för att börja lära mig Spring Boot 4 – komplett guide?

Du behöver inga förkunskaper. Utbildningen i Spring Boot 4 – komplett guide på CoddyKit är upplagd för allt från nybörjare till avancerade elever, så att du kan börja här eller från början och gå fram i din egen takt. Detta är lektion 4 av 4.

Hur lång tid tar lektionen ”Partitionering och parallell körning av steg”?

De flesta CoddyKit-lektioner tar cirka 5–10 minuter. Varje lektion är kort och interaktiv, så att du gör stadiga framsteg och kan fortsätta precis där du slutade – på webben eller i appen.

Kan jag skriva och köra kod i den här Spring Boot 4 – komplett guide-lektionen?

Ja. Varje Spring Boot 4 – komplett guide-lektion innehåller en inbyggd kodredigerare, så att du kan skriva och köra riktig kod direkt i webbläsaren och få omedelbar AI-feedback – utan lokal installation.

Alla lektioner i den här kursen

  1. Jobb, steg och JobRepository-modellen
  2. Chunk-baserade Reader-Processor-Writer-flöden
  3. Feltolerans samt skip- och retry-policyer
  4. Partitionering och parallell körning av steg
← Tillbaka till Spring Boot 4 – komplett guide