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

如何在Observable内协程完成所有数据插入后停止前台服务?

分页API数据同步+前台服务停止问题解决

问题背景

在前台服务中通过Observable拉取分页API数据并写入数据库,同时展示带进度的通知。当前存在两个核心问题:

  • insertProducts内的stopService()无法稳定触发,savedItems == totalItems的判断存在竞态
  • 将stopService()放在Observable的onComplete()中时,会在所有API请求结束后立即执行,不等数据库插入完成

需求:确保所有Observable返回的数据都完成数据库插入后,再停止前台服务。

解决方案

核心思路是把数据库插入的协程操作纳入Observable的数据流生命周期,让Observable的onComplete()仅在所有数据插入完成后触发。

1. 改造数据库插入函数为挂起函数

将原来的协程启动逻辑改为挂起函数,确保插入操作全部完成后才返回,同时修正通知更新的线程(必须在主线程操作NotificationManager):

private suspend fun insertProducts(totalItems: Int?, products: List<ProdottoBarcode>?) {
    products ?: return
    for (product in products) {
        repository.insert(product)
        savedItems += 1

        val notification = totalItems?.let { items ->
            NotificationCompat.Builder(baseContext, "progress_channel")
                .setSmallIcon(R.drawable.ic_box)
                .setContentTitle("Sincronizzati: $savedItems prodotti su $totalItems")
                .setProgress(items, savedItems, false)
                .setOngoing(true)
                .build()
        }

        // 切换到主线程更新通知
        withContext(Dispatchers.Main) {
            notification?.let { notificationManager.notify(notificationId, it) }
        }
    }
}

2. 整合协程操作到Observable数据流

使用concatMapSingle将数据库插入的挂起函数转换成Observable的一部分,确保每一页数据的插入完成后,才会处理下一页的API请求:

override fun onCreate() {
    super.onCreate()
    ...
    getAllProducts()
        .subscribeOn(Schedulers.io()) // 网络请求放在IO线程
        .concatMapSingle { response ->
            // 用Single.fromCoroutine将挂起函数包装为Observable流的一部分
            Single.fromCoroutine {
                if (response.isSuccessful) {
                    val products = response.body()
                    val totalItems = response.headers().get("items")?.toInt()
                    insertProducts(totalItems, products)
                }
            }
        }
        .subscribeWith(object : DisposableObserver<Unit>() {
            override fun onNext(unit: Unit) {
                // 无需额外操作,插入逻辑已在concatMapSingle中完成
            }

            override fun onError(e: Throwable) {
                stopService()
                // 可添加错误状态的通知更新
            }

            override fun onComplete() {
                // 所有API请求+数据库插入全部完成,停止服务并更新通知
                val completedNotification = NotificationCompat.Builder(baseContext, "progress_channel")
                    .setSmallIcon(R.drawable.ic_box)
                    .setContentTitle("Sincronizzazione completata!")
                    .setOngoing(false)
                    .build()
                notificationManager.notify(notificationId, completedNotification)
                stopService()
            }
        })
}

3. 优化分页API的线程调度(可选)

确保网络请求在IO线程执行,避免阻塞主线程:

private fun getAllProducts(): Observable<Response<List<ProdottoBarcode>>> {
    val lastId = intArrayOf(0)

    return Observable.range(1, Integer.MAX_VALUE - 1)
        .subscribeOn(Schedulers.io())
        .concatMap { currentPage -> getProducts(currentPage, lastId[0]) }
        .takeUntil { response -> response.body()?.isEmpty() == true }
        .doOnNext { response ->
            lastId[0] = response.headers().get("lastId")?.toInt()!!
        }
}

原理说明

  • concatMapSingle会按顺序处理每一页的API响应,等待当前页的数据库插入操作完成后,才会订阅下一页的API请求
  • Observable的onComplete()会在所有页的API请求和数据库插入都完成后触发,此时调用stopService()能确保所有数据都已写入数据库
  • 通知更新切换到主线程,符合Android的线程规则,避免ANR或异常

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 19:35:38