← Back to list
SW Engineering & Management
#리액티브스트림#배압#Backpressure#논블로킹#Reactor
Last updated · 2026-09-27

Reactive Streams and Backpressure Control

1. Overview

A. Definition

Reactive Streams is a specification and programming model that processes asynchronous, non-blocking data streams between producers and consumers while standardizing a backpressure signal so that data flows only at a rate the consumer can handle.

The essence of Reactive Streams lies in combining "push" and "pull" into a single protocol. Traditional event-based processing is a push model in which the producer unilaterally pushes data to the consumer, so when the consumer is slow, buffers overflow or messages are lost. Conversely, a pull model in which the consumer requests data whenever it needs it is safe but suffers lower throughput due to increased round-trip latency. Reactive Streams introduces a request(n) signal by which the consumer tells the producer "I can accept up to n items right now," realizing a dynamic push-pull model that preserves the efficiency of push in normal times yet automatically throttles the flow when the consumer becomes saturated.

Here, backpressure refers to a control signal that, when the consumer's processing capacity is lower than the producer's production rate, sends that shortfall back upstream to slow the production rate itself. Just as blocking the downstream of a water pipe raises the upstream pressure, in a software pipeline the downstream delay must be propagated upstream so that the system as a whole does not collapse. A system without backpressure, when load spikes, ends in one of two failures: the queue grows without bound until the process terminates with out-of-memory, or, if the queue is bounded, messages are silently dropped.

B. Background and Necessity

The years 2013–2015, when Reactive Streams was established as a standard, were a period when streaming data and microservices exploded. As libraries such as RxJava, Akka Streams, and Project Reactor each handled asynchronous streams in their own way, the problem of incompatible backpressure signals recurred when combining different libraries. To solve this, Netflix, Lightbend, Pivotal, and others agreed on a minimal specification of four interfaces — Publisher, Subscriber, Subscription, and Processor — which was later incorporated into the standard library as Java 9's java.util.concurrent.Flow.

The necessity is explained along two axes: resource efficiency and stability. First, the thread-per-request model occupies a thread per request, so when concurrent connections reach tens of thousands, thread stack memory and context-switching costs explode. For example, if one thread uses a 1 MB stack, ten thousand threads consume 10 GB, and since most threads are waiting on I/O, the CPU idles while memory is exhausted. The non-blocking reactive model handles tens of thousands of concurrent connections with a small number of event-loop threads (usually the number of CPU cores), greatly saving memory.

Second, in a distributed pipeline, failure is amplified if the delay of one stage does not propagate upstream. If a Kafka consumer slows down due to database write latency yet keeps pulling data from the broker, the consumer's heap overflows. Backpressure makes the consumer poll only as much as it can process in this situation, letting the delay be absorbed across the whole pipeline. Reactive Streams is therefore not merely a performance optimization but a core means of resilience design under load spikes.

C. Characteristics

The characteristics of Reactive Streams are summarized as asynchrony, non-blocking behavior, backpressure-based flow control, and composability. It describes pipelines declaratively by chaining operators in a functional style, and each operator also plays the role of reconciling upstream demand with downstream demand. The table below organizes the differences from traditional models, but the "why" of each item is explained in prose in the sections that follow.

Category Blocking synchronous model Pure push event model Reactive Streams
Flow control Call stack is natural backpressure None (loss/OOM risk) Explicit backpressure via request(n)
Thread usage Thread occupied per request Event loop Event loop + non-blocking
Delay propagation Propagated by synchronous call Not propagated Propagated by demand signal
Composability Low Medium High (operator chains)

2. The Reactive Streams Specification and Components

A. The Four Core Interfaces

The Reactive Streams specification defines the interaction contract with four interfaces: Publisher, Subscriber, Subscription, and Processor. A Publisher is the source that produces data and connects a consumer via a subscribe() call. A Subscriber consumes data and has four callbacks: onSubscribe, onNext, onError, and onComplete. A Subscription represents a single connection between producer and consumer, providing request(n) by which the consumer requests demand and cancel() to sever the connection. A Processor is an intermediate stage that is both a Publisher and a Subscriber, corresponding to each node in an operator chain.

The overall structure can be expressed as a concept diagram as follows.

flowchart LR
    P["Publisher(producer)"] -->|onSubscribe| SUB["Subscription"]
    SUB -->|request n| P
    P -->|onNext data| PR["Processor(operator chain)"]
    PR -->|onNext transform| S["Subscriber(consumer)"]
    S -->|request n demand| PR
    PR -->|request n readjust| P
    S -.->|onError / onComplete| END["termination signal"]

The key to this structure is that while data flows from left to right, the demand signal (request(n)) flows in reverse, from right to left. When the consumer requests only as many items as it can handle, that signal is relayed through the operators to the producer, and the producer emits onNext only up to the requested count. Thanks to this contract, the producer never pushes data beyond the requested amount even without knowing the consumer's processing capacity.

