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

基于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();
            }
        });

关键注意点

  1. 错误处理隔离:把每个请求的 retryWhen 放在请求内部,不要放到全局流中,避免一个请求失败影响所有后续请求。
  2. 线程安全:如果多个线程提交请求,一定要用 toSerialized() 包装 Subject,保证队列线程安全。
  3. 资源清理:在 Presenter 销毁时,记得调用调度器的 dispose() 方法,避免内存泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 04:10:30