分块列表的可控并行性实现: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() } } } }
二、现有代码的问题
- 未等待所有协程完成:
runMainJob只是启动协程就返回,导致cleanup会在请求还未全部发送时就执行。必须等待所有协程结束后再清理。 - 类型错误:
getUsersFromDatabase声明返回Users,但实际需要返回List<Users>,这会导致编译错误。 - 缺少TPS控制:仅用
Semaphore(2)控制并发数,如果API响应时间小于0.5秒,实际TPS会超过2,违反服务限制。 - 类设计不规范:
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
相关产品推荐
相关产品推荐

