Posts Tagged ‘Backpressure’
[DevoxxPL2019] Constructing Custom Reactive Publishers: Insights into Project Reactor Internals
Lecturer
Oleh Dokuka contributes to Project Reactor as a committer, authoring books on reactive programming with Spring and serving as a software engineer at Superhuman. Based in Kyiv, he actively participates in conferences and communities focused on asynchronous systems.
Abstract
This inquiry delves into the intricacies of building a reactive publisher compliant with Reactive Streams specifications, drawing from Project Reactor’s design. It covers the rationale behind the spec, naive implementations, concurrency patterns like work-in-progress, and verification via TCK. Through iterative coding, it analyzes challenges in non-blocking data flows, backpressure, and thread safety, pondering effects on debugging, customization, and library extension.
Demystifying Reactive Streams: Specification and Purpose
Reactive Streams standardize asynchronous, non-blocking data processing with backpressure, addressing overflow in producer-consumer scenarios. Oleh commences by recalling the spec’s origins, crafted to unify libraries like RxJava and Akka Streams, ensuring interoperability.
Core interfaces—Publisher, Subscriber, Subscription, Processor—define interactions: publishers emit items, subscribers consume, subscriptions mediate requests and cancellations. The spec mandates rules for thread safety and signal ordering, preventing races.
Contextually, adoption surged with Java 9’s Flow API, embedding reactivity natively. Analytically, backpressure—subscribers requesting items—prevents buffering overloads, crucial in unbounded sources like networks.
Implications: enables composable, resilient pipelines, but demands adherence to 50+ rules, tested via TCK. For developers, understanding facilitates debugging; for extenders, it unlocks optimizations.
Naive Publisher Construction: Initial Steps and Pitfalls
Commencing with a basic array publisher, Oleh demonstrates emitting elements on subscription. Yet, naivety ignores concurrency: parallel subscriptions risk duplicates or misses.
Methodologically, extend TCK’s PublisherVerification for rule checks. Initial failures highlight needs for atomic operations and request tracking.
A subscription class manages emissions:
class ArraySubscription<T> implements Subscription {
private final Subscriber<? super T> subscriber;
private final T[] array;
private int index = 0;
private boolean canceled = false;
public ArraySubscription(Subscriber<? super T> subscriber, T[] array) {
this.subscriber = subscriber;
this.array = array;
}
@Override
public void request(long n) {
if (n <= 0 && !canceled) {
subscriber.onError(new IllegalArgumentException("Non-positive request"));
canceled = true;
return;
}
for (long i = 0; i < n && !canceled; i++) {
if (index < array.length) {
subscriber.onNext(array[index++]);
} else {
subscriber.onComplete();
canceled = true;
break;
}
}
}
@Override
public void cancel() {
canceled = true;
}
}
This handles basics but falters under concurrency, necessitating refinements.
Incorporating Concurrency Safeguards: Work-in-Progress and Atomicity
To thread-safely accumulate requests, introduce work-in-progress (WIP)—an atomic counter tracking processing state. Oleh explains: increment WIP to claim emission exclusivity; if non-zero, another thread processes, so defer.
Requests add to a requested counter atomically. On WIP decrement to zero, check if more requests pend, resuming if so.
This pattern, akin to semaphores, ensures single-threaded emission despite multi-threaded requests, averting races.
Analytically, it balances responsiveness and safety, though overflows (Long.MAX_VALUE) signal unbounded requests, potentially overwhelming subscribers.
Implications: facilitates non-blocking I/O, vital for high-throughput, but debugging requires tracing atomics.
Verification and Iterative Refinement: Ensuring Spec Compliance
Leverage TCK for exhaustive testing: extend PublisherVerification, supplying working and failing publishers. Tests validate signals, backpressure, and edge cases like negative requests.
Oleh iterates: failures prompt guards, like canceling on invalid requests. Post-fixes, all pass, confirming robustness.
Methodologically, TCK simulates parallelism, exposing flaws early. For custom operators, similar suites verify.
Consequences: empowers library creation or tweaks, as in optimizing for known guarantees, enhancing performance in specific flows.
Extending to Operators and Libraries: Building Beyond Basics
With a compliant publisher, assemble operators chaining transformations. Oleh hints at flux wrappers, where sources like arrays feed pipelines.
Analytically, operators preserve backpressure, propagating requests upstream. This composability yields expressive, efficient streams.
Implications: demystifies internals, aiding contributions to Reactor or custom variants for niches like low-latency trading.
In conclusion, mastering publishers unlocks reactive potential, transforming complex async into manageable flows.
Links:
[DevoxxFR2014] Reactive Programming with RxJava: Building Responsive Applications
Lecturer
Ben Christensen works as a software engineer at Netflix. He leads the development of reactive libraries for the JVM. Ben serves as a core contributor to RxJava. He possesses extensive experience in constructing resilient, low-latency systems for streaming platforms. His expertise centers on applying functional reactive programming principles to microservices architectures.
Abstract
This article provides an in-depth exploration of RxJava, Netflix’s implementation of Reactive Extensions for the JVM. It analyzes the Observable pattern as a foundation for composing asynchronous and event-driven programs. The discussion covers essential operators for data transformation and composition, schedulers for concurrency management, and advanced error handling strategies. Through concrete Netflix use cases, the article demonstrates how RxJava enables non-blocking, resilient applications and contrasts this approach with traditional callback-based paradigms.
The Observable Pattern and Push vs. Pull Models
RxJava revolves around the Observable, which functions as a push-based, composable iterator. Unlike the traditional pull-based Iterable, Observables emit items asynchronously to subscribers. This fundamental duality enables uniform treatment of synchronous and asynchronous data sources:
Observable<String> greeting = Observable.just("Hello", "RxJava");
greeting.subscribe(System.out::println);
The Observer interface defines three callbacks: onNext for data emission, onError for exceptions, and onCompleted for stream termination. RxJava enforces strict contracts for backpressure—ensuring producers respect consumer consumption rates—and cancellation through unsubscribe operations.
Operator Composition and Declarative Programming
RxJava provides over 100 operators that transform, filter, and combine Observables in a declarative manner. These operators form a functional composition pipeline:
Observable.range(1, 10)
.filter(n -> n % 2 == 0)
.map(n -> n * n)
.subscribe(square -> System.out.println("Square: " + square));
The flatMap operator proves particularly powerful for concurrent operations, such as parallel API calls:
Observable<User> users = getUserIds();
users.flatMap(userId -> userService.getDetails(userId), 5)
.subscribe(user -> process(user));
This approach eliminates callback nesting (callback hell) while maintaining readability and composability. Marble diagrams visually represent operator behavior, illustrating timing, concurrency, and error propagation.
Concurrency Control with Schedulers
RxJava decouples computation from threading through Schedulers, which abstract thread pools:
Observable.just(1, 2, 3)
.subscribeOn(Schedulers.io())
.observeOn(Schedulers.computation())
.map(this::cpuIntensiveTask)
.subscribe(result -> display(result));
Common schedulers include:
– Schedulers.io() for I/O-bound operations (network, disk).
– Schedulers.computation() for CPU-bound tasks.
– Schedulers.newThread() for fire-and-forget operations.
This abstraction enables non-blocking I/O without manual thread management or blocking queues.
Error Handling and Resilience Patterns
RxJava treats errors as first-class citizens in the data stream:
Observable risky = Observable.create(subscriber -> {
subscriber.onNext(computeRiskyValue());
subscriber.onError(new RuntimeException("Failed"));
});
risky.onErrorResumeNext(throwable -> Observable.just("Default"))
.subscribe(value -> System.out.println(value));
Operators like retry, retryWhen, and onErrorReturn implement resilience patterns such as exponential backoff and circuit breakers—critical for microservices in failure-prone networks.
Netflix Production Use Cases
Netflix employs RxJava across its entire stack. The UI layer composes multiple backend API calls for personalized homepages:
Observable<Recommendation> recs = userIdObservable
.flatMap(this::fetchUserProfile)
.flatMap(profile -> Observable.zip(
fetchTopMovies(profile),
fetchSimilarUsers(profile),
this::combineRecommendations));
The API gateway uses RxJava for timeout handling, fallbacks, and request collapsing. Backend services leverage it for event processing and data aggregation.
Broader Impact on Software Architecture
RxJava embodies the Reactive Manifesto principles: responsive, resilient, elastic, and message-driven. It eliminates common concurrency bugs like race conditions and deadlocks. For JVM developers, RxJava offers a functional, declarative alternative to imperative threading models, enabling cleaner, more maintainable asynchronous code.