Reactor自定义订阅等待逻辑实现问询:带重置定时器与超时控制
Alright, let's break down how to implement this nuanced "wait for subscribers with dual timeouts" logic in Reactor. The core idea is to hold off emitting data until either:
- No new subscribers join within a fixed window
tafter the last subscription, OR - The total time since the first subscription exceeds a maximum threshold
t2
Core Approach
We need to track two key timestamps: when the first subscriber joins, and when the last subscriber joined. We'll use a shared trigger signal (a Mono<Void>) that completes as soon as either timeout condition is met. Every new subscriber will reset the inactivity timer (t) and wait for this shared trigger before receiving data.
Implementation Code
Here's a reusable class that wraps any source Flux with this custom logic:
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; import java.time.Duration; import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; public class DelayedStartFlux<T> { private final Flux<T> source; private final Duration inactivityTimeout; // t: Wait after last subscription private final Duration maxWaitTimeout; // t2: Max total wait time // Thread-safe state tracking private final AtomicLong firstSubscribeTime = new AtomicLong(0); private final AtomicLong lastSubscribeTime = new AtomicLong(0); private final AtomicReference<Mono<Void>> startTrigger = new AtomicReference<>(); public DelayedStartFlux(Flux<T> source, Duration inactivityTimeout, Duration maxWaitTimeout) { this.source = source; this.inactivityTimeout = inactivityTimeout; this.maxWaitTimeout = maxWaitTimeout; } public Flux<T> get() { return Flux.defer(() -> { long now = System.currentTimeMillis(); // Update timestamps: set first subscribe time if not already set, update last firstSubscribeTime.compareAndSet(0, now); lastSubscribeTime.set(now); // Manage the shared trigger signal return startTrigger.updateAndGet(currentTrigger -> { // If there's an active inactivity timer, cancel it (reset on new subscriber) if (currentTrigger != null && !currentTrigger.isDisposed()) { currentTrigger.dispose(); } // Create a new trigger that races the two timeout conditions long timeRemainingForMaxWait = maxWaitTimeout.toMillis() - (now - firstSubscribeTime.get()); return Mono.race( // Condition 1: No new subscribers for 'inactivityTimeout' Mono.delay(inactivityTimeout) .then(), // Condition 2: Max wait time exceeded (only if time remains) timeRemainingForMaxWait > 0 ? Mono.delay(Duration.ofMillis(timeRemainingForMaxWait)).then() : Mono.empty() // If max wait already passed, trigger immediately ) .subscribeOn(Schedulers.boundedElastic()) .cache(); // Cache the trigger so all subscribers share the same signal }).thenMany(source); }); } }
Key Details Explained
Flux.defer: Ensures every new subscriber triggers the timestamp update and trigger check, rather than running this logic once when theFluxis created.- Atomic State Variables: Safely track timestamps and the shared trigger across multiple subscriber threads.
Mono.race: Simultaneously waits for either the inactivity timeout or the max wait timeout—whichever comes first triggers the start.- Trigger Reset: Every new subscriber cancels the existing inactivity timer and creates a new one, effectively resetting the
twindow. - Caching the Trigger: The
cache()operator ensures all subscribers wait on the same trigger signal, so everyone starts receiving data at the exact same time.
Usage Example
Let's test this with a simple data source and simulate subscriber timing:
public static void main(String[] args) throws InterruptedException { // Cold source (each subscriber gets their own sequence) Flux<String> dataSource = Flux.just("Item 1", "Item 2", "Item 3"); // Configure: wait 2s after last subscriber, max total wait 5s DelayedStartFlux<String> delayedFlux = new DelayedStartFlux<>( dataSource, Duration.ofSeconds(2), Duration.ofSeconds(5) ); // First subscriber joins at 0s delayedFlux.get().subscribe(item -> System.out.println("Subscriber 1: " + item)); // Second subscriber joins at 1s (resets the inactivity timer) Thread.sleep(1000); delayedFlux.get().subscribe(item -> System.out.println("Subscriber 2: " + item)); // No more subscribers—so we start emitting 2s after the last subscription (1s + 2s = 3s total) Thread.sleep(3000); }
Broadcast (Hot Stream) Variation
If you want all subscribers to receive the exact same data sequence (instead of each getting their own), convert the source to a hot stream in the constructor:
this.source = source.publish().autoConnect();
This way, the data is emitted once and broadcast to all subscribed consumers.
内容的提问来源于stack exchange,提问作者goroncy

