skip to content

How does the Partitioner interface work, and how does a worker step read its assigned slice?

level: middleimportance: must knowfreq 62%

answer

  1. partition(gridSize) -> Map<name, ExecutionContext>
  2. put minId/maxId into each ExecutionContext
  3. gridSize is a hint, not a mandate
  4. @StepScope reader + #{stepExecutionContext[...]}
  5. query MIN/MAX first to balance

basics

~10 s

Partitioner has one method: partition(int gridSize) returning Map<String, ExecutionContext>. Each entry is a named partition whose ExecutionContext holds that slice's parameters (e.g. minId/maxId). The worker's @StepScope reader reads them via SpEL like #{stepExecutionContext['minId']}.

solid answer

~40 s

You implement `Partitioner.partition(int gridSize)` to return a `Map<String, ExecutionContext>`: the key is a unique partition name, the value is that partition's own `ExecutionContext` pre-loaded with the parameters that identify its slice (min/max id, a date bucket, a filename). `gridSize` is a hint for how many partitions to make — you can use it or compute your own count. Spring Batch's `StepExecutionSplitter` turns each map entry into a separate worker `StepExecution`, seeding its step ExecutionContext with those values. The worker step's reader must be `@StepScope` (late-bound), so each partition gets a fresh instance; it pulls its bounds via SpEL, e.g. `@Value("#{stepExecutionContext['minId']}") Long minId`. Common gotchas: the reader must be `@StepScope` or all partitions read the same data; and you must compute real boundaries (query MIN/MAX first) so slices are balanced and cover the whole set.

code

java · 44 lines
java
public class ColumnRangePartitioner implements Partitioner {

    private final long min;
    private final long max;

    public ColumnRangePartitioner(long min, long max) {
        this.min = min;   // e.g. from SELECT MIN(id)
        this.max = max;   // e.g. from SELECT MAX(id)
    }

    @Override
    public Map<String, ExecutionContext> partition(int gridSize) {
        long targetSize = (max - min) / gridSize + 1;
        Map<String, ExecutionContext> result = new HashMap<>();
        long start = min;
        int i = 0;
        while (start <= max) {
            long end = Math.min(start + targetSize - 1, max);
            ExecutionContext ctx = new ExecutionContext();
            ctx.putLong("minId", start);
            ctx.putLong("maxId", end);
            result.put("partition" + i, ctx); // unique, stable name
            start += targetSize;
            i++;
        }
        return result;
    }
}

// Worker reader is @StepScope so each partition gets its own bounds:
@Bean
@StepScope
public JdbcCursorItemReader<Record> workerReader(
        @Value("#{stepExecutionContext['minId']}") Long minId,
        @Value("#{stepExecutionContext['maxId']}") Long maxId,
        DataSource dataSource) {
    return new JdbcCursorItemReaderBuilder<Record>()
            .name("workerReader")
            .dataSource(dataSource)
            .sql("SELECT id, payload FROM records WHERE id BETWEEN ? AND ?")
            .preparedStatementSetter(ps -> { ps.setLong(1, minId); ps.setLong(2, maxId); })
            .rowMapper(new RecordRowMapper())
            .build();
}

go deeper

for a junior

Know that Partitioner returns a map of named ExecutionContexts, one per slice.

for a middle

Explain gridSize-as-hint, seeding ExecutionContext, and @StepScope SpEL late binding in the reader.

for a senior

Discuss StepExecutionSplitter, balancing/skew, and deterministic partition names for restart.

for a principal

Design the slice key strategy (range vs modulo vs file), serialization limits, and restart determinism at scale.

## The `Partitioner` interface ```java public interface Partitioner { Map<String, ExecutionContext> partition(int gridSize); } ``` That's the whole contract. Your implementation decides the *division strategy*. - **Key** = a unique **partition name** (e.g. `"partition0"`). It's just an identifier used in metadata/logging. - **Value** = an **`ExecutionContext`** — a serializable key/value bag. You put into it whatever the worker needs to know to process *only its slice*: numeric bounds (`minId`, `maxId`), a `fromDate`/`toDate`, a `fileName`, a `tenantId`, etc. ## `gridSize` The `partition(int gridSize)` argument is a **hint** for how many partitions the framework expects. It typically comes from the PartitionHandler (`setGridSize(n)`). Your Partitioner is free to honor it or compute its own count (e.g. one partition per input file regardless of gridSize). The number of map entries you return is the number of worker StepExecutions that will actually be created. ## From map to worker StepExecutions When the master step runs, a **`StepExecutionSplitter`** (default `SimpleStepExecutionSplitter`) calls your `partition(...)`, then for each `(name, executionContext)` entry creates one **worker `StepExecution`** and copies that ExecutionContext into the worker's **step ExecutionContext**. So the values you set in the Partitioner become the worker's step-scoped data. ## Reading the slice in the worker — why `@StepScope` matters The worker step is defined once but run many times concurrently, each needing *different* bounds. To inject per-partition values you use **late binding** via SpEL against `stepExecutionContext`, and the bean must be **`@StepScope`** so a new instance is created per StepExecution: ```java @Bean @StepScope public JdbcPagingItemReader<Record> reader( @Value("#{stepExecutionContext['minId']}") Long minId, @Value("#{stepExecutionContext['maxId']}") Long maxId, DataSource ds) { ... } ``` If the reader is **not** `@StepScope`, the SpEL `stepExecutionContext` isn't available at instantiation and every partition would read the same (or wrong) data — a classic bug. ## Built-in Partitioner Spring provides **`MultiResourcePartitioner`**: given an array of `Resource`s (e.g. `file*.csv`), it creates one partition per file, putting the file under key `fileName` (as a URL) in each ExecutionContext. For row-range partitioning you write your own. ## Balancing the slices Good practice for range partitioning: first query `SELECT MIN(id), MAX(id)` (or count rows), then divide that span into `gridSize` roughly equal buckets. Uneven slices cause **data skew** — the slowest partition determines total wall-clock time. ## Edge cases & gotchas - **Empty partition**: if a slice matches no rows, that worker just completes with 0 items — harmless. - **Overlapping / gapped bounds**: off-by-one in min/max can double-process or skip rows; use inclusive/exclusive bounds consistently. - **Non-serializable values**: ExecutionContext values are persisted to the metadata store; keep them simple (primitives/strings). - **Partitioner runs once, up front**: it can't see mid-run progress; boundaries are fixed at split time. - **Determinism for restart**: the same partition names must be regenerated on restart so completed partitions are recognized and skipped.

  • Why must the worker's reader be @StepScope?
    Because the partition-specific values live in the step ExecutionContext, which only exists once a StepExecution starts. @StepScope defers bean creation to that point and creates a separate instance per partition, so each reader can late-bind its own #{stepExecutionContext[...]} values. A singleton reader would be built before any StepExecution and share one set of bounds.
  • Is gridSize the guaranteed number of partitions?
    No. gridSize is only a hint passed to partition(). The actual number of partitions equals the number of entries your Partitioner returns in the map, which you control.

saying these in an interview costs you the question

  • Saying the reader can be a plain singleton bean
  • Claiming gridSize forces the exact partition count
  • Putting slice params in the job ExecutionContext instead of per-partition ones
  • Ignoring MIN/MAX so slices are unbalanced or miss rows

context