【问题标题】:Queue like Subject in RxJava像 RxJava 中的主题一样排队
【发布时间】:2016-12-22 05:45:16
【问题描述】:

我正在寻找可以:

  1. 如果没有订阅者,可以接收项目并将它们保存在队列或缓冲区中
  2. 一旦我们有了订阅者,所有的项目都会被消耗掉,并且永远不会再次发出
  3. 我可以订阅/取消订阅主题/从主题

BehaviorSubject 几乎可以完成这项工作,但它保留了最后观察到的项目。

更新

根据接受的答案,我为单个观察到的项目制定了类似的解决方案。还添加了取消订阅部分以避免内存泄漏。

class LastEventObservable private constructor(
        private val onSubscribe: OnSubscribe<Any>,
        private val state: State
) : Observable<Any>(onSubscribe) {

    fun emit(value: Any) {
        if (state.subscriber.hasObservers()) {
            state.subscriber.onNext(value)
        } else {
            state.lastItem = value
        }
    }

    companion object {
        fun create(): LastEventObservable {
            val state = State()

            val onSubscribe = OnSubscribe<Any> { subscriber ->
                just(state.lastItem)
                        .filter { it != null }
                        .doOnNext { subscriber.onNext(it) }
                        .doOnCompleted { state.lastItem = null }
                        .subscribe()

                val subscription = state.subscriber.subscribe(subscriber)

                subscriber.add(Subscriptions.create { subscription.unsubscribe() })
            }

            return LastEventObservable(onSubscribe, state)
        }
    }

    private class State {
        var lastItem: Any? = null
        val subscriber = PublishSubject.create<Any>()
    }
}

【问题讨论】:

  • 你看到我的回答了吗?
  • 澄清你的意思是“一旦我们有一个订阅者,所有的项目都被消耗了,再也不会发射了”——如果你有类似的东西:yourSource.take(1).subscribe(),那应该消失吗yourSource 中的所有项目?
  • 不,应该只消耗那个物品。

标签: rx-java


【解决方案1】:

我实现了预期的结果,创建了一个自定义的 Observable,该 Observable 包装了发布主题并在没有附加订阅者的情况下处理发射缓存。看看吧。

public class ExampleUnitTest {
    @Test
    public void testSample() throws Exception {
        MyCustomObservable myCustomObservable = new MyCustomObservable();

        myCustomObservable.emit("1");
        myCustomObservable.emit("2");
        myCustomObservable.emit("3");

        Subscription subscription = myCustomObservable.subscribe(System.out::println);

        myCustomObservable.emit("4");
        myCustomObservable.emit("5");

        subscription.unsubscribe();

        myCustomObservable.emit("6");
        myCustomObservable.emit("7");
        myCustomObservable.emit("8");

        myCustomObservable.subscribe(System.out::println);
    }
}

class MyCustomObservable extends Observable<String> {
    private static PublishSubject<String> publishSubject = PublishSubject.create();
    private static List<String> valuesCache = new ArrayList<>();

    protected MyCustomObservable() {
        super(subscriber -> {
            Observable.from(valuesCache)
                    .doOnNext(subscriber::onNext)
                    .doOnCompleted(valuesCache::clear)
                    .subscribe();

            publishSubject.subscribe(subscriber);
        });
    }

    public void emit(String value) {
        if (publishSubject.hasObservers()) {
            publishSubject.onNext(value);
        } else {
            valuesCache.add(value);
        }
    }
}

希望对您有所帮助!

最好的问候。

【讨论】:

  • 工作就像一个魅力。不过不太喜欢那些静态变量:)
  • 我也不是@MartynasJurkus。我们如何改进这个解决方案?
  • 查看我更新的问题。虽然我的代码在 Kotlin 中。
  • 太棒了!我爱科尔廷。你的最终解决方案比我最初的解决方案漂亮得多。
【解决方案2】:

如果只想等待单个订阅者,请使用UnicastSubject,但请注意,如果您在中间取消订阅,则所有后续排队的项目都将丢失。

编辑:

一旦我们有一个订阅者,所有的项目都会被消耗,并且再也不会发出

对于多个订阅者,请使用ReplaySubject

【讨论】:

  • 应该有多个订阅者的可能性。
  • @MartynasJurkus 您可以将UnicastSubject.publish().autoConnect() 一起使用,但您不能在第一次退订和第二次订阅之间“缓存”项目。但是有库github.com/akarnokd/RxJavaExtensionsUnicastWorkerSubject 允许(但仍然同时有1 个订阅者)或DispatchWorkSubject 允许许多订阅者。不幸的是,您不能指定缓冲区大小,因为在create() 方法中您可以传递capacityHint,这不是指定的容量大小。这只是一个提示,因为 SpscLinkedArrayQueueSpmcLinkedArrayQueue 的实现
【解决方案3】:

我有类似的问题,我的要求是:

  • 应该支持在没有订阅观察者时重播值
  • 应该一次只允许一个观察者订阅
  • 当第一个 Observer 被释放时,应该允许另一个 Observer 订阅

我已将其实现为 RxRelay,但 Subject 的实现将类似:

public final class CacheRelay<T> extends Relay<T> {

    private final ConcurrentLinkedQueue<T> queue = new ConcurrentLinkedQueue<>();
    private final PublishRelay<T> relay = PublishRelay.create();

    private CacheRelay() {
    }

    public static <T> CacheRelay<T> create() {
        return new CacheRelay<>();
    }

    @Override
    public void accept(T value) {
        if (relay.hasObservers()) {
            relay.accept(value);
        } else {
            queue.add(value);
        }
    }

    @Override
    public boolean hasObservers() {
        return relay.hasObservers();
    }

    @Override
    protected void subscribeActual(Observer<? super T> observer) {
        if (hasObservers()) {
            EmptyDisposable.error(new IllegalStateException("Only a single observer at a time allowed."), observer);
        } else {
            for (T element; (element = queue.poll()) != null; ) {
                observer.onNext(element);
            }
            relay.subscribeActual(observer);
        }
    }
}

看看这个Gist for more

【讨论】:

  • 我受这篇文章的启发创建了一个库:github.com/vrendina/RxQueue
  • 如何将工作移至后台线程?队列中的所有操作都在主线程中执行。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-01-02
  • 1970-01-01
  • 1970-01-01
  • 2019-12-16
  • 2011-08-23
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多