ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

从线程池参数到分布式选型:生产级异步调度系统避坑指南

从线程池参数到分布式选型:生产级异步调度系统避坑指南 如果你在一个有点流量的后端团队待过大概率对“跑任务”这件事有一肚子话要说。定时对账、异步通知、数据补偿、批量导出看起来都是丢一个线程去执行的事可真到了生产环境线程就像脱缰的野马几分钟就能把服务搞成PPT。我自己做后端这些年每次聊到异步任务团队里都会蹦出一个词ax调度。这个“ax”最早是我们内部对Async eXecution的缩写后来叫顺口了就成了整个异步执行模块的代名词。今天我就把这套ax调度从里到外拆一遍从线程池参数到分布式选型把那些文档里不会写的坑一并翻出来。我下面聊的不是PPT是从生产线上抠出来的经验适合正在做异步任务、定时任务或者自研调度器的后端开发也适合刚接触线程池概念、想搞明白调度原理的人。1. 从“跑腿小哥”说起ax调度到底调度了什么1.1 一句话说清ax调度的职责ax调度最直白的理解它是任务和执行资源之间的调度中枢。你可以把它想象成外卖平台的派单系统。骑手是线程池里的线程订单是业务提交的任务调度中心负责决定谁去送、送几趟、送失败了怎么办、超时了要不要重新派给其他人。程序里的ax调度做的是同一件事接收业务方交来的任务放入队列交给合适的线程执行跟踪每一个任务的状态处理超时和失败按策略重试。很多刚踏入后端的同学会有个疑问任务不就是“丢个线程去跑”吗为什么要单独做一个调度组件因为“丢个线程”这个动作在生产环境里撑不过高并发那一关。直接new Thread的代码上线之后通常只有三种下场第一种是线程数量爆炸每个请求都开一个线程几百个并发就能把CPU和内存拉爆第二种是任务静默失败线程抛异常没人接住日志可能都没打一条第三种是高峰期系统整体卡顿大量线程阻塞在锁或IO上连请求处理的线程都被拖垮。所以ax调度要做的就是把“任务”和“执行资源”彻底分开。业务方只需要描述“我要做什么”不用关心“谁来做、什么时候做、做失败了怎么办”。这种分离看起来只是多了一层抽象但它带来的收益是实打实的线程被复用资源消耗可控任务有了统一的队列缓冲流量洪峰不会直接打在执行资源上整个执行过程有状态记录失败、超时、重试都能被观测到。1.2 为什么非要有一层“调度”而不是直接new Thread直接new Thread的问题并不只在资源层面更在于管理层面的失控。你可以想象一下如果公司里每个部门都自己拉网线、自己买交换机网络一定乱成一锅粥统一接入机房让专业的人做专业的网络规划才能有序。ax调度就是把散布在业务代码里的线程管理收拢到一个组件里让线程的创建、复用、销毁、排队、拒绝全部有章法。调度层在编程模型上也做了很关键的一件事把同步调用和异步执行解耦。业务方调用submit提交任务后立刻就能得到“已受理”的返回不需要傻等任务执行完。真正的执行被放到一个受控的环境里有超时上限、有重试次数、有状态回写。这一下就把“调用方”和“执行方”之间的耦合解掉了。调用方不用因为某个任务耗时过长而被迫延长自身请求也没有必要在自己所在线程里处理重试和异常。这一层的价值在任务越来越多、团队越来越大的时候会体现得更明显。今天你只是导个Excel明天可能是定时给几千个用户推送消息后天可能是跑一个凌晨的全量数据对账。如果没有调度层这些需求会一点一点写进各自的业务模块最后的结局就是一堆互相独立、参数不一、故障案例都不一样的线程管理代码。把调度逻辑统一抽出来等于提前给系统上了一道安全阀。2. 核心设计拆解调度器的三块基石2.1 队列任务在哪儿排队等队列是整个ax调度最容易被忽视的组件但往往事故都出在这里。任务提交之后不是立即就能被线程执行当线程都在忙的时候任务就得找个地方暂存。这个暂存区就是任务队列。做一个不恰当的比喻队列相当于餐厅的出餐口厨师炒好一道菜就端走一道炒不过来的菜先放在出餐口排着。如果出餐口无限大高峰期可以一直往里面堆但菜会放坏如果出餐口太小高峰期就会直接拒客。落在程序里就是内存被堆满和任务被丢弃两种风险。具体到ThreadPoolExecutor队列类型直接决定了系统的呼吸节奏。LinkedBlockingQueue如果不指定容量默认是Integer.MAX_VALUE等于无界任务可以无限堆积最坏情况下内存被打爆ArrayBlockingQueue是有界队列容量可以按业务峰值估算满了之后触发拒绝策略SynchronousQueue比较特殊它不存任务提交线程必须直接移交给工作线程没有缓冲适合执行特别快、任务量小的场景PriorityBlockingQueue支持按优先级出队适合有紧急任务的场景但要小心优先级反转——高优先级任务不断插队低优先级任务可能永远跑不到。队列容量怎么定我习惯先算峰值吞吐和平均执行耗时。举个例子某个导出任务平均执行300毫秒核心线程数4那么理想的处理速度大约是每秒13个任务。如果业务峰值每秒要接收50个任务每秒就有37个任务需要排队。一个能撑住30秒的缓冲队列容量至少要按1100来估算。当然这只是起步值真正上线后要监控QueueSize根据观察再动态调整。这里我还要说一个观点无界队列不是说绝对不能用而是它不能作为默认选项。它会让系统失去“拒收”的能力任何一次任务爆发都可能被无限吸收直到把堆打爆。有界队列加拒绝策略才是真正把控制权握在自己手里。对大部分业务系统来说宁可高峰期丢几个任务然后告警重跑也不能让整个进程跟着陪葬。2.2 线程池谁来干活队列解决“等待”线程池解决“执行”。线程池的核心设计思路是线程复用避免频繁创建销毁线程的开销。把线程池理解成一个骑手团队核心骑手是固定员工平时一直上班当单量暴增核心骑手忙不过来可以临时招兼职骑手但兼职骑手不能无限招而且一段时间没活干就得走人订单更多到连兼职都忙不过来就只能把订单放到出餐口排队。落实到ThreadPoolExecutor四个关键参数是核心线程数、最大线程数、空闲存活时间、等待队列。核心线程数决定平时能在岗的线程数量最大线程数决定峰值时最多能撑到多少个线程keepAliveTime决定非核心线程空闲多久可以释放队列刚才已经聊过。需要说明的是核心线程数不是越大越好。线程越多CPU上下文切换就越频繁反而不如少而精。粗略的起步建议是CPU密集型任务用Ncpu1IO密集型任务用2*Ncpu左右Ncpu取Runtime.getRuntime().availableProcessors()的结果。但这只是入门公式真实系统的线程数要靠压测调整看吞吐量和响应时间的变化曲线去收敛。keepAliveTime这个参数常常被忽略。设得太短比如10秒高峰期刚招进来的兼职线程还没帮上多少忙就被释放了下个峰值又要重新创建反复横跳反而增加开销设得太长比如一两个小时空闲线程会持续占用内存和栈空间。常见生产值在60到120秒之间具体看任务的波动频率。如果业务每天只有固定的几个高峰就按高峰间隙来定。还有一个很容易踩的细节线程工厂一定要给线程起名字。默认的线程池创建出来的线程叫pool-3-thread-1线上用jstack看现场时你根本不知道它是哪个业务模块的排查效率极其低下。命名成ax-worker-1、ax-worker-2这样的格式出问题一眼就能锁定是哪一批线程在搞事。这一步代价极小收益却极大。2.3 拒绝策略兜底的尊严当队列满了线程池也满了新提交的任务就会触发拒绝策略。很多人在选型阶段根本不看拒绝策略默认的AbortPolicy直接抛RejectedExecutionException调用方如果没处理任务就无声无息地丢了连挽留的机会都没有。四大内置策略各有取舍我简单梳理一下。AbortPolicy会直接抛异常适合对数据准确性要求极高、丢一个任务都是事故的场景异常可以让上层感知并做人工补偿CallerRunsPolicy会让提交任务的线程自己执行这个任务执行完再去处理别的事等于反向施加了一个自然限流——Tomcat线程被占住后续请求自然变慢系统就不会瞬间被打垮DiscardPolicy直接静默丢弃就我个人而言非常不建议因为丢任务没有任何记录DiscardOldestPolicy会把队头最老的任务丢掉给新任务腾位置但被丢的可能是更重要的老任务使用前要慎重。我生产环境的默认选择是CallerRunsPolicy同时加一个计数器每次触发拒绝就记录并定时上报。为什么选它因为它既不会丢任务又能天然限流。真正的设计思路是任务先保住系统也别死哪怕干得慢一点。如果你所在业务连被占用提交线程都不允许那就用AbortPolicy但一定要配告警让运维收到一个“拒绝率超过x%”的提示。拒绝策略不是摆设是系统濒临崩溃时最后一道闸门。3. 实操过程5步搭建一个能上生产的ax调度组件3.1 第一步定义任务状态机任务总不能一直是一个抽象的对象它在调度系统里必须有一个明确的状态否则没人知道它到底跑成什么样。我把状态划分成六个PENDING表示已提交、还没被线程拿走RUNNING表示正在执行SUCCESS表示正常完成FAILED表示执行抛异常TIMEOUT表示超时RETRY_WAIT表示失败后处于等待重试的阶段。定义状态机的好处非常明显每个阶段都有明确的入口和出口出了任何问题都可以查状态、可以人工干预、可以补偿。这套状态机我建议用枚举来建模不要用魔数或字符串否则哪天把0写成1排查起来会想骂人。枚举配合一个简单的状态流转表提交时置为PENDING开始执行置为RUNNING正常返回置为SUCCESS异常捕获置为FAILED超时置为TIMEOUT需要重试再置为RETRY_WAIT。状态变更的同时要记录时间戳方便计算每个环节耗时这也是后面监控指标的基础。public enum TaskStatus { PENDING, RUNNING, SUCCESS, FAILED, TIMEOUT, RETRY_WAIT }3.2 第二步定义任务模型与调度接口任务模型要给每个任务一个全局唯一的taskId同时保存业务类型、业务数据、超时时间、最大重试次数、当前重试次数、下次重试时间、当前状态。其中taskId非常关键它不仅是查找任务的凭证还是去重和幂等的依据。业务数据payload可以用JSON字符串保存方便跨模块传递。超时时间和重试次数要放在任务级别而不是全局统一配置因为导出Excel的任务和推送通知的任务对超时的容忍度完全不同。调度器的对外接口不需要复杂核心就一个submit方法。submit做的事情也不多先查taskId是否已经存在存在就返回“重复提交”的结果不存在则落库并置为PENDING然后用CompletableFuture把执行动作交给线程池立刻返回受理结果。这里有一个很多人忽略的点本地业务代码里只要调用submit就能异步执行但任务落地持久化非常重要。如果只是一个内存任务服务一重启队列里的任务就全没了真正生产级的调度器一定要有持久化至少把任务表落到数据库重启后可以恢复未完成任务。public class Task { private String taskId; private String bizType; private String payload; private int timeoutSeconds; private int maxRetry; private int currentRetry; private long nextRetryTime; private TaskStatus status; } public class AxScheduler { private final ThreadPoolExecutor executor; public ScheduleResult submit(Task task) { if (taskRepository.exists(task.getTaskId())) { return ScheduleResult.duplicated(); } task.setStatus(TaskStatus.PENDING); taskRepository.save(task); CompletableFuture.runAsync(() - execute(task), executor); return ScheduleResult.accepted(); } }3.3 第三步构建可监控的线程池实例接下来是核心执行器的构建。我给出的参数是一个相对稳的生产起点结合前面提到的计算思路按需调整。核心线程4最大8队列用ArrayBlockingQueue容量10000线程工厂命名为ax-worker拒绝策略用CallerRunsPolicy。这组参数适合中等流量、任务单次执行耗时在几百毫秒以内的场景。线程池构建之后一定要把监控指标暴露出来。最常用的三个是并发活跃线程数、队列积压数量、完成任务总数。ThreadPoolExecutor本身提供了getActiveCount、getQueue().size、getCompletedTaskCount等方法定时上报到监控系统即可。我见过不少团队线程池配了但从不看指标直到线上出问题才去查那已经太晚了。ThreadPoolExecutor executor new ThreadPoolExecutor( 4, 8, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue(10000), new ThreadFactory() { private final AtomicInteger counter new AtomicInteger(1); Override public Thread newThread(Runnable r) { Thread t new Thread(r, ax-worker- counter.getAndIncrement()); t.setDaemon(false); return t; } }, new ThreadPoolExecutor.CallerRunsPolicy() );3.4 第四步超时、重试与回调处理任务光会执行不行还得会处理超时和重试。超时控制的常见做法是把执行丢给Future再用带超时参数的get方法等待结果。如果超时调用future.cancel(true)发送中断信号。这里必须强调一点cancel(true)只是告诉线程“你可以停了”如果任务内部不响应中断线程依然会继续跑下去。所以任务代码里最好定期检查Thread.currentThread().isInterrupted()实在不行宁可让这个线程跑完也要确保任务的最终状态被正确回写不然后续补偿逻辑会陷入混乱。重试逻辑我推荐指数退避。第一次失败等1秒第二次等2秒第三次等4秒这样既给了系统恢复时间又不会在故障期间疯狂重试。重试次数必须有限制超过最大次数就置为FAILED并告警。回调这里的设计也比较讲究onSuccess、onFailure、onTimeout分别处理三种结局但回调本身也可能失败所以回调里不能抛异常要用try/catch包住并单独记录日志。private void execute(Task task) { task.setStatus(TaskStatus.RUNNING); long start System.currentTimeMillis(); try { businessHandler.execute(task); task.setStatus(TaskStatus.SUCCESS); } catch (TimeoutException e) { task.setStatus(TaskStatus.TIMEOUT); } catch (Exception e) { task.setStatus(TaskStatus.FAILED); } if (task.needRetry()) { task.setStatus(TaskStatus.RETRY_WAIT); long delay 1000L * (long) Math.pow(2, task.getCurrentRetry()); task.setNextRetryTime(System.currentTimeMillis() delay); taskRepository.save(task); } }3.5 第五步接入业务时的三个硬件级注意点这一步是实操里最容易踩雷的地方。第一个注意点不要在数据库事务里直接提交异步任务。事务没提交时数据在别的事务里还看不到任务线程一旦抢先执行可能查不到刚插入的数据。正确的做法是先提交事务再提交任务或者确保任务能容忍一定延迟再来消费数据。第二个注意点线程上下文要传递。异步任务不会自动继承主线程的ThreadLocal链路追踪的traceId、登录用户信息都会丢。推荐的方案是封装一个Runnable装饰器在提交前把需要的ThreadLocal快照下来执行前还原执行完再清理。如果用的中间件多可以直接用阿里开源的TransmittableThreadLocal。第三个注意点任务方法内必须写finally状态的回写不能依赖业务代码是否成功。哪怕业务异常了finally里也要想尽办法把状态写出去这是调度器自愈能力的基础。另外要补一个监控点任务执行时长。单次任务跑太久往往比失败更可怕因为线程被它占着后面的任务全得排队。给每次执行记录耗时超过阈值就告警长期性能瓶颈就会浮出水面。4. 常见问题与排查技巧实录4.1 任务“卡死”既不执行也不报错怎么查这一类问题的排查思路有固定的路线。先看线程池状态看ActiveCount是不是一直等于最大线程数队列size是不是只涨不降其次看线程都在干嘛jstack之后搜ax-worker线程看它们阻塞在什么位置。常见的情况有三种一是业务方法里查数据库而数据库连接池耗尽线程全部阻塞在获取连接上二是任务里互相等待锁造成了死锁三是任务代码里有不可中断的循环或者第三方调用没有设置超时线程被外部接口拖死。每个情况都能在jstack的线程栈里看到痕迹所以前面强烈建议线程命名排查时真的救命。4.2 任务重复提交幂等防线怎么做业务方调用submit时常常会在一段重试逻辑里重复提交同一个任务如果调度器不做幂等控制同一个taskId可能被多个线程同时执行出现重复扣款、重复发短信这类事故。幂等防线做到两层比较稳第一层是在代码入口查taskId是否已存在已经存在就拒绝第二层是在数据库层面加唯一索引让并发重复插入直接报错。如果有多个服务实例第一层内存判断不可靠必须依赖第二层或Redis分布式锁。锁的key用业务类型加业务ID锁的有效期要大于任务最大执行时间否则任务没跑完锁就过期照样重复执行。4.3 内存悄悄上涨无界队列的锅这类事故我自己就碰到过一回。任务提交速度稍微大于消费速度无界队列开始悄悄堆积一天两天不察觉到了三四天内存直接OOM。用有界队列之后队列满了会触发拒绝策略虽然系统会丢弃任务但至少能暴露问题。这里给三个有效手段队列容量设上限这是第一道防线QueueSize加监控和告警设置一个阈值比如8000就报警这是第二道防线任务入队前做一次预估任务总量超过预期宁可先拒绝也不要无脑塞进队列。4.4 shutdown与shutdownNow优雅停机比想象中复杂服务发布的时候线程池怎么关是很有讲究的。直接调用shutdownNow执行中的任务会被强制中断队列里还没执行的任务可能会直接丢失只调用shutdown服务又可能一直被长任务拖住等不到发布完成的信号。我的做法是先shutdown然后awaitTermination等待比如等60秒如果还在执行说明有任务确实太长再调用shutdownNow并把返回的未执行任务列表记录下来后续通过补偿机制重新入队。这套流程写成一个方法发布脚本只管调用不需要人工介入。这里有一个经验不要写一个while循环不停awaitTermination如果线程池里的任务真的饿死了那个循环会把发布流程拖到天荒地老。设置一个合理的上限时间超时就进入强制中断分支发布拥有确定性。5. 进阶方向单机调度到分布式调度5.1 单机调度撑不住的场景单机调度处理核心问题已经足够但一旦服务一上多实例问题就来了。多个实例各自运行一套调度器同一批任务会被重复消费。比如定时对账任务两个实例同时去扫同一个任务表就会跑出两份对账单。另外单机调度意味着任务全在进程内存里进程一挂未完成任务就地蒸发。所以当你的服务开始多实例部署并且任务有严格的可靠性要求就需要考虑分布式调度了。分布式调度的核心思想是把“任务所有权”从进程级别提升到集群级别每个任务在同一时刻只能被一个实例认领并执行。认领的机制可以基于数据库、基于Redis、基于ZooKeeper也可以直接用一个成熟的调度框架。5.2 数据库分布式锁成本最低的认领方式小型项目我推荐先用数据库实现任务认领方案简单到只需要一个UPDATE语句将某个PENDING任务的状态改成RUNNING并打上owner标记时带条件更新只有更新成功的实例才拥有任务。用到了当前任务状态作为条件天然防止多个实例同时拿到同一个任务。UPDATE task_schedule SET owner instance-01, status RUNNING, start_time NOW() WHERE id 123 AND status PENDING这个方案的问题在于如果任务时间太长数据库连接被长时间占用连接池很容易被拖垮。另外任务结束时释放状态也要带条件更新防止把别人刚认领的同ID任务误操作。更稳一点可以用乐观锁版本号每次更新带上version字段version变了就放弃。等任务量到了几万、几十万这个级别数据库锁的方案就开始吃力那时候再升级到Redis锁或者专业框架也不迟。5.3 消息队列加定时任务异步调度的另一种形态很多系统的异步调度并不一定要有一个独立的调度中心用消息队列也能实现很优雅。定时扫表找出到期的任务发一条延迟消息到消息队列消费者收到消息后执行任务执行失败再发一条延迟消息做重试。RocketMQ原生支持延迟消息Kafka没有内置延迟消息需要自己做时间轮或者用MQ插件。这种架构的好处是天然削峰和解耦消息的堆积给大家一个共同的缓冲空间。坏处也很明显消息可能重复消费者要做幂等消息队列本身也可能打满需要监控消费Lag。它和主动式调度器不是竞争关系更适合像“用户注册后10分钟发通知”“订单超时30分钟自动关闭”这类延迟触发场景而主动式调度器更适合定时任务、批处理任务。5.4 框架选型Quartz、XXL-Job、PowerJob与Elastic-Job怎么选如果团队不想自研直接用现成框架是性价比最高的路径。我把常见的四个框架放在一起做个比较。框架调度方式适合规模运维成本特点Quartz数据库锁/集群部署中小规模中等老牌成熟功能简洁无控制台XXL-Job调度中心执行器中大规模较低可视化界面完善接入简单PowerJob调度中心Worker中大规模中等支持工作流、MapReduce性能高Elastic-JobZooKeeper分片中大规模较高分片能力强适合数据切片场景我的建议是小团队、业务规模可控优先考虑XXL-Job它的调度中心和执行器设计容易理解文档多遇到问题也容易找到答案需要对任务编排和复杂工作流或者想获得更好的扩展性可以看PowerJobElastic-Job最擅长的是按分片把一堆数据拆给多台机器并行跑比如全量数据清洗、按用户ID分片处理这种场景它很在行。框架只是工具关键是它能不能和你们团队的运维能力、业务模型匹配上。我自己的经验是调度系统最怕的不是并发高而是状态乱。线程池炸了监控能看到任务状态乱了业务就全瞎了。所以不管你是自研还是用现成框架第一优先级永远是状态机清晰、监控告警完整、人工补偿入口齐全。先单机跑稳再上分布式这个顺序千万别反。希望这些踩过的坑能让你少走一点弯路后面真遇到问题至少知道从哪里下手。
返回列表