如何让Observable.zip及Interval运行在IntentService线程(RxJava2)
Alright, let's break down your problem and fix it. You're seeing RxJava operations like zip and interval running on RxJava's own thread pool (e.g., RxSingleScheduler-1) instead of the IntentService thread, and you want everything to stick to the IntentService's worker thread. Here's what's going on and how to fix it:
What's Causing the Thread Switch?
RxJava defaults to using its own thread pools for operations that require scheduling (like interval). Even though you tried a custom Executor, it wasn't enough to force all operations onto the IntentService thread. Additionally, if you use asynchronous subscriptions (subscribeWith), RxJava will offload work to its threads, and the IntentService thread might finish before your Rx stream completes—leading to unexpected behavior or errors.
Solution: Force All Rx Operations to Run on IntentService Thread
To make sure every part of your Rx stream runs on the IntentService thread, you need two key changes:
- Run the Rx stream synchronously so it blocks the IntentService thread until completion (no thread switches).
- Adjust the
intervalscheduler to use the current thread instead of RxJava's pool.
Modified Code
@Override protected void onHandleIntent(@Nullable Intent intent) { LOG.debug(Thread.currentThread().getName()); Call<String> call = mRepository.getInfo(); try { retrofit2.Response<String> response = call.execute(); if (response.isSuccessful()) { LOG.debug("Response body " + Thread.currentThread().getName()); getInfoAboutUser(); } } catch (Exception e) { LOG.error("Failed to fetch initial info", e); } } public void getInfoAboutUser() { LOG.debug("getInfoAboutUser " + Thread.currentThread().getName()); Observable.zip( Observable.fromIterable(array), // Use Schedulers.trampoline() to run interval on the current thread Observable.interval((mRandom.nextInt(7) + 5) * 1000, TimeUnit.MILLISECONDS, Schedulers.trampoline()) .take(array.size()), new BiFunction<String, Long, String>() { @Override public String apply(String s, Long aLong) throws Exception { LOG.debug("Result " + Thread.currentThread().getName()); return s; } } ).flatMapMaybe(new Function<String, MaybeSource<String>>() { @Override public MaybeSource<String> apply(String s) throws Exception { // Since getInfoAboutUser(s) has no subscribeOn, it runs on the current (IntentService) thread return mRepository.getInfoAboutUser(s); } }) // Use blockingSubscribe to run the entire stream synchronously on the current thread .blockingSubscribe( result -> { // Handle each result here }, error -> { LOG.error("Error in user info stream", error); }, () -> { LOG.debug("User info stream completed successfully"); } ); }
Key Details
blockingSubscribe(): This replaces your asynchronoussubscribeWithcall. It blocks the IntentService thread until the entire Rx stream finishes, ensuring every operation runs on this thread instead of RxJava's pools.Schedulers.trampoline(): This tellsintervalto trigger its timed events on the current thread (IntentService thread) instead of switching to RxJava's computation scheduler. Note that this will block the IntentService thread during each interval wait, which is exactly what you want since you're restricted to using only this thread.- No extra
Executor: Your customExecutorwasn't necessary becauseSchedulers.trampoline()already handles running tasks on the current thread.
After these changes, all your log statements should show the IntentService thread name for every step of the Rx stream.
内容的提问来源于stack exchange,提问作者Delphian

