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

分块列表的可控并行性实现:Kotlin协程方案验证与优化

Kotlin协程并发处理用户导出问题解答

一、验证并行执行的方法

  • 打印线程与时间戳:在callSecondService函数中添加日志,输出当前线程名和时间戳:
    suspend fun callSecondService(group: List<Users>) {
        println("处理分组,线程:${Thread.currentThread().name},时间:${System.currentTimeMillis()}")
        // 模拟网络请求
        delay(200)
    }
    
    若日志中出现不同线程名,或同一时间段内有多条输出,说明是并行执行。
  • 统计总执行时长:10000个用户分为200组,按2TPS的要求,理论总耗时约100秒(200组 ÷ 2组/秒)。如果实际耗时接近这个值,说明并行符合预期;若耗时远超100秒,大概率是串行执行。
  • 跟踪活跃任务数:用原子变量记录同时运行的任务数,查看最大值是否接近设置的并发度(这里是2):
    import java.util.concurrent.atomic.AtomicInteger
    
    suspend fun runMainJob(userList: List<Users>) = coroutineScope {
        val activeJobs = AtomicInteger(0)
        val semaphore = Semaphore(2)
        userList.chunked(50).forEach { group ->
            launch {
                semaphore.withPermit {
                    val current = activeJobs.incrementAndGet()
                    println("当前活跃任务数:$current")
                    callSecondService(group)
                    activeJobs.decrementAndGet()
                }
            }
        }
    }
    

二、现有代码的问题

  1. 未等待所有协程完成:runMainJob只是启动协程就返回,导致cleanup会在请求还未全部发送时就执行。必须等待所有协程结束后再清理。
  2. 类型错误:getUsersFromDatabase声明返回Users,但实际需要返回List<Users>,这会导致编译错误。
  3. 缺少TPS控制:仅用Semaphore(2)控制并发数,如果API响应时间小于0.5秒,实际TPS会超过2,违反服务限制。
  4. 类设计不规范:Users类用var声明属性,不符合Kotlin优先使用不可变对象的原则,建议改为val并使用data class。

三、更简便的实现方式

使用coroutineScope管理协程生命周期,结合joinAll等待所有任务完成,同时添加速率控制保证不超过2TPS:

import kotlinx.coroutines.*
import kotlinx.coroutines.sync.Semaphore
import kotlinx.coroutines.sync.withPermit

// 改用不可变数据类
data class Users(
    val id: Long,
    val name: String,
    val age: Long,
    val email: String,
)

fun main() = runBlocking {
    val userList = getUsersFromDatabase()
    runMainJob(userList)
    cleanup()
}

// 用suspend函数+coroutineScope管理协程
suspend fun runMainJob(userList: List<Users>) = coroutineScope {
    val semaphore = Semaphore(2)
    // 启动所有协程并收集Job对象
    val jobs = userList.chunked(50).map { group ->
        launch {
            semaphore.withPermit {
                callSecondService(group)
                // 控制速率:每组至少间隔0.5秒,保证2TPS
                delay(500)
            }
        }
    }
    jobs.joinAll() // 等待所有请求完成
}

fun cleanup() {
    println("执行数据库清理操作")
}

// 模拟获取10000个用户列表
fun getUsersFromDatabase(): List<Users> {
    return (1..10000).map {
        Users(it.toLong(), "用户$it", (20..50).random().toLong(), "user$it@example.com")
    }
}

// 模拟API调用(挂起函数)
suspend fun callSecondService(group: List<Users>) {
    // 模拟网络请求延迟
    delay((100..300).random().toLong())
    println("完成分组处理,组内用户数:${group.size}")
}

实现优势

  • 用coroutineScope自动管理子协程,joinAll确保所有请求完成后才执行cleanup
  • delay(500)保证每组至少间隔0.5秒,严格控制在2TPS以内
  • Users改用data class和不可变属性,代码更规范
  • 模拟了真实网络延迟,方便验证并行效果

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 05:00:40