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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 10:27:34