使用Flowable.merge时,如何为指定上游流设置请求响应优先级?
Great question! The default behavior of Flowable.merge() doesn’t prioritize one source over the other—it uses a round-robin strategy to alternate between available sources when both have elements ready. So with your current code:
Flowable.merge(sources, 2, 1).observeOn(Schedulers.io(), false, 1)
when downstream calls request(1), it won’t automatically prioritize the first source’s element if both sources have buffered data.
How to Add Priority for the First Source
To enforce that the first source’s elements are emitted before the second’s (whenever both have data available), you’ll need to implement a custom merging logic that explicitly checks the first source’s buffer first. Here’s a practical implementation using Flowable.create() to handle the prioritization manually:
Step 1: Define a Custom Merging Flowable
This implementation maintains separate buffers for each source, and when draining elements (in response to downstream requests), it always checks the first source’s buffer first:
import io.reactivex.rxjava3.core.Flowable; import io.reactivex.rxjava3.core.FlowableEmitter; import io.reactivex.rxjava3.disposables.Disposable; import io.reactivex.rxjava3.disposables.Disposables; import java.util.Queue; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.atomic.AtomicBoolean; public class PrioritizedMergeExample { public static <T> Flowable<T> prioritizedMerge(Flowable<T> highPrioritySource, Flowable<T> lowPrioritySource) { return Flowable.create(emitter -> { // Track completion status of each source AtomicBoolean highPriorityDone = new AtomicBoolean(false); AtomicBoolean lowPriorityDone = new AtomicBoolean(false); // Buffers to hold elements from each source Queue<T> highPriorityQueue = new ConcurrentLinkedQueue<>(); Queue<T> lowPriorityQueue = new ConcurrentLinkedQueue<>(); // Subscribe to high-priority source (first source) Disposable highDisposable = highPrioritySource.subscribe( item -> { highPriorityQueue.offer(item); drainQueues(emitter, highPriorityQueue, lowPriorityQueue, highPriorityDone, lowPriorityDone); }, emitter::onError, () -> highPriorityDone.set(true) ); // Subscribe to low-priority source (second source) Disposable lowDisposable = lowPrioritySource.subscribe( item -> { lowPriorityQueue.offer(item); drainQueues(emitter, highPriorityQueue, lowPriorityQueue, highPriorityDone, lowPriorityDone); }, emitter::onError, () -> lowPriorityDone.set(true) ); // Composite disposable to manage both subscriptions emitter.setDisposable(Disposables.composite(highDisposable, lowDisposable)); // Handle downstream requests by draining queues emitter.setRequestHandler(n -> drainQueues( emitter, highPriorityQueue, lowPriorityQueue, highPriorityDone, lowPriorityDone )); }, io.reactivex.rxjava3.core.BackpressureStrategy.BUFFER); } private static <T> void drainQueues( FlowableEmitter<T> emitter, Queue<T> highPriorityQueue, Queue<T> lowPriorityQueue, AtomicBoolean highPriorityDone, AtomicBoolean lowPriorityDone ) { while (!emitter.isCancelled()) { // First, try to emit from the high-priority queue T item = highPriorityQueue.poll(); if (item != null) { emitter.onNext(item); continue; } // If high-priority queue is empty, check low-priority item = lowPriorityQueue.poll(); if (item != null) { emitter.onNext(item); continue; } // Both queues are empty—check if both sources are complete if (highPriorityDone.get() && lowPriorityDone.get()) { emitter.onComplete(); break; } // No elements available, exit loop and wait for new items or requests break; } } }
Step 2: Use the Custom Merge with Your Subjects
If your sources are Subject instances (like PublishSubject), you can use this method directly:
import io.reactivex.rxjava3.subjects.PublishSubject; import io.reactivex.rxjava3.schedulers.Schedulers; public class Main { public static void main(String[] args) { PublishSubject<String> source1 = PublishSubject.create(); PublishSubject<String> source2 = PublishSubject.create(); PrioritizedMergeExample.prioritizedMerge(source1, source2) .observeOn(Schedulers.io(), false, 1) .subscribe(item -> System.out.println("Emitted: " + item)); // Test with both sources having data source1.onNext("Source 1 - Item 1"); source2.onNext("Source 2 - Item 1"); // Downstream requests 1 element—will emit Source 1's item first ((io.reactivex.rxjava3.core.FlowableSubscriber<String>) subscriber).request(1); } }
Key Notes
- The custom
prioritizedMergemethod ensures that whenever both sources have buffered elements, the first (high-priority) source’s elements are emitted first. - The
drainQueuesmethod handles the core logic: it always checks the high-priority buffer before the low-priority one, only falling back to the second source when the first is empty. - This implementation respects backpressure, so it will only emit elements as requested by the downstream subscriber.
内容的提问来源于stack exchange,提问作者lubo-pisk

