Android中Couchbase Lite协程队列读写同步及结果返回问题
问题解答
1. 当前方案并非最优解,更优实现思路
当前方案的核心缺陷是未真正实现数据库操作的串行化:你通过Channel存放Job并逐个join,但这些Job是在多线程的IO协程池中运行的,加上CoroutineStart.LAZY的设置,读操作的await()会直接触发Job启动,大概率和未处理完的写操作并发执行,这就是读操作提前完成返回null的根源。
推荐两种更可靠的实现方式:
方案一:单线程Dispatcher串行处理(最简洁)
直接创建专属单线程Dispatcher,所有数据库操作都在这个线程上执行,天然保证串行:
object DatabaseManager { // 单线程线程池转协程Dispatcher,命名线程便于调试 private val dbDispatcher = Executors.newSingleThreadExecutor { Thread(it, "DB-Worker").apply { isDaemon = true } }.asCoroutineDispatcher() private val scope = CoroutineScope(dbDispatcher) // 提交无返回值操作(如写入) fun submit(block: suspend () -> Unit) { scope.launch { block() } } // 提交带返回值操作(如读取),挂起函数不阻塞线程 suspend fun <T> submitWithResult(block: suspend () -> T): T = withContext(dbDispatcher) { block() } // 销毁资源,避免内存泄漏 fun shutdown() { scope.cancel() dbDispatcher.close() } }
使用示例:
// 扩展Database的写入方法 fun Database.saveDocumentQueued(document: MutableDocument) { DatabaseManager.submit { Timber.tag("quechk").d("saving :: ${document.id}") this@saveDocumentQueued.save(document) } } // 扩展Database的读取方法(挂起函数) suspend fun Database.getDocumentQueued(id: String): Document? = DatabaseManager.submitWithResult { Timber.tag("quechk").d("getting :: $id") this@getDocumentQueued.getDocument(id) }
方案二:Channel串行分发任务(轻量队列)
如果坚持用Channel实现,不要存放Job,而是直接存放任务逻辑,在同一个协程中逐个执行,确保串行:
object DatabaseQueue { private val scope = CoroutineScope(Dispatchers.IO) // 无返回值任务通道 private val taskChannel = Channel<suspend () -> Unit>(Channel.UNLIMITED) // 带返回值任务通道,封装任务和结果回调 private val resultTaskChannel = Channel<Pair<suspend () -> Any?, CompletableDeferred<Any?>>>(Channel.UNLIMITED) init { // 串行处理无返回值任务 scope.launch { for (task in taskChannel) task() } // 串行处理带返回值任务 scope.launch { for ((task, deferred) in resultTaskChannel) { try { deferred.complete(task()) } catch (e: Exception) { deferred.completeExceptionally(e) } } } } fun submit(block: suspend () -> Unit) { taskChannel.trySendBlocking(block) } fun <T> submitAsync(block: suspend () -> T): Deferred<T> { val deferred = CompletableDeferred<T>() resultTaskChannel.trySendBlocking(Pair(block as suspend () -> Any?, deferred as CompletableDeferred<Any?>)) return deferred } fun cancel() { taskChannel.cancel() resultTaskChannel.cancel() scope.cancel() } }
使用示例:
fun Database.saveDocumentQueued(document: MutableDocument) { DatabaseQueue.submit { Timber.tag("quechk").d("saving :: ${document.id}") this@saveDocumentQueued.save(document) } } suspend fun Database.getDocumentQueued(id: String): Document? = DatabaseQueue.submitAsync { Timber.tag("quechk").d("getting :: $id") this@getDocumentQueued.getDocument(id) }
2. 避免读操作返回null的关键措施
要彻底解决这个问题,核心是让读操作严格在对应写操作完成后执行:
- 所有数据库操作(读写)必须统一走串行队列/单线程Dispatcher,绝对不能有操作绕开机制直接访问数据库
- 读操作必须用挂起函数(
suspend)实现,禁止用runBlocking阻塞线程,确保协程能正确等待队列中的前置任务完成 - 如果要修复原方案,必须保证Job的执行串行化:在队列消费协程中先启动Job再等待完成,同时保留
CoroutineStart.LAZY避免提前执行(但这种方式仍不如单线程方案可靠)
原方案修复示例:
object DatabaseQueue { private val scope = CoroutineScope(IOCoroutineScope) private val queue = Channel<Job>(Channel.UNLIMITED) init { scope.launch(Dispatchers.Default) { for (job in queue) { job.start() // 先启动任务 job.join() // 等待任务完成 } } } fun submit( context: CoroutineContext = EmptyCoroutineContext, block: suspend CoroutineScope.() -> Unit ) { // 用LAZY创建Job,避免提前启动 val job = scope.launch(context, CoroutineStart.LAZY, block) queue.trySendBlocking(job) } // submitAsync同理修改,保留LAZY并在队列中启动 fun submitAsync( context: CoroutineContext = EmptyCoroutineContext, id: String, database: Database ): Deferred<Document?> { val job = scope.async(context, CoroutineStart.LAZY) { database.getDocument(id) } queue.trySendBlocking(job) return job } }
内容的提问来源于stack exchange,提问作者Mohammedsaif Kordia
相关产品推荐
相关产品推荐

