如何用Kotlin协程在Realtime Database中通过事务实现计数器递增?
将Firebase Realtime Database事务与Kotlin协程+Flow结合实现
我原本用以下代码通过Kotlin协程+Flow实现Realtime Database中quantity字段的递增,运行正常:
override fun incrementQuantity() = flow { try { heroIdRef.update("quantity", FieldValue.increment(1)).await() emit(Result.Success(true)) } catch (e: Exception) { emit(Result.Failure(e)) } }
但现在需要先判断quantity的值再操作:当quantity为1时将其设为null,否则递增。这种场景必须用事务,我已经写出了可运行的事务代码,但不知道怎么把它改成和上面一样的协程+Flow形式,支持await()并返回Flow。我的事务代码如下:
override fun incrementQuantity() { val transaction = object : Transaction.Handler { override fun doTransaction(mutableData: MutableData): Transaction.Result { val quantity = mutableData.getValue(Long::class.java) ?: return Transaction.success(mutableData) if (quantity == 1L) { mutableData.value = null } else { mutableData.value = quantity + 1 } return Transaction.success(mutableData) } override fun onComplete(error: DatabaseError?, committed: Boolean, data: DataSnapshot?) { throw error.toException() } } heroIdRef.runTransaction(transaction) }
解决方案
要把事务和协程、Flow结合,核心是用suspendCoroutine将Firebase的回调式API转换为挂起函数,再包装进Flow中。具体实现如下:
1. 封装事务为挂起函数
先写一个挂起函数,处理事务的执行和回调结果:
private suspend fun runTransactionSuspended(): Boolean { return suspendCoroutine { continuation -> val transaction = object : Transaction.Handler { override fun doTransaction(mutableData: MutableData): Transaction.Result { val quantity = mutableData.getValue(Long::class.java) ?: return Transaction.success(mutableData) if (quantity == 1L) { mutableData.value = null } else { mutableData.value = quantity + 1 } return Transaction.success(mutableData) } override fun onComplete(error: DatabaseError?, committed: Boolean, data: DataSnapshot?) { error?.let { continuation.resumeWithException(it.toException()) } ?: run { continuation.resume(committed) } } } heroIdRef.runTransaction(transaction) } }
这里用suspendCoroutine捕获事务的回调结果:如果有错误就抛出异常,否则返回事务是否提交成功(committed)。
2. 包装成Flow并返回
接下来就可以像第一个示例那样,把挂起函数包装进Flow中处理异常和结果:
override fun incrementQuantity() = flow { try { val committed = runTransactionSuspended() emit(Result.Success(committed)) } catch (e: Exception) { emit(Result.Failure(e)) } }
关键说明
suspendCoroutine的作用是将基于回调的异步操作转换为挂起函数,让代码保持协程的线性写法。- 事务完成后,
onComplete中的committed参数表示事务是否成功提交(即使事务因数据冲突重试后最终成功,该值也为true)。 - 异常会被捕获并包装成
Result.Failure,和你原来的Flow逻辑完全一致。
内容的提问来源于stack exchange,提问作者Always Learner
相关产品推荐
相关产品推荐

