以下代码
package com.inthemoon.snippets.rxjava;
import io.reactivex.*;
public class HelloWorld {
public static void main(String[] args) {
Flowable.just("Hello world").subscribe(System.out::println);
}
}
Run Code Online (Sandbox Code Playgroud)
导致以下编译错误
错误:(9、15)Java:无法访问org.reactivestreams.Publisher的org.reactivestreams.Publisher类文件
POM依赖关系如下
<dependencies>
<!-- https://mvnrepository.com/artifact/io.reactivex.rxjava2/rxjava -->
<dependency>
<groupId>io.reactivex.rxjava2</groupId>
<artifactId>rxjava</artifactId>
<version>2.0.4</version>
</dependency>
</dependencies>
Run Code Online (Sandbox Code Playgroud) 我有以下类,我想返回Subscription对象或其他东西,所以我可以从我引用subscribe()方法的地方取消请求,但是subscribe(observer)返回void!我怎样才能做到这一点?
public abstract class MainPresenter<T> {
protected <T> Disposable subscribe(Observable<T> observable, Observer<T> observer) {
observable.subscribeOn(Schedulers.newThread())
.observeOn(AndroidSchedulers.mainThread())
.subscribe(observer);
}
Run Code Online (Sandbox Code Playgroud)
[ 新更新 ]我用这种方式暂时,我在等待更好的解决方案:
protected <T> DisposableMaybeObserver subscribe(final Maybe<T> observable,
final Observer<T> observer) {
return observable.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribeWith(new DisposableMaybeObserver<T>() {
@Override
public void onSuccess(T t) {
observer.onNext(t);
}
@Override
public void onError(Throwable e) {
observer.onError(e);
}
@Override
public void onComplete() {
observer.onComplete();
}
});
}
Run Code Online (Sandbox Code Playgroud)
[ 新更新2 ] [![截图] [ https://i.stack.imgur.com/mioth.jpg]]
[ 新更新3 ] [
] 1
我正在使用Retrofit2
fun create(): MyApiService {
return Retrofit.Builder()
.addCallAdapterFactory(RxJava2CallAdapterFactory.create())
.addConverterFactory(GsonConverterFactory.create())
.baseUrl(BASE_URL)
.build()
.create(MyApiService::class.java)
}
Run Code Online (Sandbox Code Playgroud)
隐式转换以下Json
[
{
"id": 1,
"name": "John",
}, {
"id": 2,
"name": "Mary",
}
]
Run Code Online (Sandbox Code Playgroud)
进入Kotlin数据类
object Model {
data class Person(val id: Int, val name: String)
}
Run Code Online (Sandbox Code Playgroud)
但是,我Expected BEGIN_OBJECT but was BEGIN_ARRAY在尝试时遇到错误
@GET("/people.json")
fun getPeople() : Observable<Model.Person>
Run Code Online (Sandbox Code Playgroud)
我已经尝试将Model对象更改为从List扩展(正如您通常在使用Java的Retrofit 1中所做的那样)或创建一个人员List字段,但无济于事.
我正在尝试使用.distinctUntilChanged()它,并且不会switchmap()在第一次之后传递价值。
RxTextView.textChanges(etUserQuery).debounce(300, TimeUnit.MILLISECONDS)
.observeOn(AndroidSchedulers.mainThread()).filter(charSequence -> {
if (charSequence.toString().isEmpty()) {
etUserQuery.setHint("Please type username");
return false;
} else
return true;
}).distinctUntilChanged()
.switchMap(charSequence -> vm.dataFromNetwork(charSequence.toString()))
.subscribe(fetchUserResponce -> {
noDataText.setVisibility(fetchUserResponce.getItems().size() == 0 ? View.VISIBLE : View.GONE);
mUsersListAdapter.updateData(fetchUserResponce.getItems());
}));
Run Code Online (Sandbox Code Playgroud)
是正确的使用地点.distinctUntilChanged()吗?
在RxJava2中,flatMap()和之间有什么区别flatMapIterable()?
背后的逻辑是flatMapIterable()什么?
io.reactivex.rxjava2:rxjava:2.1.13
kotlin_version = '1.2.30'
Run Code Online (Sandbox Code Playgroud)
我具有以下Observable,并且尝试引发异常以测试OnError中异常的捕获。但是,当我将以下内容传递给时,将onExceptionResumeNext(Observable.just(10))得到以下输出:
1
2
10
onComplete
fun main(args: Array<String>) {
Observable.fromArray(1, 2, 0, 4, 5, 6)
.doOnNext {
if (it == 0) {
throw RuntimeException("Exception on 0")
}
}
.onExceptionResumeNext(Observable.just(10))
.subscribe(
{
println(it)
},
{
println("onError ${it.message}")
},
{
println("onComplete")
} )
}
Run Code Online (Sandbox Code Playgroud)
但是,如果将lambda表达式传递给该方法,则会得到以下输出:
1
2
Observable.fromArray(1, 2, 0, 4, 5, 6)
.doOnNext {
if (it == 0) {
throw RuntimeException("Exception on 0")
}
}
.onExceptionResumeNext { Observable.just(10) }
.subscribe(
{
println(it)
},
{
println("onError …Run Code Online (Sandbox Code Playgroud) 我正在实现一个observable,在延迟5s后重试错误.我正在使用改造网络.我面临的问题是,当API返回错误时会有很多重试.我想仅在5秒后重试,但是重试以疯狂的速度发生(几乎是一秒钟的三次).知道为什么吗?
userAPI.getUsers()
.filter { it.users.isNotEmpty() }
.subscribeOn(Schedulers.io())
.retryWhen { errors -> errors.flatMap { errors.delay(5, TimeUnit.SECONDS) } }
.observeOn(AndroidSchedulers.mainThread())
.subscribe({}, {})
Run Code Online (Sandbox Code Playgroud)
其中userAPI.getUsers()返回一个可观察的.
疯狂的API请求数量:
08-13 12:31:31.308 26277-26453/com.app.user.dummy D/OkHttp: --> GET https://userapi.com/foo
08-13 12:31:31.825 26277-26453/com.app.user.dummy D/OkHttp: --> GET https://userapi.com/foo
08-13 12:31:32.370 26277-26453/com.app.user.dummy D/OkHttp: --> GET https://userapi.com/foo
08-13 12:31:32.897 26277-26453/com.app.user.dummy D/OkHttp: --> GET https://userapi.com/foo
08-13 12:31:33.436 26277-26453/com.app.user.dummy D/OkHttp: --> GET https://userapi.com/foo
08-13 12:31:33.952 26277-26453/com.app.user.dummy D/OkHttp: --> GET https://userapi.com/foo
08-13 12:31:34.477 26277-26453/com.app.user.dummy D/OkHttp: --> GET https://userapi.com/foo
08-13 12:31:35.020 26277-26453/com.app.user.dummy D/OkHttp: --> GET https://userapi.com/foo …Run Code Online (Sandbox Code Playgroud) 我有一个活动,每次用户输入更改时我都会向其发出网络请求.
api定义如下:
interface Api {
@GET("/accounts/check")
fun checkUsername(@Query("username") username: String): Observable<UsernameResponse>
}
Run Code Online (Sandbox Code Playgroud)
然后是管理它的服务:
class ApiService {
var api: Api
init {
api = retrofit.create(Api::class.java)
}
companion object {
val baseUrl: String = "https://someapihost"
var rxAdapter: RxJava2CallAdapterFactory = RxJava2CallAdapterFactory.create()
val retrofit: Retrofit = Retrofit.Builder()
.baseUrl(baseUrl)
.addConverterFactory(GsonConverterFactory.create())
.addCallAdapterFactory(rxAdapter)
.build()
}
fun checkUsername(username: String): Observable<UsernameResponse> {
return api.checkUsername(username)
}
}
Run Code Online (Sandbox Code Playgroud)
然后在我的活动中,每当EditText内容发生变化时,我都会进行此调用:
private fun checkUsername(username: String) {
cancelSubscription()
checkUsernameDisposable = ApiService()
.checkUsername(username)
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe {
updateUi(it)
}
}
Run Code Online (Sandbox Code Playgroud)
因此,每次输入变化时,这都会创建一个新的一次性用品.这显然是不正确的.我想要做的是使用新网络调用的结果更新现有订阅.
当生产者生产事件的速度快于客户消费时.
我想用可流动与onBackpressureLatest() ,我可以得到最新的情况下发出的.
但事实证明有一个128的默认缓冲区.我得到的是之前缓冲的过时事件.
那么如何才能获得最新的实际活动?
这是示例代码:
Flowable.interval(40, TimeUnit.MILLISECONDS)
.doOnNext{
println("doOnNext $it")
}
.onBackpressureLatest()
.observeOn(Schedulers.single())
.subscribe {
println("subscribe $it")
Thread.sleep(100)
}
Run Code Online (Sandbox Code Playgroud)
我的期望:
doOnNext 0
subscribe 0
doOnNext 1
doOnNext 2
subscribe 2
doOnNext 3
doOnNext 4
doOnNext 5
subscribe 5
doOnNext 6
doOnNext 7
subscribe 7
doOnNext 8
doOnNext 9
doOnNext 10
subscribe 10
...
Run Code Online (Sandbox Code Playgroud)
我得到了什么:
doOnNext 0
subscribe 0
doOnNext 1
doOnNext 2
subscribe 1
doOnNext 3
doOnNext 4
doOnNext 5
subscribe 2
doOnNext 6
doOnNext 7 …Run Code Online (Sandbox Code Playgroud) 我想在我的项目中将RxJava绑定API用于Android UI小部件。
因此,请遵循本网站' https://github.com/JakeWharton/RxBinding ' 的指导
但是我无法在Kotlin文件中导入任何Android UI小部件。 如果我在Java File中使用这些小部件,那么在哪里工作也很好。 因此,一直没有找到这个问题的实际。
作为参考,以下是在同一项目中使用的gradle文件和类文件(kotlin和Java)
build.gradle
dependencies {
implementation fileTree(dir: 'libs', include: ['*.jar'])
implementation"org.jetbrains.kotlin:kotlin-stdlib-jdk7:$kotlin_version"
implementation 'androidx.appcompat:appcompat:1.0.0-beta01'
implementation 'androidx.core:core-ktx:1.2.0-alpha01'
testImplementation 'junit:junit:4.12'
androidTestImplementation 'androidx.test:runner:1.1.0-alpha4'
androidTestImplementation 'androidx.test.espresso:espresso-core:3.1.0-alpha4'
implementation 'io.reactivex.rxjava2:rxjava:2.2.8'
implementation 'io.reactivex.rxjava2:rxandroid:2.1.1'
//RxBinding
implementation 'com.jakewharton.rxbinding3:rxbinding:3.0.0-alpha2'
}
Run Code Online (Sandbox Code Playgroud)
BindingExample.java类
RxBindingExample.kt类
已尝试在SO上探索此问题,但对于lib'com.jakewharton.rxbinding2:rxbinding'的先前版本可用问题或解决方案
rx-java2 ×10
android ×7
java ×5
rx-java ×5
kotlin ×3
retrofit2 ×2
backpressure ×1
list ×1
maven ×1
observable ×1
rx-android ×1
rx-binding ×1