代码之家  ›  专栏  ›  技术社区  ›  Jongz Puangput

FlatMap在主线程上运行,甚至使用Scheuler.io

  •  0
  • Jongz Puangput  · 技术社区  · 6 年前

    我已经测试过了,肯定是主线程发出的值导致了这个问题。但是,我想知道我是否有这个用例从主线程接收一些值来继续Rx流。为了使flatMap在不同于main的线程中运行,应该做些什么。

    class MainActivity : AppCompatActivity() {
    
        private lateinit var emitter: ObservableEmitter<String>
    
        override fun onCreate(savedInstanceState: Bundle?) {
            super.onCreate(savedInstanceState)
            setContentView(R.layout.activity_main)
    
            btnFlatMap.setOnClickListener {
    
                val obs = Observable.create<String> {
                    emitter = it
                    logThread("inside observable")
    
                    // TODO: fetch some configuration from the internet or local db
    
                    // TODO: then call startActivityForResult()
                }
    
                obs
                    .flatMap {
                        logThread("flatMap, Banana")
                        Observable.just("$it, 1 Item")
                    }
                    .subscribeOn(Schedulers.io())
                    .observeOn(AndroidSchedulers.mainThread())
                    .subscribe({ next ->
                        logThread("onNext")
                    }, { error ->
                        logThread("onError")
                    }, {
                        logThread("onComplete")
                    })
            }
    
            btnEmitter.setOnClickListener {
                // TODO: simulate that onActivityResult is called
                emitter.onNext("Banana")
                emitter.onComplete()
            }
        }
    
        private fun logThread(operation: String) {
            Log.e("THREAD", "$operation run at [${Thread.currentThread().name}]")
        } }
    

    当前日志

    在[RxCachedThreadScheduler-1]上运行inside observable

    平面图,香蕉跑在[主要]

    onNext运行于[main]

    在[main]完成运行

    需要Logcat

    在[RxCachedThreadScheduler-1]上运行inside observable

    flatMap,香蕉运行在[RxCachedThreadScheduler-1]

    onNext运行于[main]

    在[main]完成运行

    1 回复  |  直到 6 年前
        1
  •  0
  •   Jongz Puangput    6 年前

    再加一个 onserverOn(Schedulers.io()) 诀窍是告诉Rx将线程切换回工作线程

            obs
                .observeOn(Schedulers.io()) <------ additional 
                .flatMap {
                    logThread("flatMap, Banana")
                    Observable.just("$it, 1 Item")
                }
                .subscribeOn(Schedulers.io())
                .observeOn(AndroidSchedulers.mainThread())
                .subscribe({ next ->
                    logThread("onNext")
                }, { error ->
                    logThread("onError")
                }, {
                    logThread("onComplete")
                })
    

    因为从主线程发出的数据,它导致RX从工作线程(Schedulers.io)切换整个下游,而工作线程是由 subscribeOn() 同时订阅主线程。为了再次将整个下游的(切换)执行线程更改为工作线程 observeOn() 就是为了这个目的。

    简而言之

    • observeOn ,将下游更改为在特定线程上执行。

    • subscribeOn ,将上游(根源)设置为在特定线程上执行。

    现在是Logcat

    flatMap,香蕉运行在[RxCachedThreadScheduler-1]

    onNext运行于[main]

    在[main]完成运行