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

实现对调用层透明的批量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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 21:36:03