You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用Flowable.merge时,如何为指定上游流设置请求响应优先级?

Prioritizing One Source in Flowable.merge() When Both Have Buffered Elements

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 prioritizedMerge method ensures that whenever both sources have buffered elements, the first (high-priority) source’s elements are emitted first.
  • The drainQueues method 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 09:50:39