【问题标题】:Is there a doAfterSubscribe equivalent in RxJava2?RxJava2 中是否有 doAfterSubscribe 等价物?
【发布时间】:2020-04-04 13:33:33
【问题描述】:

我有一个可观察到的流式响应所有请求。我想在发出请求时创建该 observable 的过滤器,以便我可以与多个订阅者共享输出。以下是一些示例代码。

    PublishSubject<String> publishSubject = PublishSubject.create();
    Observable<String> fooObservable = publishSubject.filter(value -> value.startsWith("foo"))
            .doOnSubscribe(disposable -> {
                publishSubject.onNext("foobar");
            })
            .replay(1)
            .refCount();

    fooObservable.subscribe(val -> log.info("A val : <{}>", val));

我使用 PublishSubject 作为我的模拟服务,因为有时该服务会立即返回响应。

我的发现是,因为当有即时结果时没有当前订阅,所以我的 fooObservable 没有被填充。即我想查看时没有日志输出:

A val : <foobar>

请注意,我使用此代码得到相同的结果:

PublishSubject<String> publishSubject = PublishSubject.create();
Observable<String> fooObservable = publishSubject.filter(value -> value.startsWith("foo"))
        .replay(1)
        .refCount();

publishSubject.onNext("foobar");
fooObservable.subscribe(val -> log.info("A val : <{}>", val));

所以问题是 fooObservable 直到订阅后才订阅 PublishSubject,

有没有办法在第一次订阅 fooObservable 后立即运行代码?

编辑: 我想到了类似的东西:

PublishSubject<String> publishSubject = PublishSubject.create();
BehaviorSubject<String> fooObservable = BehaviorSubject.create();
publishSubject.filter(value -> value.startsWith("foo")).subscribe(fooObservable);
publishSubject.onNext("foobar");
fooObservable.subscribe(val -> log.info("A val : <{}>", val));

但是我有 2 个订阅,我不确定如何清理,因为过滤器不返回一次性订阅后的订阅。

编辑 2:后台任务的描述。

我的代码需要订阅第三方服务。该服务在我的代码中调用一个 onResponse 方法,其参数包含我的原始请求和响应。可以随时通过对 onResponse 的新调用来更新响应。

我想为这个提供方法的服务创建一个包装器:

public Observable<Response> getObservable(Request req);

如果请求匹配一个已经订阅的请求,那么 observable 应该在订阅时立即提供最新的匹配值。

当没有订阅者时,我需要取消订阅我正在打包的服务。

【问题讨论】:

  • 好奇,你怎么想把值的发射绑定到订阅事件上呢?为什么不使用BehaviorSubject 之类的东西作为您的模拟服务?你可以直接调用BehaviorSubject.onNext(...),它会将该值发送给后续订阅者。
  • @Trogdor :你是在建议我用 BehaviourSubject 替换 fooObservable 吗?如果是这样,问题编辑是我对这种方法的问题。

标签: java rx-java rx-java2


【解决方案1】:

我认为publish()connect() 是您所需要的。 Here 你可以阅读更多关于它的信息。

在你的情况下,它会是这样的:

PublishSubject<String> publishSubject = PublishSubject.create();
Observable<String> fooObservable = publishSubject
    .filter(value -> value.startsWith("foo"))
    .publish();

fooObservable.subscribe(val -> log.info("A val : <{}>", val));
Observable<Object> o2 = fooObservable.map { new Object() }
Observable<Object> o3 = fooObservable.map { /* Something here*/ }

Disposable disposable = fooObservable.connect()

但请记住disposable.dispose() 以防泄漏

编辑

BehaviorSubject<String> publishSubject = BehaviorSubject.create();
Observable<String> fooObservable = publishSubject.filter(value ->         
value.startsWith("foo"));
fooObservable.subscribe(val -> log.info("A val : <{}>", val));
publishSubject.onNext("foobar");

