Kotlin中Dispatchers.IO.limitedParallelism未限制并行的问题排查
问题原因
你对limitedParallelism(1)的理解有偏差:这个API限制的是同一时间处于运行状态(非挂起)的协程数量,而非协程的启动顺序或整个任务的串行执行。
当协程执行到delay(100)时会挂起,此时它会释放占用的线程,调度器会立即启动下一个等待的协程。你的代码中,counter.incrementAndGet()是极快的原子操作,所有1000个协程都会在第一个协程挂起前(或挂起后迅速)完成计数器自增,导致计数器直接升到1000,最终max值超出预期。
修正方案
如果需要让协程逐个完整执行(即一个协程完成自增→delay→自减的全流程后,下一个才开始),可以通过Mutex加锁,将整个协程逻辑包裹在临界区内:
@Test fun `testParallelism`(){ runBlocking { val maxParallel = 1 val ctx = Dispatchers.IO.limitedParallelism(maxParallel) val mutex = Mutex() val max = AtomicLong(Long.MIN_VALUE) val counter = AtomicLong(0) (0..1000).map { launch(ctx) { mutex.withLock { val c = counter.incrementAndGet() max.set(Math.max(max.get(), c)) delay(100) counter.decrementAndGet() } } }.joinAll() assertTrue { max.get() <= maxParallel } } }
另外,代码中无需每次创建新的CoroutineScope,直接使用runBlocking的作用域启动协程即可,避免不必要的作用域创建。
内容的提问来源于stack exchange,提问作者Mangat Rai Modi
相关产品推荐
相关产品推荐

