优雅替代单线程池实战)
我先说一个我自己的经历。早期做订单系统时为了保证同一个用户的操作日志不乱序代码里到处是Executors.newSingleThreadExecutor()。当时觉得这个方案简单可靠直到有一天线上日志顺序错乱排查下来才发现是有个地方每次调用都新建了一个线程池根本没有复用。后来切换到了Kotlin协程我才彻底想明白问题出在哪线程池太“重”了而且它把“串行”这件事绑定在了一个线程的生命周期上。Kotlin协程里恰好有一个专门干这个事的API——limitedParallelism尤其是Dispatchers.IO.limitedParallelism(1)在业务上能替代单线程池实现串行、轻量、又不用手动关线程池的效果。这篇文章我就从原理到实战把为什么它可以替代单线程池、具体怎么用、有哪些坑一次性讲清楚。1. 单线程池为什么会让Java代码又慢又乱1.1 单线程池真正解决的是“乱”不是“慢”先理清一个概念很多并发问题的本质是多个线程同时操作同一个资源导致数据“乱”了而不是系统本身不够快。比如多个请求同时修改同一个账户余额、多个线程同时追加同一条日志文件、多个请求同时访问一个不安全的第三方SDK。解决方案无非就几类加锁、CAS、或者干脆串行化。而Java里最简单粗暴的串行化方案就是单线程池。Executors.newSingleThreadExecutor()内部其实是一个核心线程数为1、队列无界的ThreadPoolExecutor。提交给它的所有任务都会进入队列由唯一一个线程依次执行。这种设计从源头消灭了并发不存在多线程竞争也不需要考虑锁的粒度、死锁、活性问题项目里很多资深的工程师都偏爱用它来处理一些“需要严格保证顺序”的场景。在Java面试里如果被问到“什么时候用单线程池”很多人的第一反应是“并发量低的时候”。这个回答其实不太准确。并发量低但任务之间没有共享资源用普通线程池甚至直接同步执行都可以。真正适合单线程池的场景只有一个核心特征多个任务之间必须按提交顺序执行且任务之间不能并发。比如我要保证同一个用户的订单创建日志一定出现在支付日志之前那我只需要把这两个日志提交到同一个单线程池里剩下的交给队列就好。1.2 单线程池的硬伤一个线程堵住全队列卡死单线程池好用但代价也摆在明面上。最大的问题就是这一个线程承担了所有串行任务一旦某个任务长时间阻塞队列里所有后续任务全部被卡住。举个例子如果单线程池里的某个任务调用了外部的HTTP接口而这个接口响应很慢甚至超时等待了30秒那么这30秒内队列里其他任务即使都是五六毫秒就能跑完的轻量任务也只能排队等着。这个现象在Java异步里面有一个自嘲的说法单线程池的“队头阻塞”。线程虽然只有一个但是无界队列可以无限接任务最后堆积的是延迟和内存。第二个问题是线程资源太昂贵。JVM里每个线程默认栈大小通常在512KB到1MB之间创建100个单线程池光是线程栈就占掉差不多100MB的虚拟内存。更麻烦的是线程不是“用完即走”的如果没有正确调用shutdown()线程会一直在后台存活。在Web容器或者Android应用里线程泄漏是常有的事最终导致OutOfMemoryError。热词里经常出现的java: outofmemoryerror: insufficient memory有相当一部分就是这种乱创建线程导致的。第三个问题是“一个池只能串行一种资源”。如果系统里有十几种需要分别保证顺序的资源你难道要建十几个单线程池线程数量会直接爆炸。所以有人会做分桶按某个key取模分到固定数量的池子里比如ExecutorService[] pools new ExecutorService[4]; for (int i 0; i pools.length; i) { pools[i] Executors.newSingleThreadExecutor(); } void submitOrderEvent(long userId, Runnable task) { pools[(int) (userId % pools.length)].submit(task); }分桶可以缓解线程数量问题但同时引入了一个新的风险任务之间如果存在依赖关系比如A池的任务提交给了B池并等待结果而B池的线程又在等待A池执行完就容易出现线程饥饿甚至死锁。单线程池本身不会产生死锁但它和多个池的编排组合在一起就很容易踩到这些坑。1.3 为什么切到协程后单线程池就不香了如果项目里还是纯Java上面这些坑你可以用各种工程手段规避。但如果项目已经引入了Kotlin协程再回去用单线程池就会觉得特别拧巴。最大的拧巴点是单线程池的execute(Runnable)没法直接执行suspend函数。你要么在里面套一层runBlocking要么额外包一个CoroutineScope。而runBlocking放在单线程池里相当于那个唯一的线程直接在等待协程结束一旦协程内部再有挂起操作这个线程就被占用后面排队的任务全部卡住。这种“两层世界”的别扭感很多人应该都有体会。协程最核心的思维是“挂起而不是阻塞”。同样的串行需求在协程里可以把“限制并发”这件事交给调度器而不是交给一个独占线程。这也是limitedParallelism出现的原因它希望用协程的方式来解决并发限制问题又不想让业务方去关心线程生命周期。2. Kotlin协程调度器一个更好的“串行化底座”2.1 协程不上线程挂起不等于阻塞先给不太熟悉协程的朋友补个基础。协程本质上是一种“可挂起”的计算单元它运行在线程之上但又不和线程绑定。线程在运行一个协程时如果协程执行到suspend函数线程并不会傻傻等着而是会被释放出来去执行别的协程。等挂起的条件满足了协程再从上次挂起的地方继续执行。这有点像银行柜台。线程是柜台窗口每个窗口配一个营业员。传统线程模型里一个客户办理业务办了半小时营业员就一直陪着后面所有人排队。协程模型里这个客户可以先填单、交资料然后回家等短信通知营业员马上接待下一个人。过一会儿短信来了客户再回到某个窗口继续办理。所以协程可以用非常少的线程承载非常多的并发任务。JVM上创建一个线程可能要分配512KB到1MB的栈空间而创建一个协程对象可能只需要几百字节到几KB。你可以放心地创建几万个协程但如果创建几万个线程内存直接扛不住。这也是为什么协程特别适合IO密集场景。举个最简单的例子fun main() runBlocking { val start System.currentTimeMillis() repeat(10_000) { launch { delay(1000) } } println(cost ${System.currentTimeMillis() - start} ms) }这段代码创建一万个协程每个挂起1秒。如果换成一万个线程同样逻辑内存和CPU都会很吃力。但协程几乎零压力地跑完。原因就是delay是挂起操作不会占用线程。2.2 Dispatchers.IO 和 Dispatchers.Default别用错调度器协程最终还是要跑在线程上的这个“决定协程运行在哪类线程上”的东西就是CoroutineDispatcher。Kotlin协程在JVM上提供了几个内置调度器需要先弄清楚它们的分工。调度器适合任务默认并发度Dispatchers.MainUI主线程更新界面1Dispatchers.DefaultCPU密集型计算通常与CPU核心数相关Dispatchers.IOIO密集型、阻塞调用默认64左右可调整Dispatchers.Default主要用于大量计算、排序、解析等纯CPU密集操作默认并发度跟CPU核心数有关通常不会超过核数太多因为计算密集任务一旦超过核数反而会因为上下文切换变慢。Dispatchers.IO则是为数据库访问、网络请求、文件读写这类任务设计的。它的默认并发上限早期是64后来版本允许通过kotlinx.coroutines.io.parallelism系统属性调整上限可以非常大理论最大能到65535。需要说明的是JVM上Dispatchers.IO和Dispatchers.Default底层会共享线程资源这样可以减少线程切换开销。很多初学者会把“调度器”和“线程池”混为一谈其实不完全是。调度器是一套任务分配策略底层确实有线程池但协程的挂起能力让它有了更多主动权。Dispatchers.IO本身也只是一个特殊的“受限调度器”它限制并发量但底层线程是按需扩缩的。理解了这一点再看limitedParallelism就会非常顺。3. limitedParallelism(1) 的正确打开方式3.1 limitedParallelism 内部做了什么limitedParallelism是kotlinx.coroutines从1.6.0开始提供的一个扩展函数定义很简单fun CoroutineDispatcher.limitedParallelism(parallelism: Int): CoroutineDispatcher它会基于当前调度器返回一个新的调度器视图这个视图只允许指定数量的任务同时执行。本质上它在原来的调度器前面加了一道“限量闸门”。我拿Dispatchers.IO.limitedParallelism(1)举例。这个表达式返回的调度器底层仍然使用Dispatchers.IO的线程池但同一时间最多只有一个任务在这个调度器上运行。其余任务会在闸门外排队等前一个任务执行完再进入。实现上这个调度器使用类似信号量或原子计数的方式来控制并发。任务执行前尝试获取许可获取到了就执行执行完释放。所以你可以同时创建多个受限调度器val serialA Dispatchers.IO.limitedParallelism(1) val serialB Dispatchers.IO.limitedParallelism(1) launch(serialA) { log(A) } launch(serialB) { log(B) }serialA和serialB各自最多只有一个任务在跑但它们底层共享Dispatchers.IO的线程池。两个受限调度器之间互不影响。这比创建两个单线程池要节省资源因为单线程池是两个固定线程即使任务全部挂起线程也一直在那里而协程的受限调度器在线程空闲时底层线程可以被其他协程复用。在Kotlin协程出现早期如果大家想限制并发度通常会自己写一个信号量val gate Semaphore(1) suspend fun T serial(block: suspend () - T): T gate.withPermit { block() }这种写法能用但它把限制逻辑散落到业务代码里而且只控制“临界区”并不参与调度器对线程的选择。一旦线程阻塞而不是挂起信号量保护不了线程资源。limitedParallelism则把并发限制内聚到了调度器这一层使用方不需要额外包一层代码更干净语义也更清晰。顺带提一个容易踩坑的细节limitedParallelism(0)是非法的源码会直接抛IllegalArgumentException。0没有意义因为并发度至少是1。传1表示“单飞”传n表示“最多n飞”。3.2 为什么是 limitedParallelism(1)而不是 limitedParallelism(0)很多刚接触这个API的同学会有一个疑问我要求“单线程”为什么不传0原因很简单并发度0意味着什么任务都不能执行这显然不合理。limitedParallelism里的参数是并发任务数的上限而不是线程编号。所以“单线程串行”对应的参数就是1。那limitedParallelism(1)和直接用Dispatchers.Default或者Dispatchers.IO有什么区别区别在于底层调度器允许的并发量很大可能同时有64个IO任务在执行。而limitedParallelism(1)给这整条链路加了一道闸门同一时间只放一个任务进去。你可以想象成一条多车道的公路你的调度器原本允许60辆车并行现在你给某个方向单独开了条“单车道”只允许一辆车依次通行。还有一点很有用limitedParallelism可以突破底层调度器的默认并发限制。比如你的机器只有4核Dispatchers.Default默认并发度可能只有4如果你执行一些需要更多并行IO的场景可以用Dispatchers.Default.limitedParallelism(8)显式把并发上限放大到8。当然并发放大不意味着性能一定提升线程一多上下文切换成本也上来了。但是在某些场景下这确实比自己去维护一个8线程的线程池简洁得多。3.3 对比Java单线程池一次全面的对照把Executors.newSingleThreadExecutor()和Dispatchers.IO.limitedParallelism(1)放在一起做一个直接对比维度Executors.newSingleThreadExecutor()Dispatchers.IO.limitedParallelism(1)串行能力同一队列内严格串行并发度限制为1等效串行线程占用固定占用一个线程不固定占用挂起时让出底层线程任务类型Runnable/Callable任意suspend函数生命周期需要手动shutdown()随协程作用域自动结束取消任务Future.cancel()Job.cancel()结构化取消与协程集成需要asCoroutineDispatcher()桥接原生支持内存开销一个线程栈约512KB~1MB一个轻量Dispatcher对象从业务效果上看两者都解决了“串行执行”的问题。但从资源占用和运维成本来看差别非常明显。单线程池是“我为这个串行需求专门雇一个营业员”不管业务忙不忙这个人都得在窗口坐着limitedParallelism(1)是“我允许这个通道同一时间只能有一个客户办理”但营业员是从共用团队里临时抽调过来的办完可以立刻去忙别的。所以在协程项目里我基本不再创建单线程池了。凡是需要串行化的地方要么用limitedParallelism(1)要么按业务维度分片后每个片一个limitedParallelism(1)。4. 实战改造把单线程池换成limitedParallelism4.1 改造场景按用户串行写审计日志假设现在有一个订单系统每个用户会触发多个审计事件比如“创建订单”“支付成功”“退款申请”。审计日志必须保证同一个用户的事件顺序是准确的否则后续对账排查都会出问题。如果同一个用户的两个事件同时到达而底层写入是线程不安全的就会乱序。最朴素的Java实现通常会给每个用户分配一个单线程池MapLong, ExecutorService userExecutors new ConcurrentHashMap(); void writeLog(long userId, LogEntry entry) { ExecutorService executor userExecutors.computeIfAbsent( userId, id - Executors.newSingleThreadExecutor() ); executor.execute(() - saveLog(entry)); }这个实现有几个问题。用户量大时每个用户一个线程线程数爆炸。如果忘记shutdown()线程会一直存活造成泄漏。而且saveLog如果是一个阻塞的IO操作那个线程就会被一直占用。更麻烦的是这种代码在协程环境里用起来很别扭没法直接写suspend逻辑。4.2 改造后代码从单线程池到协程调度器换成协程后同样的逻辑可以这样写