B. Signal Contract and Invariants

The specification requires strict invariants on the order and count of signals. onSubscribe must be called exactly once first, onNext must never exceed the requested count, and no signal may occur after onError or onComplete. Such rules are the minimal contract for safely connecting different libraries, and compliance is verified by a standard test suite called the TCK (Technology Compatibility Kit). Because a non-compliant implementation causes race conditions or resource leaks when composed, the practical recommendation is to use the factory methods of a verified library rather than implementing a Publisher directly.

C. Cold Streams and Hot Streams

Streams are divided into cold and hot depending on the moment of subscription. A cold stream starts data production anew from the beginning each time a subscriber attaches, and suits work that completes per request, such as HTTP requests, file reads, and database queries. A hot stream flows continuously regardless of whether subscribers exist, corresponding to broadcast-like sources such as sensor events, stock quotes, and user clicks. It is difficult to apply backpressure directly to a hot stream because the production source does not wait for the consumer's demand. In this case, one must co-design loss-tolerant strategies such as buffering, sampling, and keeping only the latest value, described later.

3. Backpressure Control Mechanism and Operating Procedure

A. Demand-Based Flow Control Procedure

The standard mechanism of backpressure is a round-trip dialogue in which the consumer calls request(n) for as much as its capacity allows, and the producer emits only within that range. The consumer can choose to request a large demand at once and replenish it as it is drained (e.g., request 256 and add 128 when half is consumed), or to request exactly one at a time to match its processing speed. The sequence diagram below shows how demand is adjusted when the consumer slows down.

sequenceDiagram
    participant S as Subscriber(consumer)
    participant O as Operator(buffer)
    participant P as Publisher(producer)
    S->>O: request(256)
    O->>P: request(256)
    P-->>O: onNext x256
    O-->>S: onNext x256
    Note over S: processing delay occurs
    S->>O: replenish only request(64)
    O->>P: request only as buffer allows
    Note over P: emission rate auto-decreases
    S->>O: cancel() or complete

What is notable in this procedure is that when the consumer is delayed, its replenishment request is delayed, and as a result the producer's emission naturally slows down too. Without any explicit "rate-limit command," the speed of the entire pipeline is matched to the consumer solely through the delay of the demand signal. In practice, Project Reactor's Flux automates this dialogue by keeping a prefetch buffer of size 256 by default and requesting the next batch when 75% is consumed.

B. Collapse Scenario When Backpressure Is Absent

If backpressure is not properly propagated, the system collapses in one of two directions. Using an unbounded buffer, the queue keeps growing on a load spike, the heap is exhausted, and the process terminates with OutOfMemoryError. Using a bounded buffer and throwing on overflow, requests fail with errors such as MissingBackpressureException. In the post-incident analyses of several streaming services in the early 2020s, cases where consumption delay failed to propagate upstream and the broker–consumer queue ballooned were reported repeatedly. Backpressure design must therefore be verified not by "does it run well under normal load" but by "how gracefully it degrades when peak load exceeds consumption capacity."

C. Loss-Tolerant Strategies and Buffering

For hot streams that cannot wait for the consumer indefinitely, an overflow strategy that intentionally drops or summarizes some data is needed. Representatively, buffer accumulates up to a certain size and hands off in batches, drop discards new data when the consumer is busy, and latest keeps only the most recent value and overwrites previous ones. Sample or window groups data by time or count to lower downstream load. The table below compares the characteristics of each strategy.

Strategy Behavior Data loss Suitable situation
buffer Load into a bounded queue, batch Error/block on overflow Brief bursts, loss unacceptable
drop Discard the excess Yes When order matters more than freshness
latest Keep only the latest 1 item Yes Dashboards, quote displays
error Immediate failure signal - Batches where fast failure is better

Strategy selection starts from business requirements. Data that must not lose even a single item, such as payment events, must guarantee no loss with buffer + a durable queue + reprocessing, whereas data for which only the latest value is meaningful, such as a real-time temperature gauge, is better served by latest, which discards stale values for the benefit of both resource efficiency and user experience.

4. Comparison and Practical Application Cases

A. Trade-offs with the Imperative Model

The reactive model gives high throughput and resource efficiency but carries the cost of debugging difficulty and a learning curve. In an asynchronous pipeline, stack traces fail to capture the actual execution path, making it hard to trace the cause of errors, and a single blocking call that stalls the event loop causes throughput to plummet. The thread-per-request model, by contrast, has intuitive code and is easy to debug, and recently Java's virtual threads (Project Loom) have emerged, making it possible to obtain high concurrency even with blocking code. Therefore, what matters is not "reactive unconditionally" but the judgment to apply it selectively to segments that are I/O-bound and inherently need streaming and backpressure control.

