标签: rx-java2

从RxJava编译示例时找不到org.reactivestreams.Publisher的类文件?

以下代码

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)

java maven rx-java rx-java2

1
推荐指数
2
解决办法
5251
查看次数

如何取消订阅rxJava请求

我有以下类,我想返回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 ] [Screenshot2] 1

java android rx-java2

1
推荐指数
1
解决办法
2452
查看次数

JsonArray使用Retrofit到Kotlin数据类(预期BEGIN_OBJECT但是BEGIN_ARRAY)

我正在使用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字段,但无济于事.

android list kotlin retrofit2 rx-java2

1
推荐指数
1
解决办法
2019
查看次数

第一次后Rxjava2distantUntilChanged()无法正常工作

我正在尝试使用.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()吗?

java android rx-java rx-java2

1
推荐指数
1
解决办法
1404
查看次数

RxJava2 flatMap和flatMapIterable

RxJava2中,flatMap()和之间有什么区别flatMapIterable()

背后的逻辑是flatMapIterable()什么?

java android rx-java2

1
推荐指数
1
解决办法
1365
查看次数

与onExceptionResumeNext的混淆传递了一个可观察到的lambda表达式

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)

kotlin rx-java2

1
推荐指数
1
解决办法
188
查看次数

使用retryWhen()重试observable

我正在实现一个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)

java android observable rx-java rx-java2

1
推荐指数
1
解决办法
244
查看次数

RxJava2如何在请求参数更改时更新现有订阅

我有一个活动,每次用户输入更改时我都会向其发出网络请求.

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)

因此,每次输入变化时,这都会创建一个新的一次性用品.这显然是不正确的.我想要做的是使用新网络调用的结果更新现有订阅.

android rx-java retrofit2 rx-java2

1
推荐指数
1
解决办法
220
查看次数

如何在Rxjava2中获得有关背压的实际最新事件?Flowable.onBackpressureLatest()未按预期工作

当生产者生产事件的速度快于客户消费时.

我想用可流动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)

backpressure rx-java rx-java2

1
推荐指数
1
解决办法
45
查看次数

使用最新的com.jakewharton.rxbinding3:rxbinding:3.0.0-alpha2库时找不到RxTextView和其他小部件

我想在我的项目中将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'的先前版本可用问题或解决方案

android kotlin rx-android rx-binding rx-java2

1
推荐指数
1
解决办法
497
查看次数