You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RxJava Observable链触发onErrorResumeNext后停止执行问题求助

解决Observable超时后链停止执行的问题

Hey there, let's break down why your Observable chain goes silent when the timeout kicks in, and how to fix it quickly.

问题根源

When your timeout triggers onErrorResumeNext, you're returning Observable.empty() — here's the catch: empty() only sends an onCompleted signal to downstream operators, no onNext events at all. Since operators like map and flatMap only run when they receive an onNext from upstream, they just sit idle once empty() completes. That's why your chain stops with no errors or further action.

解决方案

Instead of returning empty(), you need to return an Observable that emits at least one value so downstream operators have something to process. Here are two practical options:

1. 发射默认空值(快速适配现有代码)

If you want to keep your existing map logic that returns "Ok", replace Observable.empty() with Observable.just(null). This sends a null through the chain, triggering your map and subsequent operators:

public Observable<Void> ApiCallBackObservable() { 
    return Observable.fromEmitter(new Action1<AsyncEmitter<Void>>() { 
        @Override 
        public void call(final AsyncEmitter<Void> voidAsyncEmitter) { 
            APICall.setCallBack(new CallBack() { 
                @Override 
                public void onResponseReceived() { 
                    voidAsyncEmitter.onNext(null); 
                    voidAsyncEmitter.onCompleted(); 
                } 
            });
            // 额外优化:添加回调清理,避免内存泄漏
            voidAsyncEmitter.setCancellation(new Action0() {
                @Override
                public void call() {
                    APICall.removeCallBack(); // 假设APICall提供移除回调的方法
                }
            });
        } 
    }, AsyncEmitter.BackpressureMode.NONE)
    .timeout(100, TimeUnit.MILLISECONDS)
    .onErrorResumeNext(new Func1<Throwable, Observable<? extends Void>>() { 
        @Override 
        public Observable<? extends Void> call(Throwable throwable) { 
            // 替换empty为发射null,让下游能收到事件继续执行
            return Observable.just(null); 
        } 
    }); 
}

Now when a timeout happens, the map will receive null and return "Ok", and your flatMap operations will run as expected.

2. 发射标记值(区分正常响应和超时)

If you want downstream operators to explicitly handle timeouts (instead of treating them like normal success), create a simple marker class to carry state:

// 定义一个标记类,区分正常响应和超时
static class ApiResult {
    boolean isTimeout;
    // 可按需添加其他字段,比如响应数据
}

// 修改Observable返回类型为Observable<ApiResult>
public Observable<ApiResult> ApiCallBackObservable() { 
    return Observable.fromEmitter(new Action1<AsyncEmitter<ApiResult>>() { 
        @Override 
        public void call(final AsyncEmitter<ApiResult> emitter) { 
            APICall.setCallBack(new CallBack() { 
                @Override 
                public void onResponseReceived() { 
                    ApiResult successResult = new ApiResult();
                    successResult.isTimeout = false;
                    emitter.onNext(successResult); 
                    emitter.onCompleted(); 
                } 
            });
            // 清理回调避免泄漏
            emitter.setCancellation(new Action0() {
                @Override
                public void call() {
                    APICall.removeCallBack();
                }
            });
        } 
    }, AsyncEmitter.BackpressureMode.NONE)
    .timeout(100, TimeUnit.MILLISECONDS)
    .onErrorResumeNext(new Func1<Throwable, Observable<? extends ApiResult>>() { 
        @Override 
        public Observable<? extends ApiResult> call(Throwable throwable) { 
            ApiResult timeoutResult = new ApiResult();
            timeoutResult.isTimeout = true;
            return Observable.just(timeoutResult); 
        } 
    }); 
}

Then handle the timeout case explicitly in your downstream chain:

apiCallBackObservable()
    .map(new Func1<ApiResult, String>() { 
        @Override 
        public String call(ApiResult result) { 
            return result.isTimeout ? "Timeout occurred" : "Ok"; 
        } 
    })
    .flatMap(...)
    // 后续操作

内容的提问来源于stack exchange,提问作者Katazo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 06:43:34