【讨论】:

  • 似乎不起作用,。现在它错过了在调用 subscribe 和 connect 之前发送的所有更新。
  • @PeterWilkinson 这怎么可能?你在fooObservable.connect()之前的某个地方做publishSubject.onNext(...)吗?在哪里?
  • onNext 由非 rx 服务调用以响应订阅请求。因此在某些情况下可能会同步发生。对该服务的订阅应该只发生一次,这就是为什么我认为共享 observable 上的“afterSubscribe”可能是订阅的好地方。
  • @PeterWilkinson 那么我想唯一的方法是使用BehaviorSubject 而不是PublishSubject (而不是Observable ),因此它将为每个新订户“重复”最后一个人口。
  • 这是我考虑的选项之一。有什么想法可以在我的问题的编辑部分清理订阅吗?
【解决方案2】:

我认为该问题的最新编辑有助于解释您正在尝试做什么,这个 GitHub 链接可能有您正在寻找的答案: https://github.com/ReactiveX/RxJava/issues/4675

但我仍然会分享我的测试代码。我最终遵循了上面链接中的建议并使用PublishSubject.replay(1).refCount()

假设我们有一个第三方服务接口:

interface ThirdPartyService
{
    void subscribe( Consumer<Integer> responseConsumer );

    void unsubscribe();
}

接下来让我们创建一个模拟实现,它将使用随机 int 调用消费者,直到它取消订阅:

    // Mock service to emit a random integer once per second, no Rx:
    ThirdPartyService mockService = new ThirdPartyService() {

        Timer timer = new Timer();

        @Override
        public void subscribe( Consumer<Integer> responseConsumer )
        {
            System.out.println( "Subscribe" );
            Random random = new Random();

            TimerTask task = new TimerTask() {
                @Override
                public void run()
                {
                    int i = random.nextInt( 10 );
                    System.out.println( "Producing: " + i );
                    responseConsumer.accept( i );
                }
            };

            timer.schedule( task, 1000, 1000 );
        }

        @Override
        public void unsubscribe()
        {
            System.out.println( "Unsubscribe" );
            timer.cancel();
        }
    };

接下来,实际的 Rx 管道。假设我只想过滤奇数并在没有观察者时取消订阅服务:

    // Wrap service in a PublishSubject:
    PublishSubject<Integer> subject = PublishSubject.create();
    mockService.subscribe( subject::onNext );

    // Create observable:
    Observable<Integer> observable = subject
            .doFinally( mockService::unsubscribe )
            .filter( i -> i % 2 == 1 )  // Include only odd integers
            .replay( 1 )                // Replay latest to new observers
            .refCount();

最后是手动测试:

    // Subscribe to Observable:
    Disposable sub1 = observable.subscribe( i -> System.out.println( "sub1 got: " + i ));

    // Sleep:
    Thread.sleep( 3300 );

    // Create 2nd Subscriber:
    System.out.println( "adding sub2" );
    Disposable sub2 = observable.subscribe( i -> System.out.println( "sub2 got: " + i ));

    // Sleep:
    Thread.sleep( 3300 );

    // Dispose 2nd Subscriber:
    System.out.println( "disposing sub2" );
    sub2.dispose();

    // Sleep:
    Thread.sleep( 3300 );

    // Dispose 1st Subscriber:
    sub1.dispose();

    // Sleep:
    Thread.sleep( 3300 );

输出:

Subscribe
Producing: 1
sub1 got: 1
Producing: 8
Producing: 6
adding sub2
sub2 got: 1
Producing: 3
sub1 got: 3
sub2 got: 3
Producing: 7
sub1 got: 7
sub2 got: 7
Producing: 1
sub1 got: 1
sub2 got: 1
disposing sub2
Producing: 6
Producing: 7
sub1 got: 7
Producing: 1
sub1 got: 1
Unsubscribe

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2011-07-29
    • 2013-06-21
    • 2014-01-09
    • 2012-02-18
    • 2014-05-02
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多