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

Kotlin协程并行调用AWS SageMaker异常问题排查

Kotlin协程调用AWS SageMaker时线程不切换的问题

我正在构建一个调用AWS SageMaker端点获取预测结果的服务,因请求量较大,想用Kotlin协程实现并行调用。但作为协程新手,发现运行结果和预期不符:

代码逻辑是把输入特征列表分片,通过async启动协程分别发送请求至SageMaker。按理解,协程轻量,挂起函数调用后线程应该会切换,但日志显示协程线程始终不变,只有取消invokeEndpoint方法中delay(100)的注释时,线程才会切换。

疑问:为什么会出现这种情况?担心AWS SageMaker的调用是不是一直阻塞线程和协程?

代码示例

fun forecastBatch(
    features: Set<String>
): Map<String, Float> {
    val partitions: List<List<String>> = features.chunked(MAX_BATCH_SIZE)
    val result: MutableMap<String, Float> = mutableMapOf()
    runBlocking {
        val deferreds = partitions.mapIndexed { idx, partition ->
            async(CoroutineName("Coroutine-${idx}")) { // 注:原代码漏写async,否则无法创建并行协程
                log.info("${Thread.currentThread().name} Coroutines:  ${currentCoroutineContext()[CoroutineName]} for Sending Request ${idx}X")
                val batchRequest = partition.joinToString("\n")

                val forecastResult = invokeEndpoint(batchRequest)
                
                log.info("${Thread.currentThread().name} Coroutines:  ${currentCoroutineContext()[CoroutineName]}  Receiving Request ${idx}X")

                // ...构建响应逻辑
            }
        }
        deferreds.awaitAll().forEach{
            result += it
        }
    }
    return result
}

suspend fun invokeEndpoint(input: String): String {
    // 构建请求逻辑...        

    val result: InvokeEndpointResponse = sageMakerRuntimeClient.invokeEndpoint(request)
    
    // 取消注释这一行,同一个协程的线程会切换
    // delay(100)
    return result.body().asString(Charsets.UTF_8)
}

问题原因与解决方案

核心原因

  1. AWS SDK同步调用阻塞线程:你用的SageMakerRuntimeClient.invokeEndpoint是同步阻塞式API,它会一直占用当前线程直到请求完成,期间不会释放线程给其他协程。
  2. 无真实挂起点:你的invokeEndpoint挂起函数里,除了被注释的delay外没有真正的挂起操作。Kotlin协程的线程切换只发生在释放线程的挂起函数调用时(比如delay、Future.await()等),同步阻塞调用不会触发挂起,自然不会切换线程。
  3. 代码遗漏async:原代码中mapIndexed里直接写CoroutineName但没调用async,导致根本没启动并行协程,所有逻辑都在runBlocking的当前线程串行执行。

解决方案

方案1:改用AWS异步SDK(推荐)

AWS提供了异步客户端SageMakerRuntimeAsyncClient,其invokeEndpoint方法返回CompletableFuture,配合kotlinx-coroutines-jdk8的await()扩展函数可转为挂起函数,实现真正的协程挂起与线程复用:

// 初始化异步客户端
private val sageMakerAsyncClient = SageMakerRuntimeAsyncClient.create()

suspend fun invokeEndpoint(input: String): String {
    // 构建请求...
    val request = InvokeEndpointRequest.builder()
        .endpointName("your-endpoint-name")
        .contentType("text/plain")
        .body(SdkBytes.fromString(input, Charsets.UTF_8))
        .build()
    
    // 用await()挂起协程,释放线程给其他协程
    val result = sageMakerAsyncClient.invokeEndpoint(request).await()
    return result.body().asString(Charsets.UTF_8)
}
方案2:将同步调用移至IO调度器

如果必须使用同步SDK,需把阻塞调用放到Dispatchers.IO调度器(专门处理IO阻塞的线程池),避免阻塞协程线程:

suspend fun invokeEndpoint(input: String): String = withContext(Dispatchers.IO) {
    // 构建请求...
    val result = sageMakerRuntimeClient.invokeEndpoint(request)
    result.body().asString(Charsets.UTF_8)
}

补充说明

  • runBlocking默认使用当前线程作为协程执行线程,如果协程内无挂起操作,所有逻辑都会串行执行,无法实现并行。
  • 协程的“轻量并行”本质是靠挂起操作释放线程,让多个协程复用同一线程池中的线程,没有挂起的协程和普通串行代码无区别。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 14:52:54