Android Kotlin中Room数据库协程同步问题:Socket文件分块下载的并行处理与竞态条件解决
解决Room数据库与协程的并行处理及竞态条件问题
我来帮你拆解这两个问题,一步步搞定:
1. 实现文件分块的并行处理
你当前的messageResolver里每次收到FilePart消息时,都会通过scope.launch(Dispatchers.IO)启动新协程,但实际没达到并行效果——大概率是因为你的Socket消息接收逻辑是串行的(比如在同一个线程里依次处理每个消息),导致messageResolver被串行调用,协程只能排队启动。
要真正实现并行处理,你需要把消息接收和业务处理解耦,用Channel做消息缓冲池,再启动多个消费协程并行处理分块:
修改MessageHandlerImpl的实现:
class MessageHandlerImpl : MessageHandler, KoinComponent { // 用Channel缓冲FilePart消息,设置容量避免内存溢出 private val filePartChannel = Channel<FilePartArgs>(capacity = Channel.UNLIMITED) private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob()) init { // 启动多个消费协程(数量可根据业务调整,比如8个) repeat(8) { scope.launch { for (args in filePartChannel) { filesHandler.onFilePartReceived(args...) } } } } fun messageResolver(message: String, socketCallback: (message: String) -> Unit) { scope.launch { if (message.isFilePart()) { // 将消息发送到Channel,而非直接处理 val filePartArgs = parseToFilePartArgs(message) // 自行实现消息解析逻辑 filePartChannel.send(filePartArgs) } // 其他类型消息的处理逻辑 } } // 记得在类销毁时清理资源 fun destroy() { filePartChannel.close() scope.cancel() } }
这样Socket接收线程只需要把消息扔进Channel,多个消费协程会自动并行处理每个分块,充分利用IO线程池的能力。
2. 解决竞态条件导致的重复"Complete"打印
问题根源:多个协程同时更新数据库状态,各自查询已完成数量时,可能出现"协程A刚更新未提交,协程B查询到旧数据"的情况,导致多次触发completeCount == total的判断。
这里有两种可靠的解决方式:
方式一:用Room事务保证操作原子性
把「更新状态+检查完成情况」的逻辑放进Room事务,让整个操作成为原子操作,不会被其他协程打断:
首先在你的Dao中添加事务方法:
@Transaction suspend fun updatePartAndCheckCompletion( fileHash: String, segment: Int, newStatus: Int, totalSegments: Int ): Boolean { // 先更新当前分块状态 updateFilePartStatus(fileHash, segment, newStatus) // 再查询已完成数量 val completeCount = getFilePartCountByStatus(fileHash, newStatus) // 返回是否全部完成 return completeCount == totalSegments }
然后修改onFilePartReceived的逻辑:
suspend fun onFilePartReceived(args...) { val filePart = filesRepository.getFilePartBySegment(receiver.filePart.fileHash, receiver.filePart.segment) ?: return if (filePart.status == FILE_PART_RECEIVED) return println("## File part received ${receiver.filePart.segment}") filesRepository.createAndWriteFilePart(server.serverId, receiver.filePart.fileHash, receiver.filePart.segment, receiver.filePart.data) // 用事务方法更新并检查完成状态 val isComplete = filesRepository.updatePartAndCheckCompletion( receiver.filePart.fileHash, receiver.filePart.segment, FILE_PART_RECEIVED, receiver.filePart.total ) if (isComplete) { println("## Complete -> ${receiver.filePart.total} parts received") } }
方式二:用Mutex按文件Hash加锁
如果不想修改Dao,也可以给每个文件Hash单独加锁,确保同一文件的分块处理串行执行,避免竞态:
在你的业务类(比如FilesHandler)中维护锁缓存:
class FilesHandler { // 用ConcurrentHashMap存储每个文件Hash对应的锁 private val fileLocks = ConcurrentHashMap<String, Mutex>() private fun getFileLock(fileHash: String): Mutex { return fileLocks.computeIfAbsent(fileHash) { Mutex() } } suspend fun onFilePartReceived(args...) { val fileHash = receiver.filePart.fileHash // 获取当前文件的锁,同一文件的分块处理串行执行 getFileLock(fileHash).withLock { val filePart = filesRepository.getFilePartBySegment(fileHash, receiver.filePart.segment) ?: return if (filePart.status == FILE_PART_RECEIVED) return println("## File part received ${receiver.filePart.segment}") filesRepository.createAndWriteFilePart(server.serverId, fileHash, receiver.filePart.segment, receiver.filePart.data) filesRepository.updateFilePartStatus(server.serverId, fileHash, receiver.filePart.segment, FILE_PART_RECEIVED) val completeCount = filesRepository.getFilePartCountByStatus(fileHash, FILE_PART_RECEIVED) if (completeCount == receiver.filePart.total) { println("## Complete -> $completeCount parts received") // 完成后移除锁,节省内存 fileLocks.remove(fileHash) } } } }
这种方式下,不同文件的分块仍能并行处理,只有同一文件的分块串行执行,既保证线程安全,又不损失并行性能。
内容的提问来源于stack exchange,提问作者james04
相关产品推荐
相关产品推荐