B. Web Service Case — Spring WebFlux

Suppose an API gateway must handle tens of thousands of requests per second with a small number of event loops. With Spring MVC (thread-per-request), 200 Tomcat threads are all occupied waiting for external API responses, so even though actual CPU utilization is 20%, requests wait in the queue and time out. Switching to Spring WebFlux (based on Reactor Netty), event loops numbering as many as the CPU cores handle tens of thousands of connections non-blockingly, and when a downstream API slows, request(n) propagates to the upstream client and naturally slows the inflow. In measured cases, on the same hardware the concurrent-connection capacity increases several-fold and P99 latency stabilizes.

C. Data Pipeline Case — Kafka Consumer

Suppose that in a Kafka-based event pipeline the consumer can write 5,000 records per second to the downstream database while 20,000 records per second flow into the topic. Polling without limit and without backpressure piles unprocessed records on the consumer heap, leading to GC storms and OOM. The Kafka connectors of Reactor Kafka or Akka Streams adjust the poll calls and offset-commit rate to the downstream sink's demand and limit per-partition prefetch so the consumer pulls only as much as it can process. As a result, even when momentary inflow exceeds consumption capacity, data is safely retained on the broker while the consumer catches up at a steady rate, and the whole pipeline buffers the inflow surge.

5. Deep Dive — Recent Trends and Related-Technology Linkage

The Reactive Streams standard was incorporated into Java 9's Flow API to become part of the language standard, and backpressure is expanding across all layers of the stack, as with R2DBC (reactive relational database access) and reactive gRPC. If a database driver does not support backpressure, one point of the pipeline blocks and the whole benefit collapses, so "end-to-end non-blocking" has become a core practical concern.

Meanwhile, virtual threads, finalized in Java 21, have risen as an alternative to the reactive model. Virtual threads execute code that looks blocking in a non-blocking way, letting one obtain the reactive model's resource efficiency with imperative code, but it is important that virtual threads themselves do not provide backpressure. That is, virtual threads solve the concurrency problem (thread scarcity), but the flow-control problem (mismatch between production and consumption rates) still needs a separate mechanism such as Reactive Streams or semaphores/bounded queues. Therefore the two technologies are settling into a complementary relationship — "handle simple I/O with virtual threads, and streaming that needs backpressure with Reactive Streams" — rather than competing.

The concept of backpressure is not confined to Reactive Streams. TCP's receive window, gRPC's HTTP/2 flow control, Kafka's consumer-lag-based throttling, and the circuit breaker's load shedding all belong to the broad family of backpressure and flow control. From a professional engineer's perspective, tying these into a single "flow control" lineage lets one compose an answer centered on principles rather than memorizing individual technologies.

6. Considerations and Implications

First, one must judge the scope of application soberly. Reactive Streams gives large benefits in I/O-bound, high-concurrency, streaming segments, but it only adds complexity to CPU-bound batches or simple CRUD. Now that virtual threads are mature, one needs a decision criterion that first distinguishes "is blocking code the problem, or is flow control the problem," and chooses reactive only in the latter case.

Second, one must make end-to-end non-blocking a design principle. If a blocking call (JDBC, synchronous file I/O, etc.) runs on the event loop at even a single point in the pipeline, the small number of event loops stalls and total throughput collapses. Unavoidable blocking work must be isolated on a separate bounded elastic scheduler, and it is desirable to unify the data-access layer on a non-blocking driver such as R2DBC.

Third, one must explicitly select the overflow strategy to match the business risk level. For domains where data loss is unacceptable, combine a lossless buffer with a durable queue and reprocessing; for domains where only freshness matters, save resources with a loss-tolerant strategy — leaving the overflow policy itself as an Architecture Decision Record (ADR) so operators can understand it.

Fourth, one must design observability and testing together. Because backpressure is invisible under normal load and only surfaces at peak load, one must continuously collect metrics such as prefetch size, buffer occupancy, emission relative to request, and consumer lag, and verify in advance "the degradation curve when consumption capacity is exceeded" through load testing and fault injection. Checkpoints and context propagation to complement asynchronous stack traces are also part of operational readiness.

Fifth, one must pursue organizational capability and standardization in parallel. Because reactive code has a steep learning curve and can even worsen performance if misused, it is safer to set the idioms of a verified library as the team standard and to include in code review the rule of composing factories and operators rather than implementing a Publisher directly.

References


In one line: Reactive Streams is an asynchronous, non-blocking stream-processing standard that controls the production rate by sending request(n)-based backpressure signals back upstream to match the consumer's processing capacity, achieving both the resource efficiency of non-blocking behavior and resilience under load spikes at once.