skip to content

How does the remote-chunking master track completion and fault tolerance across asynchronous workers, and what are the operational limits?

level: principalimportance: should knowfreq 22%

answer

  1. Reader null is necessary, not sufficient
  2. Sent-count == received-response-count to finish
  3. Failure flag in ChunkResponse fails step
  4. throttleLimit = in-flight cap / back-pressure
  5. Skip/retry runs on worker, not master

basics

~20 s

The master's ChunkMessageChannelItemWriter counts chunks sent versus ChunkResponses received; the step only completes when all outstanding responses arrive. A failure flag in a response fails the step. The main limits: the single reader, throttle/timeout on outstanding chunks, and worker-side skip/retry not being visible to the master.

solid answer

~50 s

Because workers reply asynchronously, the master can't finish just because the reader is exhausted — it must reconcile acknowledgements. `ChunkMessageChannelItemWriter` keeps a running count of dispatched chunks and consumed `ChunkResponse`s and blocks step completion until they balance; a response carrying a failure/exception marks the step failed. Two operational controls matter: a **throttle limit** on how many chunks may be in flight (bounding memory and queue depth / providing back-pressure), and a **max-wait / timeout** for responses so a dead worker doesn't hang the job forever. Key limits to raise in design review: (1) the **reader is single-threaded** — the ceiling on throughput; (2) **skip/retry semantics** normally applied inside a chunk step now happen on the worker, so the master's view of errors is coarser and restart/idempotency must be designed deliberately; (3) requests can outrun workers, so throttling and durable queues are essential; (4) the master is a single coordinator whose failure needs restart handling.

code

java · 18 lines
java
// Master step: bound in-flight chunks and cap response wait so a
// dead worker can't hang the job. (Builder methods surface these controls.)
@Bean
public TaskletStep managerStep(RemoteChunkingManagerStepBuilderFactory managers,
                               ItemReader<Order> reader,
                               MessageChannel requests,
                               PollableChannel replies) {
    return managers.get("managerStep")
            .<Order, Order>chunk(100)
            .reader(reader)
            .outputChannel(requests)
            .inputChannel(replies)
            .throttleLimit(20)      // max unacknowledged chunks in flight (back-pressure)
            .maxWaitTimeouts(3)     // give up waiting after N empty polls of replies
            .build();
}
// Note: the writer only lets the step complete once sent-chunk count
// == received-ChunkResponse count and none reported failure.

go deeper

for a junior

Not expected to answer in depth.

for a middle

Should know the master waits for worker acknowledgements before finishing.

for a senior

Should explain sent-vs-received reconciliation, failure via ChunkResponse, and throttle/timeout controls.

for a principal

Should reason about back-pressure sizing, restart+idempotency interplay, the single-reader ceiling, SPOF, and observability across master/worker.

This is the systems-design question: how does a fundamentally asynchronous, distributed step reach a deterministic conclusion, and where are its edges? **Completion tracking:** - A normal step ends when the reader returns `null`. In remote chunking that's necessary but **not sufficient** — the master has sent chunks whose fates are still unknown. - `ChunkMessageChannelItemWriter` maintains local state: number of chunks **sent** and number of `ChunkResponse`s **received/expected**. Each cycle it drains the reply channel, applies each response's `StepContribution` to the master `StepExecution` (keeping write/skip counts accurate), and decrements the outstanding count. - The step is allowed to complete only when **all expected responses have been received** and none signalled failure. So even after the reader is exhausted, the writer keeps polling the reply channel until the ledger balances. - A `ChunkResponse` with `successful=false` (or carrying a throwable) causes the master to fail the step. **Fault-tolerance controls & realities:** - **Throttle limit / max in-flight chunks:** bounds how many unacknowledged chunks may exist. This provides **back-pressure** (master doesn't flood the broker), caps memory, and limits how much work is 'at risk' in flight. Too low starves workers; too high risks large redelivery/reprocessing on failure. - **Response timeout / max wait:** without a bound, a crashed worker (its response never arrives) would hang the step forever. A timeout lets the master fail (or, with durable queues + another worker, the redelivered chunk can still be processed). - **Skip/retry:** the classic chunk-step fault tolerance (`faultTolerant()`, skip/retry policies) executes **on the worker** where process+write run. The master doesn't directly observe per-item skips except via what the `ChunkResponse`/`StepContribution` reports. This means error accounting is coarser and you must decide restart semantics carefully. - **Restart & idempotency:** on job restart the master re-reads from where it left off (reader state in the `ExecutionContext`), but in-flight or unacknowledged chunks may be reprocessed. Combined with at-least-once delivery, **idempotent writes are effectively mandatory**. **Operational limits to name in review:** 1. **Single reader = throughput ceiling.** If the reader can't keep the workers fed, adding workers yields nothing. If reading is the bottleneck at all, this is the wrong pattern (use partitioning). 2. **Master is a single coordinator / SPOF for the step.** Its failure requires job restart; design for that. 3. **Serialization + broker capacity:** items on the wire cost CPU and queue space; large chunks/records stress the broker. 4. **Ordering/correlation:** relies on sequence counters and a well-behaved durable transport; a flaky broker undermines completion accounting. 5. **Observability:** worker-side failures surface as responses, not master stack traces — plan logging/metrics on both sides. **When it shines vs. when to avoid:** ideal when a cheap sequential read feeds expensive, parallelizable process/write and you have a durable broker. Avoid when the read is the bottleneck, when you can't guarantee delivery, or when per-item fine-grained fault tolerance visibility on the master is a hard requirement (partitioning's self-contained worker steps may model that better).

  • Why can't the master complete the step the moment its ItemReader returns null?
    Because chunks already dispatched are still being processed asynchronously by workers. The master must keep draining the reply channel until every sent chunk has a matching ChunkResponse; otherwise it might finish while work is unacknowledged or has failed.
  • What does the throttle limit protect against, and what's the downside of setting it too high?
    It caps how many chunks are unacknowledged/in-flight, giving back-pressure so the master doesn't flood the broker and bounding memory and at-risk work. Set too high, a failure or restart means many chunks may be reprocessed and the queue can bloat; too low starves workers and hurts throughput.
  • Where does chunk-level skip/retry fault tolerance actually run, and why does that matter?
    On the worker, where process and write execute. The master only learns outcomes via ChunkResponse/StepContribution, so its error view is coarser. That shifts responsibility for idempotency and restart correctness onto your worker writers and delivery semantics.

saying these in an interview costs you the question

  • Claiming the step ends as soon as the reader is exhausted
  • Thinking the master applies skip/retry to worker-side processing
  • Ignoring the need for a response timeout (dead worker hangs the job)
  • Believing more workers always increases throughput despite the single reader
  • Assuming no back-pressure is needed

context