【问题标题】:Have some trouble with thread work线程工作有一些问题
【发布时间】:2018-02-25 10:20:48
【问题描述】:

我的事件发射器类有代码:

 private val socketListeners: ArrayList<SocketContentListener> = ArrayList()

  //add listener here
 override fun subscribe(socketListener: SocketContentListener) {
        socketListeners.add(socketListener)
    }

       private fun getSocketConnectListener()
                : SocketContentListener {
            /**
             * Post received messages to listeners via Handler
             * because handler helps to set all messages in order on main thread.
             */
            return object : SocketContentListener {

                override fun onUdpServerListenerCreated(inetAddress: InetAddress?, port: Int) {
                    val subscribers = ArrayList<SocketContentListener>(socketListeners)
                    for (listener in subscribers) {
                        Handler(Looper.getMainLooper()).post({ listener.onUdpServerListenerCreated(inetAddress, port) })
                    }
            }
        }

我尝试创建 Observable:

val udpObservable = Observable.create<Int> { emitter ->
        val listener = object : SocketListener() {
            override fun onUdpServerListenerCreated(inetAddress: InetAddress, port: Int) {
                emitter.onNext(port)
                emitter.onComplete()
            }

        }
        //add listener here
        socketSource.subscribe(listener)
        emitter.setCancellable { socketSource.unSubscribe(listener) }
    }.subscribeOn(Schedulers.io())
            .doOnNext { Log.d("123-thread", "current is: " + Thread.currentThread().name) }
            .onErrorReturn { throw ConnectionException(it) }
            .subscribe()

但在测试期间,我看到的不是预期的RxCachedThreadScheduler-1 thread 工作

  D/123-thread: current is:-> main

那你能帮帮我吗?请。我的错误在哪里?如何为 rx 链实现所需的 RxCachedThreadScheduler 线程?

【问题讨论】:

    标签: java android multithreading kotlin rx-java2


    【解决方案1】:

    创建 Observable 的代码会在调度器上执行,没有隐式的上下文变化。

     您的事件从主线程上的侦听器到达。然后将它们发送到同一线程上的发射器。

    所以,订阅在调度程序 io 线程上进行,但发射器在主线程上进行

    所以解决方法是在 Observable create 之后添加observerOn(Schedulers.newThread())。 Like this

    val udpObservable = Observable.create<Int> { emitter ->
            val listener = object : SocketListener() {
                override fun onUdpServerListenerCreated(inetAddress: InetAddress, port: Int) {
                    emitter.onNext(port)
                    emitter.onComplete()
                }
    
            }
            //add listener here
            socketSource.subscribe(listener)
            emitter.setCancellable { socketSource.unSubscribe(listener) }
        }.subscribeOn(Schedulers.io())
    
          //need add this for work with emmit data on background 
         .observerOn(Schedulers.newThread())
                .doOnNext { Log.d("123-thread", "current is: " + Thread.currentThread().name) }
                .onErrorReturn { throw ConnectionException(it) }
                .subscribe()
    

    【讨论】:

    • 感谢@Sergey,下次尝试使用代码而不是图像。您可以使用 cmets 来注释您使用红色箭头所做的事情。
    • @GoRoS 我会考虑你的fidbek,认为你是对的
    猜你喜欢
    • 1970-01-01
    • 2011-06-18
    • 2020-05-18
    • 2023-04-02
    • 1970-01-01
    • 2011-04-18
    • 1970-01-01
    • 2013-09-16
    • 2010-11-04
    相关资源
    最近更新 更多