如何在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
相关产品推荐
相关产品推荐

