【问题标题】:Asynchronous, early exiting, concatenated Observable异步、提前退出、串联的 Observable
【发布时间】:2014-06-28 00:41:27
【问题描述】:

假设我们有 3 个可观察对象,ABC。我需要同时运行所有 3 个(异步,对于外行),但是:

  1. 如果我从 A 得到任何东西,请发出它...不要发出任何其他东西
  2. 如果 A 完成后没有发出任何内容,则将规则 1 应用于 B
  3. 如果 B 完成时没有发出任何内容,则从 C 发出项目。
  4. 如果 C 完成后没有发出任何内容,则发出默认项。

我昨天花了几个小时试图解决这个问题,但 RxJava 中似乎没有任何操作组合可以让我做到这一点。

你可以想到从左到右的级联值:

A --> B --> C

而且,级联被阻塞,而每个级联运行异步并缓存它们的值。

A(无)--> B(无)--> C(无)--> 默认项

明确地说,A 必须在任何其他观察者发出任何内容之前完成。 B 和 C 的逻辑相同,如果 A、B、C 未能发出任何内容,则默认为默认值。

显然涉及到缓存,我绝对不想重播 observable。我将需要重播缓存的值。每个门口都挂着。

该行为与concat() 极为相似,只是如果之前有排放,则不会释放链的下一部分。

