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

RxJava2实现队列文件上传前逐项校验WIFI并阻塞等待

解决方案:为每个上传项前置WIFI连接检查

我完全理解你的需求——你希望每个待上传文件在启动上传操作前,都能确保当前处于WIFI连接状态;如果WIFI断开,就暂停队列处理,直到WIFI重新连接后再继续处理当前项,而不是像现有代码那样只在启动时检查一次WIFI。

核心思路是把「文件队列发射」作为上游数据源,然后对每个文件项单独做WIFI连接等待,确保只有在WIFI可用时才执行上传。这里用concatMap来保证队列的顺序处理(如果你的场景允许并行上传,可以换成flatMap,但上传队列通常需要顺序执行)。

第一步:封装复用的WIFI连接观察器

先把WIFI连接检查的逻辑抽成一个可复用的Observable,这样每个文件项都能复用这段逻辑:

private Observable<Connectivity> waitForWifiConnection() {
    return ReactiveNetwork.observeNetworkConnectivity(application)
            // 过滤出已连接状态
            .filter(ConnectivityPredicate.hasState(NetworkInfo.State.CONNECTED))
            // 过滤出WIFI类型
            .filter(ConnectivityPredicate.hasType(ConnectivityManager.TYPE_WIFI))
            .take(1); // 拿到一次有效连接就停止当前观察,避免持续发射
}

第二步:重构主流程,让每个文件项先等WIFI再上传

把queueItemAvailable()移到顶层作为上游,然后通过concatMap对每个文件项执行「等待WIFI → 上传」的逻辑:

queueItemAvailable()
    .subscribeOn(Schedulers.io()) // 让队列项的发射在IO线程执行
    .concatMap(fileItem -> 
        // 对每个文件,先等待WIFI连接,再执行上传
        waitForWifiConnection()
            .flatMap(connectivity -> audioSampleUploader.uploadFile(fileItem))
            .subscribeOn(Schedulers.single()) // 上传操作指定线程
    )
    .observeOn(AndroidSchedulers.mainThread()) // 结果回调切回主线程
    .subscribe(
        uploadResult -> {
            // 处理单个文件上传成功的逻辑
        },
        throwable -> {
            // 处理错误(比如上传失败、网络异常等)
        }
    );

关键逻辑解释

  1. queueItemAvailable()作为上游:现在它可以正常发射所有待上传文件项,每个项都会进入concatMap的处理流程。
  2. concatMap保证顺序:前一个文件的上传流程(包括等待WIFI、上传完成/失败)结束后,才会处理下一个文件,完美符合队列的顺序性要求。
  3. waitForWifiConnection()的阻塞效果:如果当前WIFI断开,这个Observable会一直等待,直到检测到WIFI重新连接,才会触发后续的上传操作。
  4. 线程调度:通过subscribeOn和observeOn合理分配线程,避免阻塞主线程。

可选优化:上传失败后的重试逻辑

如果希望在上传失败(比如中途WIFI断开)时自动重试,可以给上传逻辑加上retryWhen,结合WIFI检查实现智能重试:

.concatMap(fileItem -> 
    waitForWifiConnection()
        .flatMap(connectivity -> audioSampleUploader.uploadFile(fileItem))
        .retryWhen(errors -> errors.flatMap(error -> {
            // 判断是否为网络相关错误,是的话等待WIFI恢复后重试
            if (isNetworkRelatedError(error)) {
                return waitForWifiConnection();
            }
            // 非网络错误直接抛出,终止当前项的处理
            return Observable.error(error);
        }))
        .subscribeOn(Schedulers.single())
)

这里的isNetworkRelatedError()需要你根据实际的错误类型(比如IOException、自定义的网络异常等)来实现判断逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 08:21:51