skip to content

How are StepExecution counts aggregated safely across concurrency models — multi-threaded steps versus partitioned steps?

level: principalimportance: nice to knowfreq 20%

answer

  1. multi-threaded: one StepExecution, synchronized apply()
  2. partitioned: N worker StepExecutions, independent counts
  3. StepExecutionAggregator sums worker counts into manager
  4. reader must be thread-safe in multi-threaded step
  5. partition totals roll up at the end, not live

basics

~20 s

In a multi-threaded step, many chunks share one StepExecution and merge counts via the synchronized apply(StepContribution). In partitioning, each partition is its own child StepExecution with independent counts; the manager step sums them for reporting.

solid answer

~40 s

Two different concurrency models aggregate counts differently. A multi-threaded step (a TaskExecutor on a single step) runs many chunks concurrently against ONE StepExecution; each chunk stages counts in its own StepContribution and merges via StepExecution.apply(), which is synchronized so the shared counters don't race. In partitioning, the manager creates N worker StepExecutions (one per partition, each persisted separately), so their readCount/writeCount/etc. accumulate independently with no shared-counter contention; there is no single 'live' aggregate during the run. Tooling and the PartitionStep aggregate the worker counts (e.g. via StepExecutionAggregator, which sums counts and status/exitStatus into the manager StepExecution) for reporting. The principal trade-off: multi-threaded steps need thread-safe readers and share metadata contention, while partitioning gives isolation and horizontal scaling (even remote workers) at the cost of coordinating and rolling up separate StepExecutions.

code

java · 17 lines
java
// Multi-threaded step: one StepExecution, apply() synchronized
Step mt = new StepBuilder("mt", jobRepository)
    .<Row, Row>chunk(100, txManager)
    .reader(new SynchronizedItemStreamReader<>(delegateReader)) // thread-safe!
    .writer(writer)
    .taskExecutor(new SimpleAsyncTaskExecutor())  // concurrent chunks
    .build();

// Partitioned step: N worker StepExecutions, counts summed by aggregator
Step manager = new StepBuilder("manager", jobRepository)
    .partitioner("worker", new ColumnRangePartitioner())
    .step(workerStep)                 // each partition -> its own StepExecution
    .gridSize(8)
    .taskExecutor(new SimpleAsyncTaskExecutor())
    .build();
// DefaultStepExecutionAggregator sums worker read/write/skip/commit/rollback
// counts (and folds status/exitStatus) into the manager StepExecution.

go deeper

for a junior

Know Spring Batch can scale steps with threads or partitions.

for a middle

Know multi-threaded steps share one StepExecution and partitions each have their own.

for a senior

Explain synchronized apply() for multi-threaded and independent worker StepExecutions for partitioning.

for a principal

Weigh contention vs isolation, restart granularity, observability, and how StepExecutionAggregator rolls partition counts into the manager.

### The core issue StepExecution counters are mutable shared state. How they're kept correct depends on **which** Spring Batch scaling model you use. ### 1) Multi-threaded step (single step, TaskExecutor) You attach a `TaskExecutor` to a chunk step so multiple chunks run concurrently. There is **one** `StepExecution`. Each chunk gets its **own** `StepContribution` (thread-local staging — hot per-item increments are lock-free), and on commit calls `stepExecution.apply(contribution)`, which is **`synchronized`** on the StepExecution. That serialization is exactly what keeps `readCount/writeCount/filterCount/skip*` correct under concurrency. Caveat: the `ItemReader` must be **thread-safe** (e.g. wrapped in `SynchronizedItemStreamReader`), and ordering/restart is harder. All counts land in that single StepExecution. ### 2) Partitioned step (manager + workers) With `PartitionStep` / `PartitionHandler`, a **manager** step splits work into partitions via a `Partitioner` and launches N **worker** StepExecutions — **each a distinct, separately persisted `StepExecution`** (child rows in `BATCH_STEP_EXECUTION`, tied to the manager). Each worker counts **independently**; there is **no shared live counter**, so no cross-thread contention on counters at all. Workers can even run on **remote JVMs** (remote partitioning), where a shared in-memory counter would be impossible anyway. #### Roll-up After workers finish, their counts are **aggregated** into the manager's StepExecution for reporting. Spring Batch provides `StepExecutionAggregator` (default impl `DefaultStepExecutionAggregator`) which **sums** the workers' read/write/filter/commit/rollback/skip counts and combines their `BatchStatus`/`ExitStatus` (manager status = the max/most-severe of the workers). So the manager StepExecution shows the totals, while each worker retains its own breakdown. ### 3) Remote chunking (contrast) Here the manager reads and sends chunks to remote workers that process/write. Write counts come back and are applied to the manager's StepExecution — closer to the single-StepExecution model but distributed. (Mentioned for completeness; the aggregation nuance is mainly multi-threaded vs partitioned.) ### Principal-level trade-offs - **Contention vs isolation:** multi-threaded steps share one StepExecution (synchronized merges, thread-safe reader required); partitioning isolates state per worker (no counter contention, natural for large datasets and remote scale-out). - **Restart granularity:** partitioning restarts only the failed partitions (each worker StepExecution has its own status/ExecutionContext); a multi-threaded step restarts as a whole and is trickier to make restartable. - **Observability:** partition workers give per-shard counts (useful to spot a skewed/bad partition); the multi-threaded step gives only the blended total. - **Correctness pitfall:** if you wrote a custom aggregator or read counts mid-run in a partitioned job, remember the manager total isn't 'live' — it's meaningful only after workers complete and aggregation runs. ### Gotchas - Don't assume a partitioned manager's counts update continuously during the run; they're rolled up at the end. - A multi-threaded step with a non-thread-safe reader will corrupt both data and counts — the synchronized apply() protects the counters, not your reader. - commitCount/rollbackCount still bypass StepContribution and are incremented on the (respective) StepExecution directly in all models.

  • In a partitioned job, is there a single StepExecution whose counts update live as workers run?
    No. Each partition worker has its own persisted StepExecution counting independently. The manager's aggregate totals are produced by a StepExecutionAggregator after the workers finish, so they aren't a live running total during execution.
  • What breaks if you use a multi-threaded step with a stateful, non-thread-safe reader?
    The reader's shared cursor/state gets corrupted across threads, causing skipped/duplicated items and wrong read/write counts. The synchronized apply() protects the aggregate counters, not the reader — wrap it (e.g. SynchronizedItemStreamReader) or use partitioning instead.

saying these in an interview costs you the question

  • Claiming a partitioned manager StepExecution has one shared live counter across workers
  • Assuming multi-threaded steps work with any reader without thread-safety
  • Not knowing partition worker counts are aggregated (summed) into the manager

context