Kotlin协程:如何高效处理并发重复请求?
Kotlin协程实现并发重复请求复用处理结果
你的需求核心是并发相同请求共享正在执行的处理结果,避免重复执行耗时操作,同时在处理完成后允许新请求重新计算。借助Kotlin协程的Deferred和Mutex可以高效实现这个逻辑,核心思路是维护一个"正在处理请求"的缓存,让重复请求直接等待已有任务的结果。
核心实现逻辑
- 请求标识:为每个请求生成唯一标识(比如基于请求参数的哈希或字符串),用来区分"相同请求"
- 任务缓存:用映射表存储正在处理的请求对应的
Deferred<Response>,代表未完成的协程任务 - 线程安全保护:用
Mutex保证对映射表的操作是原子的,避免并发竞争 - 复用/新建任务:收到请求时,先检查缓存中是否有正在处理的相同任务:
- 有则直接等待该任务的结果
- 没有则创建新协程执行耗时操作,将任务存入缓存,执行完成后移除缓存条目
示例代码
import kotlinx.coroutines.* import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock // 示例请求类:根据业务定义相等性(参数相同则视为相同请求) data class Request(val requestId: String, val businessParams: Map<String, String>) // 示例响应类 data class Response(val payload: String, val processedAt: Long) class RequestProcessor { // 存储正在处理的请求任务:请求标识 -> 未完成的协程任务 private val pendingTasks = mutableMapOf<String, Deferred<Response>>() // 保证pendingTasks的线程安全访问 private val mutex = Mutex() suspend fun processRequest(request: Request): Response { // 生成请求唯一标识:可根据业务调整(比如用businessParams的哈希值) val requestKey = request.businessParams.toString().hashCode().toString() return mutex.withLock { // 检查是否已有正在处理的相同请求 pendingTasks[requestKey]?.let { existingTask -> return@withLock existingTask.await() } // 创建新协程执行耗时操作 val newTask = CoroutineScope(Dispatchers.IO).async { try { // 模拟耗时操作:比如数据库查询、API调用、文件IO等 delay(2000) println("Executing new task for request ${request.requestId}") Response( payload = "Processed result for params: ${request.businessParams}", processedAt = System.currentTimeMillis() ) } finally { // 无论成功/失败,处理完成后移除缓存条目 mutex.withLock { pendingTasks.remove(requestKey) } } } // 将新任务存入缓存 pendingTasks[requestKey] = newTask // 等待任务执行结果 newTask.await() } } } // 测试代码 fun main() = runBlocking { val processor = RequestProcessor() // 构造一个重复请求:businessParams相同则视为相同请求 val repeatedRequest = Request( requestId = "REQ-001", businessParams = mapOf("user_id" to "123", "action" to "get_profile") ) // 模拟5个并发的重复请求 println("Starting 5 concurrent repeated requests...") val concurrentJobs = List(5) { jobIndex -> launch { val start = System.currentTimeMillis() val response = processor.processRequest(repeatedRequest) val duration = System.currentTimeMillis() - start println("Job $jobIndex finished in ${duration}ms, response: $response") } } concurrentJobs.forEach { it.join() } // 模拟处理完成后,新的相同请求(会重新执行耗时操作) println("\n--- Sending new request after first batch ---") val start = System.currentTimeMillis() val newResponse = processor.processRequest(repeatedRequest) val duration = System.currentTimeMillis() - start println("New request finished in ${duration}ms, response: $newResponse") }
关键细节说明
- 异常处理:如果耗时操作抛出异常,所有等待该任务的并发请求都会收到相同的异常,符合业务预期(同一个请求处理失败,所有重复请求都应该得到失败结果)
- 内存管理:
pendingTasks只会存储未完成的任务,处理完成后会立即移除,不会造成内存泄漏 - 协程上下文:示例中用
Dispatchers.IO处理IO密集型操作,CPU密集型任务可替换为Dispatchers.Default - 请求标识优化:实际业务中可以用更高效的方式生成请求标识,比如对关键参数进行哈希(如
Objects.hash(params.values)),避免长字符串的性能损耗
内容的提问来源于stack exchange,提问作者vivi
相关产品推荐
相关产品推荐

