skip to content

What do buffer and window do, and how does transform differ from transformDeferred? Give a use case for each.

level: seniorimportance: should knowfreq 40%

answer

  1. buffer -> Flux<List<T>> (materialized batches)
  2. window -> Flux<Flux<T>> (composable sub-streams)
  3. bufferTimeout = size OR time
  4. transform = once at assembly time
  5. transformDeferred (was compose) = per subscriber at subscribe time

basics

~20 s

buffer groups incoming elements into Lists (by size or time) and emits Flux<List<T>>. window is similar but emits Flux<Flux<T>> — each group is itself a stream. transform lets you factor a reusable chain of operators into a function applied once at assembly time; transformDeferred applies it per subscriber.

solid answer

~40 s

buffer collects elements into batches and emits them as a List — buffer(n) by count, buffer(Duration) by time, or bufferTimeout(n, duration) whichever first — turning a Flux<T> into a Flux<List<T>>. It's for batching, e.g. bulk-inserting DB rows N at a time. window does the same partitioning but emits a Flux<Flux<T>>: each batch is a live inner Flux you can compose further, useful for streaming aggregation without materializing whole lists. transform applies a function that takes and returns a Publisher, letting you package a reusable operator chain and insert it with .transform(myChain); it runs ONCE at assembly time. transformDeferred (formerly composeaka) runs that function per subscriber at subscription time, so it can produce subscriber-specific behavior or read state captured at subscribe time.

code

java · 19 lines
java
// buffer: batch inserts, flush at 500 rows or every second
Flux<Reading> readings = sensor.stream();
readings.bufferTimeout(500, Duration.ofSeconds(1))
        .flatMap(batch -> repository.saveAll(batch))
        .subscribe();

// window: streaming aggregation per 10-second window (no full List in memory)
readings.window(Duration.ofSeconds(10))
        .flatMap(win -> win.map(Reading::value).reduce(0.0, Double::sum));

// transform: reusable cross-cutting chain applied once at assembly
Function<Flux<String>, Flux<String>> logged =
        f -> f.doOnNext(v -> log.info("item {}", v)).timeout(Duration.ofSeconds(2));
Flux.just("a", "b").transform(logged);

// transformDeferred: chain rebuilt for each subscriber at subscription time
AtomicBoolean upper = new AtomicBoolean();
Flux<String> f = Flux.just("a")
        .transformDeferred(p -> upper.get() ? p.map(String::toUpperCase) : p);

go deeper

for a junior

Probably only knows buffer batches into Lists; window/transform are beyond scope.

for a middle

Should explain buffer vs window (List vs sub-Flux) and give a batching use case.

for a senior

Core target: bufferTimeout semantics, window as composable streams, and transform vs transformDeferred = assembly vs subscription time.

for a principal

Reasons about memory/back-pressure of window inners, factoring cross-cutting chains via transform for maintainability, and deferring behavior to subscription time deliberately.

## buffer — batch into Lists `buffer` accumulates upstream elements and emits them grouped as a `List`, converting `Flux<T>` → `Flux<List<T>>`. - **By size:** `buffer(int maxSize)` emits a list every N elements. - **By time:** `buffer(Duration)` emits whatever accumulated each interval. - **Hybrid:** `bufferTimeout(int, Duration)` flushes on size OR time, whichever comes first — the practical choice for batching to avoid unbounded latency on slow streams. - **Use:** bulk operations — accumulate 500 records then do one batch insert; or coalesce bursty events. - **Gotcha:** on completion, a partial (non-full) buffer is still emitted; watch memory if the batch size/time is large and elements are big. ```java sensorReadings .bufferTimeout(500, Duration.ofSeconds(1)) .flatMap(batch -> repository.saveAll(batch)); // one write per batch ``` ## window — batch into sub-Fluxes `window` partitions the same way but emits `Flux<Flux<T>>`: each group is itself a **Flux** you can apply operators to (reduce, count, etc.) **without collecting the whole group into memory** first. Variants: `window(n)`, `window(Duration)`, `windowTimeout(...)`. - **buffer vs window:** buffer gives you the finished `List` (materialized); window gives you a live stream per group (composable, streaming). - **Use:** per-window aggregation, e.g. compute a rolling average per time window while streaming. ```java ticks.window(Duration.ofSeconds(10)) .flatMap(win -> win.reduce(0, Integer::sum)); // sum per 10s window ``` ## transform — reusable assembly-time chain `transform(Function<Publisher, Publisher>)` lets you extract a chain of operators into a named function and splice it into a pipeline: `flux.transform(this::withRetryAndLog)`. The function is evaluated **once, at assembly time** (when the pipeline is built), so the same operator instances are shared by all subscribers. - **Use:** DRY — apply a common cross-cutting chain (logging, retry, timeout) across many pipelines. ## transformDeferred — per-subscriber chain `transformDeferred` (renamed from the old `compose`) runs the function **once per subscriber, at subscription time**. So each `subscribe()` can get a freshly-built chain, and any state/variable the function reads is captured at subscribe time, not assembly time. - **Use:** when behavior must vary per subscription — e.g. flip a feature toggle or counter that's read each time someone subscribes. ```java AtomicInteger calls = new AtomicInteger(); Flux<String> f = Flux.just("x").transformDeferred(p -> calls.incrementAndGet() % 2 == 0 ? p.map(String::toUpperCase) : p); // each subscribe() re-evaluates the branch ``` ## The assembly-vs-subscription connection `transform` = assembly time (built once). `transformDeferred` = subscription time (built per subscriber). This mirrors the broader Reactor distinction: building the pipeline (assembly) is separate from running it (subscription), and deferred operators let you push work to subscription time. ## Gotchas - `transform` capturing a mutable variable shares ONE value across all subscribers because it's evaluated once; use `transformDeferred` if you need per-subscribe evaluation. - window's inner Fluxes must be consumed; ignoring/not subscribing to them can stall or leak because they carry back-pressure. - buffer/window with time need a scheduler; be aware they can emit empty buffers/windows on quiet intervals (depending on variant).

  • When would you choose window over buffer?
    When you want to process each group as a live stream — applying operators like reduce/scan/count per group — without materializing the whole group into a List in memory. buffer is simpler when you actually need the List (e.g., a batch save call that takes a Collection).
  • You capture a counter inside a transform function and expect it to change per subscription, but it doesn't. Why?
    transform evaluates the function once at assembly time, so the counter is read a single time and shared by all subscribers. Use transformDeferred, which re-invokes the function per subscriber at subscription time.

saying these in an interview costs you the question

  • Thinking window emits Lists like buffer
  • Believing transform runs per subscriber
  • Not knowing bufferTimeout flushes on size OR time
  • Forgetting partial buffers flush on completion

context