kotlin让订阅者使用RxJava2观察可观察量

ant*_*009 2 kotlin rx-java2

Android Studio 3.0 Beta2
Run Code Online (Sandbox Code Playgroud)

我创建了2个方法,一个创建了observable,另一个创建了订阅者.

但是,我遇到了一个问题,试图让订阅者订阅observable.在Java中,这将起作用,我试图让它在Kotlin中工作.

在我的onCreate(..)方法中,我试图设置它.这是正确的方法吗?

class MainActivity : AppCompatActivity() {

    override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        setContentView(R.layout.activity_main)

        /* CANNOT SET SUBSCRIBER TO SUBCRIBE TO THE OBSERVABLE */
        createStringObservable().subscribe(createStringSubscriber())
    }


    fun createStringObservable(): Observable<String> {
        val myObservable: Observable<String> = Observable.create {
            subscriber ->
            subscriber.onNext("Hello, World!")
            subscriber.onComplete()
        }

        return myObservable
    }

    fun createStringSubscriber(): Subscriber<String> {
        val mySubscriber = object: Subscriber<String> {
            override fun onNext(s: String) {
                println(s)
            }

            override fun onComplete() {
                println("onComplete")
            }

            override fun onError(e: Throwable) {
                println("onError")
            }

            override fun onSubscribe(s: Subscription?) {
                println("onSubscribe")
            }
        }

        return mySubscriber
    }
}
Run Code Online (Sandbox Code Playgroud)

非常感谢任何建议,

hom*_*man 7

密切关注类型.

Observable.subscribe() 有三个基本变体:

  • 一个不接受任何论据的人
  • 几个接受一个 io.reactivex.functions.Consumer
  • 一个接受一个 io.reactivex.Observer

您在示例中尝试订阅的类型org.reactivestreams.Subscriber(定义为Reactive Streams规范的一部分).您可以参考文档来更全面地了解这种类型,但足以说它与任何重载Observable.subscribe()方法都不兼容.

这是一个createStringSubscriber()允许代码编译的方法的修改示例:

fun createStringSubscriber(): Observer<String> {
        val mySubscriber = object: Observer<String> {
            override fun onNext(s: String) {
                println(s)
            }

            override fun onComplete() {
                println("onComplete")
            }

            override fun onError(e: Throwable) {
                println("onError")
            }

            override fun onSubscribe(s: Disposable) {
                println("onSubscribe")
            }
        }

        return mySubscriber
    }
Run Code Online (Sandbox Code Playgroud)

改变的是:

  1. 这会返回一个Observer类型(而不是Subscriber)
  2. onSubscribe()通过a Disposable(代替Subscription)

..正如'Vincent Mimoun-Prat'所提到的,lambda语法可以真正缩短您的代码.

    override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        setContentView(R.layout.activity_main)

        // Here's an example using pure RxJava 2 (ie not using RxKotlin)
        Observable.create<String> { emitter ->
            emitter.onNext("Hello, World!")
            emitter.onComplete()
        }
                .subscribe(
                        { s -> println(s) },
                        { e -> println(e) },
                        {      println("onComplete") }
                )

        // ...and here's an example using RxKotlin. The named arguments help
        // to give your code a little more clarity
        Observable.create<String> { emitter ->
            emitter.onNext("Hello, World!")
            emitter.onComplete()
        }
                .subscribeBy(
                        onNext     = { s -> println(s) },
                        onError    = { e -> println(e) },
                        onComplete = {      println("onComplete") }
                )
    }
Run Code Online (Sandbox Code Playgroud)

我希望有所帮助!