【问题标题】:Create one observable from several tasks从多个任务中创建一个 observable
【发布时间】:2015-09-07 19:37:47
【问题描述】:

我想创建一个接受S 并返回一个Observable<T> 的方法,该方法在三个(异步)任务完成后以一个值完成。

这三个任务都在队列上运行,它们自己的消费者在不同的线程中。

这就是我的想法:

  • 任务 1 通过队列接收 <S, Subject<T>> 并计算来自 S 的值并使用 onNext 设置在 Subject<T> 中,然后调用 onComplete()
  • 任务 2 和 3 收到 <S, Subject<Void>> 并在完成工作后致电 onComplete()

任务 1、2 和 3 可以独立运行,并以随机顺序完成。

现在我有三个主题,一个T,两个Void。我想返回第一个,但只有在 all 任务完成后才让它发出值T。这是因为在所有任务完成之前,我不希望任何 observable 的订阅者对 T 做任何事情。

结合主题以实现此行为的正确方法是什么?我可以使用CountDownLatch 等轻松破解这个问题,但希望有一种 rx-native 的方式来解决这个问题。

我计划通过队列使用主题作为回调是正确的方法吗?我曾经为此使用CompletableFuture<T>,但我想迁移到RX。

【问题讨论】:

  • 作为一个快速的经验法则 - 如果您在设计中使用主题,您可能做错了什么。如果可能,您应该始终避免主题。

标签: system.reactive rx-java


【解决方案1】:

我不完全确定主题在哪里发挥作用,但您可以使用 When + And + Then 同步三个任务

public IObservable<T> MyMethod<S, T>(S incoming) {

  //Create a new plan
  return Observable.When(

   //Start with Task one which will return an T from an S
   Observable.FromAsync(async () => await SomeTaskToTurnSIntoT(incoming))

   //Add in Task two which returns a System.Reactive.Unit
   .And(Observable.FromAsync(() => /*Do Task 2*/))

   //Same for Task 3
   .And(Observable.FromAsync(() => /*Do Task 3*/))

   //Only emit the item from the first Task.
   .Then((ret, _, __) => ret))

   //Finally we only want this to process once, then we will reuse the
   //existing value for subsequent subscribers
   .PublishLast().RefCount();
}

上面将等到所有三个项目都完成后才会发出。需要注意的一点是,在 Rx 中,Void 对象实际上是 System.Reactive.Unit,因此如果没有值,您应该返回它。

【讨论】:

    【解决方案2】:

    不要为此使用Subjects。而是使用异步调度器合并 observables。所以你有:

    T task1(S s);
    void task2(S s);
    void task3(S s);
    

    然后

    <S,T> Observable<T> get(S s) {
        return Observable.merge(
            Observable.just(s)
                .map(x -> task1(x))
                .subscribeOn(Schedulers.computation()),
            (Observable<T>) Observable.just(s)
                .doOnNext(x -> task2(x))
                .ignoreElements()
                .cast(Object.class)
                .subscribeOn(Schedulers.computation()),
            (Observable<T>) Observable.just(s)
                .doOnNext(x -> task3(x))
                .ignoreElements()
                .cast(Object.class)
                .subscribeOn(Schedulers.computation()))
            // wait for completion before emitting the single value
            .last();
    }   
    

    【讨论】:

    • 我需要 Subjects 的原因(我认为)是另一个任务从队列中运行,因此它会在将来的某个地方计算结果,并且需要将结果值设置为某些东西。所以task(x) 立即完成,因为它只是添加到队列中,但结果稍后会设置到我传递到队列中的主题中。
    • 对,没有正确阅读。合并方法仍然有效。您只需要将 Subjects 封装到正在合并的 observables 中。老实说,在这方面我会全力以赴地使用 Rx,我只会使用调度程序进行排队。您在队列使用方面是否受到限制?
    【解决方案3】:

    您可以编写自己的 Subject 来为您执行此操作,但出于多种原因通常不鼓励这样做。相反,您也可以连接您的任务产生的三个 Observables/Subjects。这个操作会产生一个Observable,它会发出这些任务产生的所有值,并且只有在所有输入的 Observables 都完成后才会完成。

    由于它们的类型略有不同,您需要使用 map() 更改后两个任务生成的 Observable 的签名。

    Observable<T> output = Observable.concat(t1, 
        t2.map(in -> null).ignoreElements(), 
        t3.map(in -> null).ignoreElements());
    

    如果您想等待任何订阅者使用生成的值,直到所有 Observable 完成,您可以在此 Observable 上调用 last() 方法。

    Observable<T> output = Observable.concat(t1, 
        t2.map(in -> null).ignoreElements(), 
        t3.map(in -> null).ignoreElements()).last();
    

    【讨论】:

    • 这不会同时运行三个任务。你会得到序列化的调用。此外,.last() 的值将是 null,而不是来自 t1 的所需值。
    • 这个想法是三个任务通过 Observables t1-t3 发出值/完成。这不会运行这三个任务,我假设他已经自己处理好了。该解决方案简单描述了如何在所有任务完成后获取价值。现在确实会发出null。我会解决的。
    • 我的理解是 OP 说任务 2 和 3 在他们收到 S 时开始。因此,所有三个任务必须在通过S 时同时启动。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-12-29
    • 2017-08-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多