Flow separates publishers, subscribers and subscriptions so a subscriber can request a number of items rather than accepting an unspecified push rate.
Java Flow: request items and observe publisher completion
This complete program targets Java 11. Its displayed output is checked by the tutorial validation script.
Completion belongs to the subscription
The receipt subscriber requests one amount at subscription time and requests another after each received amount. The publisher sends ten and twenty, then closes. The completion future is joined before the total is read, making successful completion distinct from merely calling submit twice.
SubmissionPublisher owns buffering and delivery behavior. This fixture chooses a direct executor for deterministic local observation; it does not represent a worker pool or network transport. A one-at-a-time request policy here does not prove that a remote importer retains only one record.
Keep errors visible
onError completes the same future exceptionally so join cannot return a successful total after a delivery failure. Production callbacks should define cancellation and downstream failure handling without throwing errors into unrelated application work. Do not treat normal close as a way to repair a failed record.
The atomic total avoids relying on an undocumented thread choice in the subscriber. A real payment amount needs a unit and overflow policy too. Compare this API-level contract with Reactor demand before assuming they provide an identical set of operators or runtime services.
Working program
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
public class ReceiptFlowTotal {
public static void main(String[] args) {
AtomicInteger total=new AtomicInteger();
CompletableFuture<Void> finished=new CompletableFuture<>();
try(SubmissionPublisher<Integer> publisher=new SubmissionPublisher<>(Runnable::run,8)) {
publisher.subscribe(new Flow.Subscriber<Integer>() {
Flow.Subscription subscription;
public void onSubscribe(Flow.Subscription value){subscription=value;subscription.request(1);}
public void onNext(Integer amount){total.addAndGet(amount);subscription.request(1);}
public void onError(Throwable failure){finished.completeExceptionally(failure);}
public void onComplete(){finished.complete(null);}
});
publisher.submit(10);publisher.submit(20);
}
finished.join();System.out.println(total.get());
}
}Output
30Costs and boundaries
The fixture processes two integers with a configured buffer capacity and direct delivery. Application storage depends on retained subscriptions, buffered items and callbacks. A blocking submit or slow subscriber can affect caller latency; no throughput measurement is provided.
Common Mistakes
- Request demand through the subscription.
- Read results after observing completion or failure.
- Do not infer a network buffer limit from this local sequence.
Read next
Java CompletableFuture: composition, failures and executor ownership, Java blocking queues: bounded capacity and backpressure, Spring reactive foundations: request values and cancel a subscription.