【问题讨论】:

    标签: algorithm rx-java


    【解决方案1】:

    这是我想出的:

    **
     * Works like {@link rx.Observable#concat} but concatenated Observables
     * are all run immediately on their given {@link rx.Scheduler}.
     *
     * This Observable is blocking in the sense that items are emitted in order
     * like {@link rx.Observable#concat} but since each Observable is run on
     * an (possibly) asynchronous scheduler, items emitted further down the chain
     * of Observables are held until items further up the chain are (possibly) emitted.
     *
     * This Observable also short-circuits and does not emit items further down
     * the chain of Observables when an Observable higher up the chain emits items.
     *
     * For example:
     *
     * Given Observable A, B, and C
     *
     * If A emits item(s) emit them... do not emit anything else.
     * If A completes without emitting anything, apply previous rule to B.
     * If B completes without emitting anything, emit items from C (if any)
     *
     * @param <T>
     */
    public class ConcatObservable<T> {
      private final List<Observable<? extends T>> observables;
    
      private ConcatObservable(List<Observable<? extends T>> observables) {
        this.observables = observables;
      }
    
      public static <T> ConcatObservable<T> from(Observable<? extends T>... observables) {
        return new ConcatObservable<T>(Arrays.asList(observables));
      }
    
      public Observable<T> asObservable() {
        final List<Subscription> subscriptions = new CopyOnWriteArrayList<Subscription>();
    
        return Observable.create(new Observable.OnSubscribe<T>() {
          @Override public void call(final Subscriber<? super T> subscriber) {
            List<Observable<? extends T>> cachedObservables = new ArrayList<Observable<? extends T>>();
            for (Observable<? extends T> observable : observables) {
    
              // tell it to cache values
              final ReplaySubject<T> subject = ReplaySubject.create();
              cachedObservables.add(subject);
    
              // run it with nobody listening
              Subscription subscription = observable.subscribe(new Observer<T>() {
                @Override public void onCompleted() {
                  subject.onCompleted();
                }
    
                @Override public void onError(Throwable e) {
                  subject.onError(e);
                }
    
                @Override public void onNext(T item) {
                  subject.onNext(item);
                }
              });
              subscriptions.add(subscription);
            }
    
            final AtomicReference<Throwable> error = new AtomicReference<Throwable>();
    
            // for the cached ones, already running
            for (Observable<? extends T> observable : cachedObservables) {
    
              final AtomicBoolean shouldExit = new AtomicBoolean(false);
              final CountDownLatch latch = new CountDownLatch(1);
              Subscription subscription = observable.subscribe(new Observer<T>() {
                @Override public void onCompleted() {
                  latch.countDown();
                }
    
                @Override public void onError(Throwable e) {
                  error.set(e);
                  shouldExit.set(true);
                  latch.countDown();
                }
    
                @Override public void onNext(T item) {
                  subscriber.onNext(item);
                  shouldExit.set(true);
                }
              });
    
              // Track each subscription
              subscriptions.add(subscription);
    
              try {
                // Wait for this one to stop emitting, or error
                latch.await();
              } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                throw new RuntimeException("Interrupted while waiting for subscription to complete", e);
              }
    
              // This one had an item(s), so we don't bother with the rest
              if (shouldExit.get()) {
                break;
              }
            }
    
            // Release inner subscriptions
            for (Subscription subscription : subscriptions) {
              subscription.unsubscribe();
            }
    
            // Obey the Observable contract...
            Throwable throwable = error.get();
            if (throwable != null) {
              subscriber.onError(throwable);
            } else {
              subscriber.onCompleted();
            }
          }
        });
      }
    }
    

    下面是相应的测试:

    public class ConcatObservableTest {
    
      @Test @SuppressWarnings("unchecked")
      public void it_onlyEmitsFromFirstObservable() {
        Observable<String> A = Observable.from(Arrays.asList("A", "A", "A"));
        Observable<String> B = Observable.from(Arrays.asList("B", "B", "B"));
        Observable<String> C = Observable.from(Arrays.asList("C", "C", "C"));
    
        Observable<String> observable = ConcatObservable.from(A, B, C).asObservable();
    
        TestSubscriber<String> testSubscriber = new TestSubscriber<String>();
        observable.subscribe(testSubscriber);
    
        assertThat(testSubscriber.getOnNextEvents()).containsExactly("A", "A", "A");
      }
    
      @Test @SuppressWarnings("unchecked")
      public void it_onlyEmitsFromSecondObservable() {
        Observable<String> A = Observable.empty();
        Observable<String> B = Observable.from(Arrays.asList("B", "B", "B"));
        Observable<String> C = Observable.from(Arrays.asList("C", "C", "C"));
    
        Observable<String> observable = ConcatObservable.from(A, B, C).asObservable();
    
        TestSubscriber<String> testSubscriber = new TestSubscriber<String>();
        observable.subscribe(testSubscriber);
    
        assertThat(testSubscriber.getOnNextEvents()).containsExactly("B", "B", "B");
      }
    
      @Test @SuppressWarnings("unchecked")
      public void it_onlyEmitsFromLastObservable() {
        Observable<String> A = Observable.empty();
        Observable<String> B = Observable.empty();
        Observable<String> C = Observable.from(Arrays.asList("C", "C", "C"));
    
        Observable<String> observable = ConcatObservable.from(A, B, C).asObservable();
    
        TestSubscriber<String> testSubscriber = new TestSubscriber<String>();
        observable.subscribe(testSubscriber);
    
        assertThat(testSubscriber.getOnNextEvents()).containsExactly("C", "C", "C");
      }
    
      @Test @SuppressWarnings("unchecked")
      public void it_shouldStartAllObservables() {
        TestObservable<String> letters = TestObservable.createTestObservable("A", "B", "C");
        TestObservable<String> numbers = TestObservable.createDelayedTestObservable(100, "1", "2", "3");
        TestObservable<String> animals = TestObservable.createDelayedTestObservable(200, "zebra", "donkey", "unicorn");
    
        Observable<String> observable = ConcatObservable.from(letters, numbers, animals).asObservable();
    
        TestSubscriber<String> testSubscriber = new TestSubscriber<String>();
        observable.subscribe(testSubscriber);
    
        assertThat(letters.isCalled()).isTrue();
        assertThat(numbers.isCalled()).isTrue();
        assertThat(animals.isCalled()).isTrue();
      }
    
      static class TestObservable<T> extends Observable<T> {
        private final TestOnSubscribe<T> onSubscribeFunc;
    
        private TestObservable(TestOnSubscribe<T> f) {
          super(f);
          onSubscribeFunc = f;
        }
    
        public boolean isCalled() {
          return onSubscribeFunc.isCalled();
        }
    
        @SuppressWarnings("unchecked")
        public static <T> TestObservable<T> createTestObservable(final T... items) {
          return createDelayedTestObservable(0, items);
        }
    
        @SuppressWarnings("unchecked")
        public static <T> TestObservable<T> createDelayedTestObservable(final long delay, final T... items) {
          return new TestObservable<T>(new TestOnSubscribe<T>(delay, items));
        }
    
        private static class TestOnSubscribe<T> implements OnSubscribe<T> {
          private final long delay;
          private final T[] items;
          private boolean isCalled;
    
          private TestOnSubscribe(long delay, T... items) {
            this.delay = delay;
            this.items = items;
          }
    
          @Override public void call(Subscriber<? super T> subscriber) {
            isCalled = true;
    
            for (T item : items) {
              if (delay > 0) {
                sleep(delay);
              }
              subscriber.onNext(item);
            }
            subscriber.onCompleted();
          }
    
          public boolean isCalled() {
            return isCalled;
          }
    
          private void sleep(long time) {
            try {
              Thread.sleep(time);
            } catch (InterruptedException e) { }
          }
        }
      }
    }
    

    【讨论】:

    • 抱歉回复晚了。 ReplaySubject 是一个聪明的主意。刚刚在您的代码中发现了一个问题:您没有正确设置订阅。请查看我的更新答案。
    • 你能告诉我他们怎么不正确吗?我的测试都表明这工作正常。我已将我的测试添加到我的答案中。
    • 我的意思是即使输入List&lt;Observable&lt;? extends T&gt;&gt; observables支持退订,你也不能退订,因为你没有将他们返回的订阅添加到Subscriber
    • 更准确地说,用户可能希望在任何List&lt;Observable&lt;? extends T&gt;&gt; observables 开始运行之前取消订阅
    【解决方案2】:
    public class ConcatObservable<T> {
    
    private final List<Observable<? extends T>> observables;
    
    private ConcatObservable(List<Observable<? extends T>> observables) {
        this.observables = observables;
    }
    
    public static <T> ConcatObservable<T> from(Observable<? extends T>... observables) {
        return new ConcatObservable<T>(Arrays.asList(observables));
    }
    
    public Observable<T> asObservable() {
        return Observable.create(new Observable.OnSubscribe<T>() {
            @Override
            public void call(final Subscriber<? super T> subscriber) {
                List<Observable<? extends T>> cachedObservables = new ArrayList<Observable<? extends T>>();
                for (Observable<? extends T> observable : observables) {
                    ConnectableObservable<? extends T> replayedObservable = observable.replay();
                    cachedObservables.add(replayedObservable);
                    subscriber.add(replayedObservable.connect());
                }
                Subscription s = Observable.concat(Observable.from(cachedObservables)).take(1).subscribe(subscriber);
                subscriber.add(s);
            }
        });
    }
    }
    

    编辑

    这很接近,但它没有通过以下测试:

    @Test @SuppressWarnings("unchecked")
    public void it_onlyEmitsFromFirstObservable() {
      Observable<String> A = Observable.from(Arrays.asList("A", "A", "A"));
      Observable<String> B = Observable.from(Arrays.asList("B", "B", "B"));
      Observable<String> C = Observable.from(Arrays.asList("C", "C", "C"));
    
      Observable<String> observable = ConcatObservable.from(A, B, C).asObservable();
    
      TestSubscriber<String> testSubscriber = new TestSubscriber<String>();
      observable.subscribe(testSubscriber);
    
      assertThat(testSubscriber.getOnNextEvents()).containsExactly("A", "A", "A");
    }
    

    【讨论】:

    • 你能在下面查看我的答案吗?我对它进行了单元测试,似乎工作正常。
    • 你能不能用你的单元测试来测试这个新的答案?
    • 我认为这不能满足我的需求。当我说如果 A、B 或 C 发出任何东西时,就不会发出任何其他东西,我的意思是只从那个可观察对象发出,而没有其他东西发出。我的意思不是发射单个项目。
    • 我通过失败的测试编辑了您的答案,以显示我正在寻找的行为。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-12-20
    • 2011-03-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多