资讯详情

资讯详情

Java 多线程总结:线程、锁、线程池与异步协作

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]二、线程怎么开始、等待和结束先运行最小例子ThreadworkernewThread(()-{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{privateintstock10;publicsynchronizedbooleandeduct(intquantity){if(quantity0)thrownewIllegalArgumentException(数量必须大于零);if(stockquantity)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 个任务。ThreadPoolExecutorpoolnewThreadPoolExecutor(2,4,30,TimeUnit.SECONDS,newArrayBlockingQueue(3),Executors.defaultThreadFactory(),newThreadPoolExecutor.AbortPolicy());假设池起初为空连续提交任务并且前面的任务都没有完成第 1、2 个创建工作线程第 35 个排队第 6、7 个触发扩容第 8 个被拒绝。因此提交顺序不等于执行顺序。几个参数里corePoolSize是核心线程数maximumPoolSize是最大线程数workQueue保存待执行任务keepAliveTime与时间单位控制多余空闲线程的保留时间ThreadFactory负责创建线程最后一个参数决定怎样拒绝任务。核心线程默认按需创建不是构造线程池时立即全部启动。[^pool]示例选择AbortPolicy会抛出RejectedExecutionException应用要将它转成明确的“系统忙”或任务失败不能假装提交成功。CallerRunsPolicy在池未关闭时会让提交者自己执行可能拖慢请求线程关闭后会丢弃任务。[^reject]newFixedThreadPool()使用无界队列积压可能耗尽内存newCachedThreadPool()允许线程数持续增长也需要评估上限风险。它们不是不能用而是默认配置未必符合业务容量要求。[^executors]业务中通常复用由应用统一管理的池不要每个请求新建一个。大文件导入和用户通知可以分开设置执行池与并发限额避免长任务占满所有工作线程但它们若共用数据库仍会竞争数据库资源。池大小没有万能公式要结合 CPU、等待耗时、数据库连接数、下游限流和压测结果决定。五、任务怎么协作导入要汇总订单要拆分导入场景多个批次完成后拿到各自结果Runnable不返回结果CallableT可以返回结果、抛出异常。提交批次任务后拿到的FutureT可以理解为“这一批的结果凭证”get()用来拿结果没完成时会等待。[^future]例如每个批次返回BatchResult(成功条数, 失败条数)。协调线程收集这些结果再相加得到整个文件的处理情况。这样工作线程不用同时修改同一个普通计数器。ExecutorCompletionService可以按任务完成的先后取结果不必一直守着最先提交、却还没完成的那一批。下面是完整示例的核心片段省略了批次读取方法和超时处理[^completion]// pool 是导入专用池source 是逐条读取的数据源。CompletionServiceBatchResultcompletednewExecutorCompletionService(pool);intpending0;intmaxInFlight8;// 已提交、但尚未回收结果的批次上限仅为示例while(source.hasNext()){if(pendingmaxInFlight){merge(completed.take().get());// 仅协调线程汇总成功与失败条数pending--;}ListRowbatchreadNextBatch(source,500);completed.submit(()-validateAndSave(batch));pending;}while(pending0){merge(completed.take().get());pending--;}把这段代码读成**先读一批交给线程池待回收批次太多就先收一个结果文件读完后收齐剩余结果。**批次完成顺序不一定与文件顺序相同需要还原顺序时可以在结果中保留批次号或行号。这里有两道边界线程池限制同时工作的线程数maxInFlight限制提交端不断囤积任务与结果。完整的BatchImportDemo.java使用模拟流式数据源不会先创建一个装满 10 万条数据的大列表也不会创建 10 万个Future它不包含真实 Excel 解析或数据库写入。订单场景前置状态完成后分别处理后续任务CompletableFuture适合表达“先做什么、哪些可以分开做、全部完成后怎么办”。看下面这个核心片段进入这里时订单已经存在支付成功状态也已经提交。[^cf]// 示例片段pool 由应用统一管理三个业务方法代表互不依赖的后续动作。CompletableFutureVoidwarehouseCompletableFuture.runAsync(()-createDeliveryTask(orderId),pool);// 安排仓库备货不是已经发货CompletableFutureVoidnoticeCompletableFuture.runAsync(()-sendPaymentSuccessNotice(orderId),pool);CompletableFutureVoidstatisticsCompletableFuture.runAsync(()-updateOrderStatistics(orderId),pool);CompletableFutureVoidallCompletableFuture.allOf(warehouse,notice,statistics);// 返回一个可继续观察成功与失败的结果不是“提交后不管”。returnall.whenComplete((unused,error)-{if(error!null){log.error(订单后续任务存在失败orderIdorderId,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理解任务为什么会被拒绝。大数据导出流程可能会出现内存溢出导出流程可以看这篇总结订单导出总是内存溢出、超时从分批查询到异步任务把大数据导出做稳
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →