【问题标题】:Custom RxJava Observable emits nothing on subscribe自定义 RxJava Observable 在订阅时不发出任何内容
【发布时间】:2018-12-24 08:10:51
【问题描述】:

我从here 得到了这个方法,效果很好:

@Throws(IOException::class)
fun readTextFromUri(ctx: Context, uri: Uri): String {
    val stringBuilder = StringBuilder()
    ctx.contentResolver.openInputStream(uri)?.use { inputStream ->
        BufferedReader(InputStreamReader(inputStream)).use { reader ->
            var line: String? = reader.readLine()
            while (line != null) {
                stringBuilder.append("$line\n")
                line = reader.readLine()
            }
        }
    }
    return stringBuilder.toString()
}

然后将其转换为使用 Observable 返回每一行的这种方法:

fun getUriAsStringObservable(ctx: Context, uri: Uri): Observable<String> {
    return Observable.create {
        try {
            ctx.contentResolver.openInputStream(uri)?.use { inputStream ->
                BufferedReader(InputStreamReader(inputStream)).use { reader ->
                    var line: String? = reader.readLine()
                    while (line != null) {
                        it.onNext(line)
                        line = reader.readLine()
                    }
                    it.onComplete()
                }
            }
        } catch (e: IOException) {
            it.onError(e)
        }
    }
}

但它并没有像我预期的那样工作,订阅后什么都没有打印出来:

getUriAsStringObservable(this, uri)
    .subscribeOn(Schedulers.io())
    .observeOn(AndroidSchedulers.mainThread())
    .doOnNext {
        print("Next: $it")
    }
    .doOnError {
        print("Error: $it")
    }
    .doOnComplete {
        print("completed")
    }
    .subscribe()

我的错误是什么?

【问题讨论】:

    标签: observable rx-kotlin


    【解决方案1】:

    我找到了三种方法来解决我的问题:

    1) 发射每个项目后使用Thread.sleep(1)

    2) 使用具有非deamon 线程link 的自定义调度程序。

    3) 将 Flowable 与 BackpressureStrategy.BUFFER 一起使用,而不是 Observable(最佳方式)。

    Flowable.create({
        try {
            ctx.contentResolver.openInputStream(uri)?.use { inputStream ->
                BufferedReader(InputStreamReader(inputStream)).use { reader ->
                    var line: String? = reader.readLine()
                    while (line != null) {
                        it.onNext(line)
                        line = reader.readLine()
                    }
                    it.onComplete()
                }
            }
        } catch (e: IOException) {
            it.onError(e)
        }
    }, BackpressureStrategy.BUFFER)
    

    谢谢贾维德

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-03-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多