Java 多线程总结:线程、锁、线程池与异步协作
从“大量数据怎样导入”和“一笔订单怎样拆成多个任务”出发,看懂多线程究竟在解决什么。
主体以Java 17 的平台线程为基线;末尾单独说明 Java 21 虚拟线程。导入数量、批次大小、库存和线程池参数都是教学示例,不是实际项目指标或通用配置。
多线程最容易让人困惑的地方,不是 API 多,而是几个不同问题混在一起:任务怎么执行?共享数据怎么保护?多个任务怎么等结果?忙不过来怎么办?
沿着这四个问题看,线程、锁、线程池和CompletableFuture就有了各自的位置。
一、为什么要多线程:从两个常见业务场景说起
写业务时,多线程通常不是为了“同时打印两句话”,而是遇到了两类问题:数据太多,一个任务处理得慢;一次操作要做很多事,没必要全部排成一条长队。
场景一:导入的数据很多,分批交给多个线程处理
假设运营人员上传一个有10 万条商品数据的文件。每条数据都要检查必填项、转换格式,再保存到数据库。
最直接的实现是一个线程从头处理到尾。小文件这样做就够了;数据量大时,就需要看看能否拆开处理。
这里假设各条商品数据互不依赖,可以这样安排:**文件按顺序读取,每 500 条组成一批,交给 4 个工作线程处理。**一个线程完成当前批次后,再接下一批,而不是给每一行数据创建一个线程。
图里最重要的不是“4 个线程”,而是三个动作:**分批、限制并发、汇总结果。**有界线程池限制执行资源;读取端也要控制已提交但尚未回收结果的批次数,达到上限就先等一个结果,不继续无限读取和堆积。[pool][completion]
后台页面可以先显示“任务已受理”,再查询导入进度;全部处理结束后,展示成功条数、失败条数和失败原因。**“提交成功”不等于“导入完成”。**实际产品中可以保存任务记录,让用户离开页面后仍能查看结果。
不过,数据多不代表一定要加线程。先做好批量写入、必要索引和重复数据校验,再根据数据库承受能力调整并发。若多个批次都在更新同一个商品,或者前一行是后一行的前提,就不能当成独立任务随意并发。
也不要让几个线程直接争抢同一个文件读取位置。这里采用“一个读取方准备独立批次,多个工作线程处理批次”的设计。**多批次各自提交,不等于整个文件一个事务。**如果要求整个文件全部成功或全部失败,就要另行设计暂存、统一校验和提交过程,不能仅靠多线程解决。
场景二:用户一次操作,要触发多项订单业务
以一个简化的、先支付后发货的订单系统为例:用户下单后,系统要校验数据、生成订单;支付成功后,还要安排发货、发送通知、更新运营统计。
这些事情不是全部同时做,也不是全部串行做。先看依赖关系:
在这个例子里,生成订单需要先完成;确认支付成功、订单状态提交后,通知仓库备货、发送支付成功通知、更新非核心统计可以分别交给后台任务。这里假设这三项业务互不依赖,不要求作为同一笔事务同时成功。
但“发送支付成功通知”和“发送发货通知”不是一回事。**创建发货任务,只表示开始安排履约;只有实际出库、物流单号等发货信息确认后,才能更新已发货状态,再发送发货通知。**具体条件由业务规则决定。
因此,用户不必在支付结果页面一直等到仓库完成发货;系统也不能在订单还没生成时,就启动依赖这笔订单的发货动作。
落到 Spring 项目时,可以把后续任务的触发点放在事务提交成功之后。@TransactionalEventListener默认对应AFTER_COMMIT;普通线程绑定事务不会自动传播到新线程,不能以为主方法加了@Transactional,子线程就都在同一个事务里。[spring-event][spring-tx]
回头看几个名词,就容易理解了
进程可以理解为正在运行的程序实例;线程是进程里的执行路径。多个导入工作线程属于同一个 Java 进程,可以共享对象和资源,但各有自己的调用栈。[^concept]
并发:多个导入批次在一段时间内都在推进。并行:多个批次在同一时刻执行,通常利用多核。异步:任务交出去后,调用方不必原地等它完成。三者不是同义词;后台异步处理也可以只有一个工作线程。
多线程可能提高处理能力,但不是越多越快。数据库连接、下游接口和 CPU 都有容量边界,线程增加后也会带来调度与竞争成本。[^pool]
二、线程怎么开始、等待和结束
先运行最小例子:
Threadworker=newThread(()->{System.out.println(Thread.currentThread().getName()+":校验一个导入批次");},"import-worker");worker.start();worker.join();// 当前线程等待 worker 结束;外层需处理或声明 InterruptedExceptionSystem.out.println("任务已结束");Runnable是任务说明,Thread是执行者。start()才会启动线程;直接调用run()只是当前线程中的普通方法调用。同一个线程对象不能再次start()。[^thread]
RUNNABLE不代表这一刻一定占着 CPU,也可能正在等待处理器。BLOCKED特指等待进入synchronized使用的监视器锁,不能把所有等待都叫成这个状态。[^state]
几个常见动作可以连起来记:
| 动作 | 意思 | 注意点 |
|---|---|---|
sleep(...) | 暂停当前线程一段时间 | 不释放已经持有的监视器锁 |
join() | 当前线程等待另一个线程结束 | 不会启动被等待的线程 |
wait() | 等待某个对象上的条件变化 | 必须持有该对象监视器;只释放这个对象的监视器锁 |
interrupt() | 发出协作式中断请求 | 不是强行杀掉线程 |
wait()返回前需要重新拿到锁;notifyAll()也只是唤醒竞争者,不会立刻把锁交出去。手写等待条件要用while重新检查,而不是只判断一次,因为存在虚假唤醒等情况。[^wait]
三、数据为什么会改错:一行代码,不一定是一步操作
回到批量导入:如果 4 个工作线程都直接修改同一个“成功条数”变量,就可能少统计。下面的count++很短,但逻辑上包含:读取旧值 → 加一 → 写回。
两个线程都读到0,分别算出1,再先后写回,最后就可能得到1,而不是2。
Java 内存模型(JMM)规定线程之间的读写何时必须可见、哪些顺序必须被遵守,不是简单的“主内存复制一份”模型。入门先抓住三个词:原子性是整体不可被交错破坏;可见性是修改能按同步规则被其他线程看到;有序性是必须遵守的操作先后约束。[^jmm]
用一把锁,保护完整规则
以“有库存才能扣减”为例,检查和扣减必须一起保护:
classInventory{privateintstock=10;publicsynchronizedbooleandeduct(intquantity){if(quantity<=0)thrownewIllegalArgumentException("数量必须大于零");if(stock<quantity)returnfalse;stock-=quantity;returntrue;}publicsynchronizedintremaining(){returnstock;}}这里锁的是同一个Inventory实例,不是整个 Java 程序。参与并发访问的路径都要遵守同一套同步规则;不同实例各自加锁,不能保护彼此。多台服务之间更不能靠这一把 JVM 内的锁解决库存一致性。[^jmm]
几种工具分别解决什么
| 工具 | 适合解决 | 不要误解 |
|---|---|---|
synchronized | 保护一段共享数据操作 | 同一把锁才会互斥,也提供相应可见性保证 |
ReentrantLock | 需要可中断、可超时的锁获取 | 成功加锁后,在finally中解锁 |
volatile | 状态标记等读写的可见性与顺序保证 | 不能让count++成为原子操作 |
AtomicInteger | 单个整数的原子增减等操作 | 不自动保护多个变量之间的业务约束 |
导入成功后,单纯计数可以用counter.incrementAndGet();也可以让每个批次返回自己的结果,最后由一个协调线程相加,减少共享写入。CAS 可理解为“当前值仍是预期值时,才替换成新值”的原子比较更新;原子类也提供这类操作,但不要把所有原子方法都想成必须手写循环。[^atomic]
ConcurrentHashMap同样不是万能保险:先get()再put()仍是两步;按一个键更新可考虑compute()、merge()等原子组合操作,其回调应保持短小。[^map]
四、线程池怎么接活:先创建核心线程,再排队,再扩容
批量导入会产生很多批次,订单系统也会不断产生通知任务。线程池可以复用工作线程,也能限制任务占用的资源。可以把它看成一个有固定接单规则的工作组,而不是“提交多少任务,就启动多少线程”。[^pool]
下面的池只用于说明规则:核心线程 2,最大线程 4,等待队列最多 3 个任务。
ThreadPoolExecutorpool=newThreadPoolExecutor(2,4,30,TimeUnit.SECONDS,newArrayBlockingQueue<>(3),Executors.defaultThreadFactory(),newThreadPoolExecutor.AbortPolicy());假设池起初为空,连续提交任务,并且前面的任务都没有完成:第 1、2 个创建工作线程;第 3~5 个排队;第 6、7 个触发扩容;第 8 个被拒绝。因此,提交顺序不等于执行顺序。
几个参数里,corePoolSize是核心线程数,maximumPoolSize是最大线程数,workQueue保存待执行任务;keepAliveTime与时间单位控制多余空闲线程的保留时间;ThreadFactory负责创建线程;最后一个参数决定怎样拒绝任务。核心线程默认按需创建,不是构造线程池时立即全部启动。[^pool]
示例选择AbortPolicy,会抛出RejectedExecutionException;应用要将它转成明确的“系统忙”或任务失败,不能假装提交成功。CallerRunsPolicy在池未关闭时会让提交者自己执行,可能拖慢请求线程;关闭后会丢弃任务。[^reject]
newFixedThreadPool()使用无界队列,积压可能耗尽内存;newCachedThreadPool()允许线程数持续增长,也需要评估上限风险。它们不是不能用,而是默认配置未必符合业务容量要求。[^executors]
业务中通常复用由应用统一管理的池,不要每个请求新建一个。大文件导入和用户通知可以分开设置执行池与并发限额,避免长任务占满所有工作线程;但它们若共用数据库,仍会竞争数据库资源。池大小没有万能公式,要结合 CPU、等待耗时、数据库连接数、下游限流和压测结果决定。
五、任务怎么协作:导入要汇总,订单要拆分
导入场景:多个批次完成后,拿到各自结果
Runnable不返回结果;Callable<T>可以返回结果、抛出异常。提交批次任务后拿到的Future<T>,可以理解为“这一批的结果凭证”;get()用来拿结果,没完成时会等待。[^future]
例如,每个批次返回BatchResult(成功条数, 失败条数)。协调线程收集这些结果,再相加得到整个文件的处理情况。这样工作线程不用同时修改同一个普通计数器。
ExecutorCompletionService可以按任务完成的先后取结果,不必一直守着最先提交、却还没完成的那一批。下面是完整示例的核心片段,省略了批次读取方法和超时处理:[^completion]
// pool 是导入专用池;source 是逐条读取的数据源。CompletionService<BatchResult>completed=newExecutorCompletionService<>(pool);intpending=0;intmaxInFlight=8;// 已提交、但尚未回收结果的批次上限;仅为示例while(source.hasNext()){if(pending==maxInFlight){merge(completed.take().get());// 仅协调线程汇总成功与失败条数pending--;}List<Row>batch=readNextBatch(source,500);completed.submit(()->validateAndSave(batch));pending++;}while(pending>0){merge(completed.take().get());pending--;}把这段代码读成:**先读一批,交给线程池;待回收批次太多,就先收一个结果;文件读完后,收齐剩余结果。**批次完成顺序不一定与文件顺序相同;需要还原顺序时,可以在结果中保留批次号或行号。
这里有两道边界:线程池限制同时工作的线程数,maxInFlight限制提交端不断囤积任务与结果。完整的BatchImportDemo.java使用模拟流式数据源,不会先创建一个装满 10 万条数据的大列表,也不会创建 10 万个Future;它不包含真实 Excel 解析或数据库写入。
订单场景:前置状态完成后,分别处理后续任务
CompletableFuture适合表达“先做什么、哪些可以分开做、全部完成后怎么办”。看下面这个核心片段:进入这里时,订单已经存在,支付成功状态也已经提交。[^cf]
// 示例片段:pool 由应用统一管理,三个业务方法代表互不依赖的后续动作。CompletableFuture<Void>warehouse=CompletableFuture.runAsync(()->createDeliveryTask(orderId),pool);// 安排仓库备货,不是已经发货CompletableFuture<Void>notice=CompletableFuture.runAsync(()->sendPaymentSuccessNotice(orderId),pool);CompletableFuture<Void>statistics=CompletableFuture.runAsync(()->updateOrderStatistics(orderId),pool);CompletableFuture<Void>all=CompletableFuture.allOf(warehouse,notice,statistics);// 返回一个可继续观察成功与失败的结果;不是“提交后不管”。returnall.whenComplete((unused,error)->{if(error!=null){log.error("订单后续任务存在失败,orderId="+orderId,error);}});runAsync()提交不返回业务结果的任务;supplyAsync()用于需要结果的任务。allOf()汇总完成状态,不会把三个操作变成一笔事务:某个通知失败,已经完成的备货安排不会自动撤销。线程池满时,提交本身还可能抛出拒绝异常,完整示例对这一点也做了处理。[cf][pool]
真正存在依赖的动作要连接起来。例如,拿到真实发货结果并保存已发货状态后,再通知用户,可以用thenRunAsync()表达“上一步正常完成才执行下一步”。thenApply()转换结果,thenCompose()连接返回另一份 Future 的后续任务,thenCombine()汇总两个结果。[^cf]
**不要在用户请求里立刻all.join(),又让用户等待全部后续动作。**是否等待,由接口语义决定;独立命令行演示为了观察输出,可以在末尾等待。join()仍会阻塞;orTimeout()让 Future 超时完成,不代表底层数据库或 HTTP 调用自动停止。[^cf]
落地还要补一层可靠性:**线程池负责“谁来执行”,不负责“进程重启后任务还能找回来”。**对不能丢的发货或通知任务,可把业务状态与待处理事件在同一数据库事务中保存,再由后台读取、执行和重试;采用 MQ 时也要处理业务提交与消息投递之间的一致性。重试可能重复执行,要用订单号与任务类型等业务标识防止重复安排发货,这就是幂等。[^outbox]
其他协作工具可以顺着场景理解:CountDownLatch等一组动作结束;Semaphore限制同时调用仓库接口的任务数量;BlockingQueue让读取方把批次交给处理方。[^coordination]
六、落地前再检查:退出、异常和容量
**线程必须能退出。**当sleep()等可中断等待抛出InterruptedException时,要么继续向上抛,要么恢复中断标记并结束当前任务,不能吞掉异常后无限重试:[^thread]
try{Thread.sleep(100);}catch(InterruptedExceptione){Thread.currentThread().interrupt();return;}线程池必须有人管理。shutdown()停止接收新任务,但允许已提交任务继续执行;它本身不等待全部结束。需要时配合awaitTermination(),超时再尝试shutdownNow();后者也只是尝试中断,不保证任务立刻停止。[^lifecycle]
**问题必须能被看见。**用submit()后不能完全不管返回的Future,否则任务异常可能无人检查。排查时同时观察活跃线程数、队列长度、拒绝次数、耗时与线程栈,而不是只看线程总数。
ThreadLocal保存线程各自的数据,不是给共享对象加锁。在线程池任务中使用请求上下文时,要在finally中remove(),避免复用线程带入上一次任务的数据。[^threadlocal]
另外,两个线程相反顺序获取两把锁可能形成死锁;小线程池中的任务又向同一池提交子任务并阻塞等待,也可能互相耗尽执行机会。设计时尽量统一锁顺序,减少持锁范围,避免在池内同步等待同池子任务。
**Java 21 的虚拟线程放在哪里?**它适合大量以等待为主的任务,不会让 CPU 计算自动变快。虚拟线程通常一任务一线程,不需要池化复用;数据库连接数和下游并发仍要单独限制,锁和共享数据问题也不会消失。[^virtual]
多线程最终只有一条主线:先判断任务能不能并发,再保护共享数据,安排协作关系,最后给资源和等待设边界。
在自己的项目里试一试:运行examples/BatchImportDemo.java看“分批处理与结果汇总”,运行OrderAsyncDemo.java看“支付确认后的任务拆分”。再用CounterDemo.java观察少统计一次的问题,用PoolRoutingDemo.java理解任务为什么会被拒绝。
大数据导出流程可能会出现内存溢出,导出流程可以看这篇总结:
订单导出总是内存溢出、超时?从分批查询到异步任务,把大数据导出做稳