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) }
问题原因与解决方案
核心原因
- AWS SDK同步调用阻塞线程:你用的
SageMakerRuntimeClient.invokeEndpoint是同步阻塞式API,它会一直占用当前线程直到请求完成,期间不会释放线程给其他协程。 - 无真实挂起点:你的
invokeEndpoint挂起函数里,除了被注释的delay外没有真正的挂起操作。Kotlin协程的线程切换只发生在释放线程的挂起函数调用时(比如delay、Future.await()等),同步阻塞调用不会触发挂起,自然不会切换线程。 - 代码遗漏
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
相关产品推荐
相关产品推荐

