实现对调用层透明的批量API请求的最优设计模式咨询
批量处理器优化实现方案
核心采用CompletableDeferred作为单输入的结果占位符,配合Channel做输入缓冲,后台异步攒批执行,完全对外屏蔽批量逻辑,保留1:1的调用形式。
优化后完整实现
import kotlinx.coroutines.* import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock class BatchService<A, B>( private val batchSize: Int, private val scope: CoroutineScope, // 批量调用逻辑注入,和业务解耦 private val batchCall: suspend (List<A>) -> List<B?> ) { // 内部请求包装类:输入参数 + 待完成的结果占位符 private data class Request<A, B>( val input: A, val result: CompletableDeferred<B?> ) private val requestChannel = Channel<Request<A, B>>(capacity = Channel.UNLIMITED) private val mutex = Mutex() private var isProcessorRunning = false // 对外暴露的1:1调用接口,业务侧完全感知不到批量逻辑 suspend fun call(input: A): B? { // 可在此处加前置过滤逻辑,满足条件直接返回null,无需进入批次 // if (filterCondition(input)) return null val result = CompletableDeferred<B?>() requestChannel.send(Request(input, result)) // 后台批处理任务仅启动一次 tryStartProcessor() return result.await() } private fun tryStartProcessor() { scope.launch { mutex.withLock { if (isProcessorRunning) return@launch isProcessorRunning = true } // 循环攒批执行 while (isActive) { val batch = mutableListOf<Request<A, B>>() // 先取1个请求避免空跑 val firstRequest = requestChannel.receiveCatching().getOrNull() ?: break batch.add(firstRequest) // 凑最多batchSize-1个已就绪的请求组成完整批次 repeat(batchSize - 1) { requestChannel.tryReceive().getOrNull()?.let { batch.add(it) } } // 异步执行当前批次,不阻塞后续批次攒取 scope.launch { val inputs = batch.map { it.input } val results = runCatching { batchCall(inputs) }.getOrDefault(List(inputs.size) { null }) // 结果映射回对应请求的占位符,执行完立刻返回 batch.forEachIndexed { index, request -> request.result.complete(results.getOrNull(index)) } } } mutex.withLock { isProcessorRunning = false } } } // 可选优雅关闭方法 fun close() { requestChannel.close() } }
问题解决说明
- 问题1(API一次性全部发起):通过通道做输入缓冲,后台会持续攒批次,每凑够指定大小就立刻发起调用,不用等所有输入提交完成,批次之间并行执行,可灵活控制并发批次上限。
- 问题2(首个输入需等所有结果返回):每个批次执行完成后立刻将结果写回对应请求的占位符,对应的
call方法会直接唤醒返回结果,不需要等其他批次执行完成,天然严格保序。 - 问题3(中间容器不规范):无全局输入/输出列表,不需要手动维护索引和清理逻辑,所有请求的生命周期和对应的
CompletableDeferred绑定,结果返回后自动被GC回收,无内存泄漏和索引错位风险。
使用示例
// 初始化批量服务 val userQueryService = BatchService<Long, User>( batchSize = 10, scope = CoroutineScope(Dispatchers.IO + SupervisorJob()), batchCall = { ids -> // 注入你的批量调用逻辑 userApi.batchGetUserInfo(ids) } ) // 业务侧直接1:1调用 suspend fun getUserInfo(userId: Long): User? { return userQueryService.call(userId) } // 原有批量处理逻辑无需修改 class Processor<A, B> { val service: BatchService<A, B> val scope: CoroutineScope fun processBatch(input: List<A>) { input.map { Pair(it, scope.async { service.call(it) }) }.map { (a, deferred) -> runBlocking { deferred.await().let { // 原有结果处理逻辑 } } } } }
内容的提问来源于stack exchange,提问作者lod
相关产品推荐
相关产品推荐

