【问题标题】:onNext not called in Retrofit and RxJava, when running in thread在线程中运行时,在 Retrofit 和 RxJava 中未调用 onNext
【发布时间】:2016-07-17 17:20:08
【问题描述】:

我正在尝试在使用 RxJava 时从 Retrofit 2.1 向 Spotify Web API 发出请求。我希望每个请求都在其自己的线程中发生,并且在它们准备好时以任何顺序打印结果。

当前代码在主线程上执行并工作。但是,当我在 Rx 链中插入 .subscribeOn(Schedulers.newThread()) 时,我得到一个空输出。 似乎从未调用过 onNext()。

我发现这通常可以通过在 Android 中插入 .observeOn(AndroidSchedulers.mainThread()) 来解决,但我如何在常规 Java 8(没有 RxAndroid)中做到这一点?

public interface SpotifyService {

    @GET("tracks/{trackId}")
    Observable<SpotifyTrack> getTrack(@Path("trackId") String id);

    Retrofit retrofit = new Retrofit.Builder()
            .baseUrl("https://api.spotify.com/v1/")
            .addConverterFactory(JacksonConverterFactory.create())
            .addCallAdapterFactory(RxJavaCallAdapterFactory.create())
            .build();
}


public class Main {

    private static SpotifyService spotifyService = SpotifyService.retrofit.create(SpotifyService.class);

    public static void main(String[] args) {
        String[] trackIds = {
            "spotify:track:2HUI2s84pkL5815G8WI1Lg",
            "spotify:track:1bZrI1KgVKr8Qfja9cnmGh",
            "spotify:track:1WP1r7fuvRqZRnUaTi2I1Q",
            "spotify:track:5kqIPrATaCc2LqxVWzQGbk",
            "spotify:track:0mWiuXuLAJ3Brin3Or2x6v",
            "spotify:track:0BF6mdNROWgYo3O3mNGrBc"
        };

        printTrackNames(trackIds);
    }

    private static void printTrackNames(String[] trackIds) {
        Observable.from(trackIds)
            .map(Main::toTrackId)
            .flatMap(spotifyService::getTrack)
            .map(SpotifyTrack::getName)
            .subscribe(System.out::println);
    }

    private static String toTrackId(String track) {
        if (track.contains("spotify:track:")) {
            return track.split(":")[2];
        }

        return track;
    }
}

【问题讨论】:

  • 那是因为您的 main() 方法返回,并且进程已关闭(调度程序使用守护线程)。出于演示目的,您可以阻止当前线程,并可能获得预期的结果。一种方法是在 printTrackNames() 之后添加 Thread.sleep() 或延迟阻塞 observable:Observable.empty().delay(10,TimeUnit.SECONDS).toBlocking().subscribe();
  • 当然。我没有想到这一点。感谢您的回答!如果你把它变成一个答案,我会把它标记为答案。

标签: rx-java retrofit2


【解决方案1】:

由于您没有像在 Android 中那样继续执行的运行循环,因此您需要阻塞主线程,以便它不会退出 main() 并在订阅完成之前死亡。您可以使用toBlocking()

【讨论】:

    【解决方案2】:

    这不是说不会去 onNext,而是因为它在另一个线程中,所以你无法调试它。

    为了证明这一点,使用 TestSubscriber 等待观察者结束。

    这里有一个单元测试来证明它。

    @Test
    public void testObservableAsync() throws InterruptedException {
        Subscription subscription = Observable.from(numbers)
                                              .subscribeOn(Schedulers.newThread())
                                              .subscribe(number -> System.out.println("Items emitted:" + total));
        System.out.println("I finish before the observable finish.  Items emitted:" + total);
        new TestSubscriber((Observer) subscription)
                .awaitTerminalEvent(100, TimeUnit.MILLISECONDS);
    }
    

    如果想看更多异步测试,请看这里https://github.com/politrons/reactive/blob/master/src/test/java/rx/observables/scheduler/ObservableAsynchronous.java

    【讨论】:

      猜你喜欢
      • 2018-09-09
      • 1970-01-01
      • 1970-01-01
      • 2020-02-15
      • 1970-01-01
      • 1970-01-01
      • 2018-02-20
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多