基于Retrofit与RxJava2的所有HTTP请求间隔调度问题
问题分析与解决方案
你的核心问题在于:两个独立的 Observable 请求链是并行运行的,哪怕你给每个链都用了单线程调度器,也只能保证单个链内部的操作串行,但两个链之间的请求依然会同时触发,导致机器人服务器接收到重叠请求。
举个例子:连接检查的 Observable 每2秒发一次请求,传感器数据的 Observable 每 refreshRate 毫秒发一次请求,它们的订阅是独立的,各自的调度器不会互相约束,所以两个请求会同时打向服务器。
正确的解决思路:全局请求串行化 + 固定间隔调度
你需要一个全局的请求管理器,把所有要发送的请求都放进一个队列,然后按 100-150ms 的固定间隔依次执行,确保同一时间只有一个请求发向机器人服务器。
下面是两种可行的实现方案:
方案一:基于 Subject 的请求队列调度器
创建一个全局的调度器类,用 PublishSubject 作为请求队列,结合 Observable.interval 定时取出请求执行:
public class RobotRequestScheduler { // 序列化的 Subject 保证线程安全 private final Subject<Callable<?>> requestQueue = PublishSubject.<Callable<?>>create().toSerialized(); private final Disposable schedulerDisposable; public RobotRequestScheduler(long intervalMs) { // 按固定间隔从队列取请求执行 schedulerDisposable = Observable.interval(intervalMs, TimeUnit.MILLISECONDS) .withLatestFrom(requestQueue, (tick, request) -> request) .observeOn(Schedulers.io()) .subscribe(request -> { try { // 执行请求 request.call(); } catch (Exception e) { e.printStackTrace(); } }); } // 提交请求到队列 public void submitRequest(Callable<?> request) { requestQueue.onNext(request); } // 销毁调度器 public void dispose() { schedulerDisposable.dispose(); requestQueue.onComplete(); } }
然后在你的 Presenter 中使用这个调度器:
1. 初始化调度器(比如在 BasePresenter 中)
// 设置150ms的请求间隔,符合服务器处理能力 RobotRequestScheduler requestScheduler = new RobotRequestScheduler(150);
2. 改造连接检查请求
// 每2秒提交一次连接检查请求到队列 Observable.interval(2000, TimeUnit.MILLISECONDS) .subscribeOn(Schedulers.io()) .subscribe(tick -> { requestScheduler.submitRequest(() -> { requestInterface.checkConnectionToRobot() .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .doOnNext(response -> action(true)) .doOnError(err -> { err.printStackTrace(); action(false); }) .subscribe(new DisposingObserver<ResponseBody>() { @Override public void onSubscribe(Disposable d) { addContinuous(d); } }); return null; }); });
3. 改造传感器数据请求
// 按refreshRate间隔提交传感器数据请求到队列 Observable.interval(refreshRate, TimeUnit.MILLISECONDS) .subscribeOn(Schedulers.io()) .subscribe(tick -> { requestScheduler.submitRequest(() -> { requestInterface.refreshBaseData() .subscribeOn(Schedulers.io()) .map(RefreshDataConverter::convertData) .observeOn(AndroidSchedulers.mainThread()) .doOnNext(data -> getViewState().setBaseData(data)) .retryWhen(errors -> errors.delay(1000, TimeUnit.MILLISECONDS)) .subscribe(new DisposingObserver<RefreshData>() { @Override public void onSubscribe(Disposable d) { addContinuous(d); baseCheckDisposable = d; } }); return null; }); });
方案二:合并请求流 + concatMap 串行执行
如果你不想单独写调度器类,可以把两个请求流合并,用 concatMap 保证串行执行,同时添加固定间隔:
// 定义连接检查的任务流 Observable<Runnable> connectionCheckTasks = Observable.interval(2000, TimeUnit.MILLISECONDS) .map(tick -> (Runnable) () -> { requestInterface.checkConnectionToRobot() .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .doOnNext(response -> action(true)) .doOnError(err -> { err.printStackTrace(); action(false); }) .subscribe(new DisposingObserver<ResponseBody>() { @Override public void onSubscribe(Disposable d) { addContinuous(d); } }); }); // 定义传感器数据的任务流 Observable<Runnable> sensorDataTasks = Observable.interval(refreshRate, TimeUnit.MILLISECONDS) .map(tick -> (Runnable) () -> { requestInterface.refreshBaseData() .subscribeOn(Schedulers.io()) .map(RefreshDataConverter::convertData) .observeOn(AndroidSchedulers.mainThread()) .doOnNext(data -> getViewState().setBaseData(data)) .retryWhen(errors -> errors.delay(1000, TimeUnit.MILLISECONDS)) .subscribe(new DisposingObserver<RefreshData>() { @Override public void onSubscribe(Disposable d) { addContinuous(d); baseCheckDisposable = d; } }); }); // 合并两个任务流,按固定间隔串行执行 Observable.merge(connectionCheckTasks, sensorDataTasks) .concatMap(task -> Observable.just(task) .delay(150, TimeUnit.MILLISECONDS)) // 控制请求间隔 .observeOn(Schedulers.io()) .subscribe(task -> { try { task.run(); } catch (Exception e) { e.printStackTrace(); } });
关键注意点
- 错误处理隔离:把每个请求的
retryWhen放在请求内部,不要放到全局流中,避免一个请求失败影响所有后续请求。 - 线程安全:如果多个线程提交请求,一定要用
toSerialized()包装 Subject,保证队列线程安全。 - 资源清理:在 Presenter 销毁时,记得调用调度器的
dispose()方法,避免内存泄漏。
内容的提问来源于stack exchange,提问作者Sprite_______
相关产品推荐
相关产品推荐

