我知道当所有项目都发出时,会调用观察者的 OnComplete 。在下面的代码中,我将数据从游标放入 flatMap 运算符中的 ArrayList。我的光标有 100 个条目( c.getCount() 给出 100 ),我的列表大小是 100。 onNext 也被调用 100 次。但 onComplete 没有被调用。我正在 onComplete 中填充列表视图。
static int i = 0;
final List<String> ar = new ArrayList<>();
ListView lv = ...;
ArrayAdapter<String> adapter = ...;
.
.
.
q.subscribeOn(Schedulers.io()).observeOn(AndroidSchedulers.mainThread())
.flatMap(new Func1<SqlBrite.Query, Observable<String>>() {
@Override
public Observable<String> call(SqlBrite.Query query) {
Cursor c = query.run();
c.moveToFirst();
Log.d("testApp", String.valueOf(c.getCount())); // prints 100
do {
ar.add(c.getString(0));
} while (c.moveToNext());
Log.d("testApp", String.valueOf(ar.size())); // prints 100
return Observable.from(ar);
}
}).subscribe(new Observer<String>() …Run Code Online (Sandbox Code Playgroud) 这是我第一次在反应范式世界中开发,我开始使用rxjava2/rxandroid2,基于我观看过的视频和我读过的文章,从2开始看起来更好,因为有很多变化,图书馆的大规模存在差异,但现在我在寻找像这样的东西时遇到了一些麻烦
unsubscribe()
Run Code Online (Sandbox Code Playgroud)
以前的rxjava/rxandroid库的方法
我的目标很简单
,基于周围的资源
dispose()
Run Code Online (Sandbox Code Playgroud)
rx2的方法,我对此的理解是它处理任何当前资源(在我的情况下,基于我理解的,调用它将使可观察的自身分离给任何观察者).
但这似乎不是我所期待的,请看一下ff代码:
public class MainActivity extends AppCompatActivity {
final Disposable disposable = new Disposable() {
@Override
public void dispose() {
Log.e("Disposed", "_ dispose called.");
}
@Override
public boolean isDisposed() {
return true;
}
};
@Override
protected void onCreate(Bundle savedInstanceState) {
super.onCreate(savedInstanceState);
setContentView(R.layout.activity_main);
Observer<Object> observer = new Observer<Object>() {
@Override
public void onSubscribe(Disposable d) {
Log.e("OnSubscribe", "On Subscribed Called");
}
@Override
public void onNext(Object value) {
Log.e("onNext", "Actual Value (On Next …Run Code Online (Sandbox Code Playgroud) 我有一个返回 a 的方法,Single<List<Item>>我想获取此列表中的每个项目并将其向下游传递给返回Completable. 我想等到每个项目成功完成并返回Completable结果。我最初的做法是分别处理每个项目使用flatMapIterable和组合使用的结果toList,但我不能把toList一个上Completable对象。有没有其他方法可以以这种方式将许多Completable任务“聚合”为一个Completable?这是我到目前为止所拥有的:
public Single<List<Item>> getListOfItems() {
...
}
public Completable doSomething(Item item) {
...
}
public Completable processItems() {
return getListOfItems()
.toObservable()
.flatMapIterable(items -> items)
.flatMapCompletable(item -> doSomething(item))
.toList() // ERROR: No method .toList() for Completable
.ignoreElements();
}
Run Code Online (Sandbox Code Playgroud) 我试图找到一种方法来并行执行请求并在每个可观察对象完成时处理它们。尽管当所有 observables 都给出响应时一切都在工作,但我没有看到在一切都完成后处理所有错误的方法。
这是 zip 运算符的示例,它基本上并行执行 2 个请求:
Observable.zip(
getObservable1()
.onErrorResumeNext { errorThrowable: Throwable ->
Observable.error(ErrorEntity(Type.ONE, errorThrowable))
}.subscribeOn(Schedulers.io()),
getObservable2()
.onErrorResumeNext { errorThrowable: Throwable ->
Observable.error(ErrorEntity(Type.TWO, errorThrowable))
}.subscribeOn(Schedulers.io()),
BiFunction { value1: String, value2: String ->
return@BiFunction value1 + value2
})
//execute requests should be on io() thread
.subscribeOn(Schedulers.io())
//there are other tasks inside subscriber that need io() thread
.observeOn(AndroidSchedulers.mainThread())
.subscribe(
{ result ->
Snackbar.make(view, "Replace with your own action " + result, Snackbar.LENGTH_LONG)
.setAction("Action", null).show()
},
{ error ->
Log.d("TAG", "Error …Run Code Online (Sandbox Code Playgroud) 我正在尝试使用Rx方式使用房间从数据库中检索数据.这就是我试图这样做的方式
override fun onStart() {
super.onStart()
disposable.add(presenter.getAllBooks()
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe({
println(it.size())
}))
}
Run Code Online (Sandbox Code Playgroud)
这是getAllBooks()演示者内部的方法
fun getAllBooks() : Flowable<List<Book>> {
val isMainThread = Looper.myLooper() == Looper.getMainLooper()
if (!isMainThread) {
updateBooks()
return db.bookDao().allBooks
}
return Flowable.empty()
}
Run Code Online (Sandbox Code Playgroud)
这里isMainThread变量总是true,我也试过observeOn(Shcedulers.io()),但同样的问题.
我正在使用 rxjava 2 并尝试使用 rxbus 传递值
接收总线代码
public class SeasonTabSelectorBus {
private static SeasonTabSelectorBus instance;
private PublishSubject<Object> subject = PublishSubject.create();
public static SeasonTabSelectorBus instanceOf() {
if (instance == null) {
instance = new SeasonTabSelectorBus();
}
return instance;
}
public void setTab(Object object) {
try {
subject.onNext(object);
subject.onComplete();
} catch (Exception e) {
e.printStackTrace();
}
}
public Observable<Object> getSelectedTab() {
return subject;
}
}
Run Code Online (Sandbox Code Playgroud)
我将值设置为
SeasonTabSelectorBus.instanceOf().setTab(20);
Run Code Online (Sandbox Code Playgroud)
这是我订阅的代码
SeasonTabSelectorBus.instanceOf().getSelectedTab().subscribe(new Observer<Object>(){
@Override
public void onSubscribe(Disposable d) {
}
@Override
public void onNext(Object o) …Run Code Online (Sandbox Code Playgroud) 实现 'io.reactivex.rxjava3:rxandroid:3.0.0' 实现 'io.reactivex.rxjava3:rxjava:3.0.0'
val TAG:String = RXKotlinDemoClass::class.java.simpleName
override fun onCreate(savedInstanceState: Bundle?) {
super.onCreate(savedInstanceState)
setContentView(R.layout.activity_main)
var observable = Observable.just("Goat","Dog","Cow")
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread()).subscribe({
value -> println(TAG+"$value")
},{
error -> println(TAG+"$error")
},{
println(TAG+"onComplete")
}
)
Run Code Online (Sandbox Code Playgroud)
}
异常:java.lang.NoSuchMethodError:没有静态方法元工厂(Ljava/lang/invoke/MethodHandles$Lookup;Ljava/lang/String;Ljava/lang/invoke/MethodType;Ljava/lang/invoke/MethodType;Ljava/lang/invoke /MethodHandle;Ljava/lang/invoke/MethodType;)Ljava/lang/invoke/CallSite; 在类 Ljava/lang/invoke/LambdaMetafactory 中;或其超类(“java.lang.invoke.LambdaMetafactory”的声明出现在 /apex/com.android.runtime/javalib/core-oj.jar 中)位于 io.reactivex.rxjava3.android.schedulers.AndroidSchedulers.(AndroidSchedulers .java:33) 在 io.reactivex.rxjava3.android.schedulers.AndroidSchedulers.mainThread(AndroidSchedulers.java:44) 在 com.android.myfirstapp.RXKotlinDemoClass.onCreate(RXKotlinDemoClass.kt:19) 在 android.app.Activity。在 android.app.Activity.performCreate(Activity.java:7791) 处执行Create(Activity.java:7802) 在 android.app.Instrumentation.callActivityOnCreate(Instrumentation.java:1299) 在 android.app.ActivityThread.performLaunchActivity(ActivityThread.java :3245)在android.app.ActivityThread.handleLaunchActivity(ActivityThread.java:3409)在android.app.servertransaction.LaunchActivityItem.execute(LaunchActivityItem.java:83)在android.app.servertransaction.TransactionExecutor.executeCallbacks(TransactionExecutor.java: 135) 在 android.app.servertransaction.TransactionExecutor.execute(TransactionExecutor.java:95) 在 android.app.ActivityThread$H.handleMessage(ActivityThread.java:2016) 在 android.os.Handler.dispatchMessage(Handler.java:107) )在 android.os.Looper.loop(Looper.java:214) 在 android.app.ActivityThread.main(ActivityThread.java:7356) 在 java.lang.reflect.Method.invoke(Native Method) 在 com.android。 Internal.os.RuntimeInit$MethodAndArgsCaller.run(RuntimeInit.java:492) 在 com.android.internal.os.ZygoteInit.main(ZygoteInit.java:930)
我在Android上使用RxJava来做一些事情,
在使用之前我总是在observable上做同样的事情:
Observable<AnyObject> observable = getSomeObservable();
// The next 2 lines are the lines that i always add them to any Observable
observable.observeOn(AndroidSchedulers.mainThread())
.subscribeOn(Schedulers.computation());
Run Code Online (Sandbox Code Playgroud)
因此,Observable是通用的,可以是任何对象,如果我想在它上面添加这两行并在Statis方法中返回它,我需要使该方法也是Generic
我试图做的是通过参数传递observable,添加设置并将其返回如下:
public class UtilsObservable<T> {
public static Observable<T> setupObservable(Observable<T> observable) {
return observable.observeOn(AndroidSchedulers.mainThread())
.subscribeOn(Schedulers.computation());
}
Run Code Online (Sandbox Code Playgroud)
我在这里遇到编译错误说:
UtilsObservable.this cannot be referenced from a static context
Run Code Online (Sandbox Code Playgroud)
我的问题是:
那么这可以做到吗?通用方法采用泛型对象修改它并返回相同的类型?
为什么RxJava 1.x flatMap()运算符是通过merge实现的?
public final <R> Observable<R> flatMap(Func1<? super T, ? extends Observable<? extends R>> func) {
if (getClass() == ScalarSynchronousObservable.class) {
return ((ScalarSynchronousObservable<T>)this).scalarFlatMap(func);
}
return merge(map(func));
}
Run Code Online (Sandbox Code Playgroud)
从flatMap()调用中,我只能返回一个符合的Observable <? extends Observable<? extends R>>。比map(func)调用将它包装到另一个Observable中,这样我们就得到了这样的东西Observable<? extends Observable<? extends R>>。这使我认为map(func)之后的merge()调用是不必要的。
据说merge()运算符执行以下操作:
将可发射Observable的Observable展平为单个Observable,该Observable发射那些Observable发射的项目,而无需进行任何转换。
现在,在平面地图内,我们只能有一个Observable发出一个Observable。为什么要合并?我在这里想念什么?
谢谢。
所以我想用rx-java2进行表单验证.我正在使用Kotlin.我遇到了两个问题.emailObservable和passwordObservable都是类型Disposable!.我尝试通过调用指定类型,val emailObservable: Observable<Boolean>但Android Studio认为它Disposable!.
其次,当我想使用方法时combineLatest出现错误:使用提供的参数不能调用以下任何函数.
emailObservable和passwordObservable都能正常工作.我是rx-java的新手,我对这种类型的东西很困惑.
val emailObservable = RxTextView.afterTextChangeEvents(textEmail)
.observeOn(AndroidSchedulers.mainThread())
.map { x -> textEmail.text.length > 3 }
.subscribe { x -> foo(x) }
val passwordObservable =RxTextView.afterTextChangeEvents(textPassword)
.observeOn(AndroidSchedulers.mainThread())
.map { x -> textPassword.text.length > 5 }
.subscribe { x -> foo(x) }
Observable.combineLatest(emailObservable,
passwordObservable,
BiFunction { x: Boolean, y:Boolean -> x && y })
Run Code Online (Sandbox Code Playgroud) rx-android ×10
rx-java ×7
rx-java2 ×6
android ×5
kotlin ×3
java ×2
android-room ×1
dispose ×1
flatmap ×1
generics ×1
reactivex ×1
rx-binding ×1
unsubscribe ×1