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

Kotlin协程:子协程触发取消同级协程并保留自身结果

Kotlin协程批量请求处理问题

需求

批量发送请求并将响应封装为Map:若某请求返回5xx响应,需取消其他未完成的请求,最终返回包含该5xx响应及此前所有成功响应的Map。

初始实现与问题

我尝试在父协程中创建多个子async协程发送请求,但遇到问题:当某子协程收到5xx响应时,调用parentScope.coroutineContext.cancelChildren()会同时取消当前子协程,导致无法获取该5xx响应。相关代码如下:

fun sendRequests(
        requestSet: Set<XxxRequest>
    ): Mono<Map<XxxRequest, XxxResponse>> =
        mono {
            // 假设requestSet中有4个请求
            if (requestSet.isNotEmpty()) {
                val requestToResponse = async {
                    val parentScope = this
                    // 假设其中2个返回200 OK,第3个返回5xx
                    val listDeferred = requestSet.map { request ->
                        async {     
                            downStreamService.sendRequest(request) // 挂起函数,用WebClient发起HTTP请求
                            .also { 
                                if (it.status.is5xxServerError) {
                                    // 希望取消其他未完成任务,但确保当前任务正常执行
                                    // 调用cancelChildren()会取消当前子Job
                                    parentScope.coroutineContext.cancelChildren() 
                                }
                            }
                        }
                    }
                    // 期望responses大小为3,但实际只有2
                    val responses = listDeferred.filter { !it.isCancelled }.awaitAll()
                    requestSet.zip(responses).toMap()
                }
                requestToResponse.await()
            } else emptyMap()
        }

尝试NonCancellable方案(未解决)

我尝试用withContext(NonCancellable)包裹结果,但问题依然存在:

fun sendRequests(
        requestSet: Set<XxxRequest>
    ): Mono<Map<XxxRequest, XxxResponse>> =
        mono {
            if (requestSet.isNotEmpty()) {
                val requestToResponse = async {
                    val parentScope = this
                    val listDeferred = requestSet.map { request ->
                        async {     
                            downStreamService.sendRequest(request)
                            .let { 
                                if (it.status.is5xxServerError) {
                                    // 调用cancelChildren()会取消当前子Job
                                    parentScope.coroutineContext.cancelChildren() 
                                    withContext(NonCancellable) { it } // 无效
                                } else it
                            }
                        }
                    }
                    val responses = listDeferred.filter { !it.isCancelled }.awaitAll()
                    requestSet.zip(responses).toMap()
                }
                requestToResponse.await()
            } else emptyMap()
        }

手动维护结果Map的修改方案

后续根据建议修改代码,手动封装结果Map,但调试发现所有HTTP请求实际已完成,仅未将剩余响应加入Map:

fun sendRequests(
        requestSet: Set<XxxRequest>
    ): Mono<Map<XxxRequest, XxxResponse>> =
        mono {
            if (requestSet.isEmpty()) return@mono emptyMap<XxxRequest, XxxResponse>()
            val mutex = Mutex()
            val resultMap = mutableMapOf<XxxRequest, XxxResponse>()
            withContext(Dispatchers.IO.limitedParallelism(4)) {
                requestSet.map { request -> 
                    launch {
                        downStreamService.sendRequest(request)
                        .let { response -> 
                            mutex.withLock {
                               resultMap += request to response
                               if (response.status.is5xxServerError) this@withContext.coroutineContext.cancelChildren()
                            }
                        }
                    }
                }.joinAll()
            }
            resultMap
        }

核心疑问

当某协程收到WebClient返回的5xx响应时,能否真正取消其他协程发送的下游请求?

下游请求的实现代码:

suspend fun sendRequest(request: XxxRequest): XxxResponse =
        runCatching {
            webClient
                .put()
                .uri {
                    it.path("/xxx").build()
                }
                .body(jsonMapper.writeValueAsString(request).toMono(), String::class.java)
                .headers { addHeaders(it) }
                .accept(MediaType.APPLICATION_JSON)
                .awaitExchange { res ->
                    res.awaitEntity(String::class)
                        .let {
                            // parseEntityString 是非挂起函数
                            parseEntityString(it.body, it.statusCodeValue)
                        }
                }
        }.getOrElse {
            val exceptionMessage = it.message
            val exceptionCause = it.cause?.message
            XxxResponse(
                status = HttpStatus.INTERNAL_SERVER_ERROR,
                error = Error(exceptionMessage, exceptionCause)
            )
        }

内容的提问来源于stack exchange,提问作者NotGoodAtEnglish

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 17:55:59