ARTICLE DETAIL

资讯详情

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

用MongoDB构建高可靠任务调度系统:原子领取、状态机与索引实践

用MongoDB构建高可靠任务调度系统:原子领取、状态机与索引实践 做任务调度系统的人多半在存储选型上纠结过。我接手过一个仿真任务调度平台日增量几百万条任务最初用Redis List当队列、MySQL存任务状态队列和状态分离之后同步问题、消息丢失、重复消费轮着来后来索性把任务状态和待执行队列统一放进MongoDB一套文档模型同时搞定存储和调度吞吐量反而比原来两套系统加起来还稳。这篇文章就把我在类似场景里用MongoDB做任务调度的核心套路、高级特性和踩过的坑一次性说清楚。文章适合三类人看准备用MongoDB做任务队列但还没定方案的已经在用但遇到并发冲突、性能瓶颈的以及想搞清楚MongoDB在调度场景里边界在哪、什么情况该换Kafka或RabbitMQ的。我会结合具体代码、索引设计和生产环境里的实测表现来讲不绕弯子。1. 先聊清楚任务调度对存储层到底要什么MongoDB为什么能接住很多团队一上来就纠结“选哪个消息队列”其实任务调度和消息投递是两件事。任务调度关注的是任务的生命周期管理——谁在什么时候执行、执行到哪一步了、失败了要不要重试、重试几次消息队列关注的是消息的可靠传递——投递了、确认了、不会丢。两者有交集但不能画等号。1.1 存储层要同时满足的六个能力我把任务调度场景对存储的要求拆成六点缺一个后面都会出事原子领取多个worker并发扫描任务表时同一条任务不能被两个worker同时拿走。这是最硬的需求后面会展开讲。状态流转任务得有一个明确的状态机pending→running→success/failed每一步都要有记录方便定位问题。延迟可用任务不一定立刻执行可能是5分钟后、明天凌晨3点存储层要能支持这种“到点再可见”的语义。优先级调度紧急任务要能插队存储层要支持按优先级排序取任务。故障恢复worker执行到一半宕机任务卡在running状态存储层要能识别并重新调度。吞吐量批量任务入库、批量回写结果的时候存储层不能成为瓶颈。1.2 对比一轮MySQL、Redis、MongoDB谁更合适我在这三个里都踩过坑直接说结论性质的对比存储方案原子领取状态管理延迟任务优先级故障恢复吞吐量MySQL行锁事务可以做但锁竞争和连接池压力大强但表结构要预先设计死支持轮询扫字段支持排序即可靠业务代码写入量上来后要分库分表Redis原子性极好LPOP/BRPOP弱状态和队列分离丢数据风险可以用ZSET但逻辑要自己拼支持多个List要自己实现数据易丢需要额外持久化方案单机吞吐很高集群运维成本高MongoDBfindOneAndUpdate天然原子文档模型灵活一个文档装下所有状态和结果支持scheduledAt字段索引支持复合索引排序通过lease机制实现写入吞吐强水平扩展方便我的结论是如果你的调度系统是中等规模——日任务量几百万到几千万、峰值TPS在几千这个量级MongoDB的性价比很高。它把队列和数据库合二为一少了一套系统的同步和运维成本。但如果你的业务是极致的FIFO、每天上亿消息、对延迟要求到毫秒级那还是老老实实用Kafka或者RabbitMQMongoDB不是干这个的。2. 任务队列的核心模型原子领取、状态机与优先级调度怎么设计这一节是整篇文章的地基。我见过太多人把MongoDB当作“能存JSON的MySQL”然后用关系型思维去建表、用两步操作去更新结果并发一起来就乱套。下面直接给出一套经过生产验证的建模方案。2.1 原子领取必须用findOneAndUpdate禁止findupdate两步走先看最典型的错误写法// 错误示范两个操作之间任务可能已经被别的worker领走 const task db.tasks.findOne({ state: pending }); db.tasks.updateOne( { _id: task._id }, { $set: { state: running, workerId: worker-01 } } );这段代码看起来没毛病但稍微有点并发量就会出问题。两个worker同时查到同一条pending任务然后先后把它改成running任务就被执行了两遍。在仿真任务调度这种场景里重复执行可能意味着重复扣资源、重复写结果后果很严重。正确做法是用findOneAndUpdate把“查询修改”合并成一个原子操作const task db.tasks.findOneAndUpdate( { state: pending, scheduledAt: { $lte: new Date() }, attempts: { $lt: 5 } }, { $set: { state: running, workerId: worker-01, leaseUntil: new Date(Date.now() 5 * 60 * 1000), startedAt: new Date() }, $inc: { attempts: 1 } }, { sort: { priority: -1, scheduledAt: 1 }, returnDocument: after } );这段代码里有几个细节要说明过滤条件直接写在查询里state: pending、scheduledAt小于当前时间、attempts小于5一次性把“可执行”“到点了”“还没超重试次数”三个条件全部卡死。修改操作是原子的MongoDB单个文档的操作天然原子findOneAndUpdate不会出现“查到了但没改成”的中间状态。多个worker同时调用只有一个人能成功把state从pending改成running其余人拿到null继续下一轮循环。sort里指定优先级和预约时间这样每次领取到的都是当前最优的一条任务而不是随机一条。returnDocument: after表示返回更新后的文档。不同驱动的写法略有差异Node.js驱动是returnDocumentJava驱动是returnDocument配合ReturnDocument.AFTERPython是return_documentReturnDocument.AFTER查一下自己用的驱动文档很容易搞混。attempts用$inc自增领取一次就加一次同时结合查询条件里的attempts: { $lt: 5 }天然实现了重试次数限制。2.2 任务文档的字段设计一个文档装下整个生命周期任务调度的数据模型我建议用“单文档全生命周期”的方式而不是把任务基本信息和执行状态拆到两张表。一个任务文档大致长这样{ _id: ObjectId(...), taskType: simulation, state: pending, priority: 5, scheduledAt: ISODate(2024-06-01T02:00:00Z), startedAt: null, finishedAt: null, leaseUntil: null, attempts: 0, maxAttempts: 5, lastError: , payload: { simulationId: sim_1024, params: { iterations: 10000, model: v3 } }, result: null, createdAt: ISODate(2024-05-31T10:00:00Z) }字段含义字段作用taskType任务类型方便按类型做不同处理逻辑state状态机核心字段pending/running/success/failed/cancelledpriority优先级数值越小越紧急排序时用sort: { priority: -1 }还是1取决于你定义的语义关键是统一scheduledAt首次可执行时间延迟任务和定时任务都靠它leaseUntil租约过期时间worker宕机后靠它恢复任务attempts / maxAttempts已重试次数 / 最大重试次数payload任务入参文档模型的优势正在于此不同任务类型可以有不同的payload结构不需要提前建表result执行结果直接塞回同一个文档查询、审计都方便状态机流转我建议限制在一张图里pending → running → successrunning → failedpending → cancelled。其中running状态如果leaseUntil过期了可以由调度进程重新置回pending这个机制下面专门讲。2.3 延迟任务和定时任务一个scheduledAt字段就够了但别用TTL索引实现延迟任务的需求很常见比如“订单支付后30分钟未支付自动关闭”“仿真结果出来后5分钟开始生成报表”。实现方式就是在文档里写死一个scheduledAt领取条件带上scheduledAt: { $lte: now }即可。没有到点的任务查询条件不命中自然不会被领走。定时任务比如每天晚上3点跑全量数据同步可以单独建一个cron_jobs集合存储调度配置{ _id: cron:daily_report, jobName: daily_report, cronExpression: 0 3 * * *, lastRunAt: ISODate(2024-05-30T03:00:00Z), nextRunAt: ISODate(2024-05-31T03:00:00Z), enabled: true }调度进程定时扫描nextRunAt到点后创建一个实际的任务文档插入tasks集合同时把nextRunAt推一天。这里要注意不要在cron_jobs里直接存“要执行的任务内容”而是让它只负责“触发”真正执行的内容放到tasks集合这样两边的职责清晰状态管理也不会乱。经常有人问我能不能用TTL索引实现延迟任务——也就是给文档加一个expireAt字段让MongoDB到期自动删除然后监听删除事件去触发业务。我明确不建议原因下面“误区”章节会详说TTL索引是后台清理数据用的它的过期删除时机不精确而且它删的是数据不是“触发信号”只有把业务逻辑建立在数据的消失上这本身就是个大坑。2.4 索引设计pending任务的检索要快历史任务不能拖后腿任务表的数据量是持续增长的索引设计直接决定吞吐量。我推荐的复合索引组合db.tasks.createIndex( { state: 1, scheduledAt: 1, priority: -1 }, { partialFilterExpression: { state: pending }, name: idx_pending_scan } );这里用了部分索引只索引state: pending的文档。因为worker领取任务时只关心pending状态索引里塞一堆running、success的文档纯属浪费空间和写入性能。这个索引能让“扫描下一批待执行任务”这个高频操作走索引覆盖性能非常可观。有的场景还需要按worker查“我正在跑哪些任务”那就再加一个db.tasks.createIndex( { workerId: 1, state: 1, leaseUntil: 1 }, { partialFilterExpression: { state: running } } );这个索引用于worker心跳上报、租约续期、以及故障恢复时扫描“哪些running任务其实已经超时了”。还有一点经验之谈不要为了“可能用到”的查询盲目建索引。任务场景里最常见的查询就那么几个索引控制在3到5个以内多了写入会明显变慢因为每次插入和更新都要同步维护索引树。3. 进阶玩法Change Streams、分布式锁与任务过期回收基础模型跑通之后会碰到几个“不上点手段就难受”的场景轮询太频繁怕把数据库查死多worker抢同一个周期任务会打架worker宕机留下的僵尸任务没人管。这三个问题分别可以用Change Streams、唯一索引锁、lease回收机制解决。3.1 用Change Streams替代高频轮询实时感知任务状态变化很多调度系统在任务状态变化后需要立刻通知外部系统比如“仿真任务跑完了把结果推给前端”。最朴素的做法是前端轮询每秒钟查一次tasks表查状态变了没。任务量小的时候没问题任务量大起来一秒一次的查询会把数据库拖垮。MongoDB 3.6开始提供Change Streams功能相当于数据库级别的“事件订阅”。任务集合里的文档发生insert、update、delete时会推送给订阅方。前提是MongoDB必须以副本集模式部署单节点实例不支持。典型用法const changeStream db.tasks.watch( [ { $match: { operationType: update, fullDocument.state: success } } ], { fullDocument: updateLookup } ); changeStream.on(change, (change) { notifyExternalSystem(change.fullDocument); });几个使用要点fullDocument: updateLookup这个配置很重要。默认情况下update操作通知里不包含完整文档只有一个变更描述你要拿到最新数据还得自己再查一遍。设置了updateLookup之后Change Streams会帮你把文档最新的完整内容拼进去。$match可以大大降低网络和业务压力只订阅自己关心的变更。Change Streams有resume机制如果你自己的服务重启了可以从上次的resumeToken继续消费不会丢事件。当然你得把token存下来一般是存到单独的集合里。不要用它替代消息队列来做可靠性投递。Change Streams的定位是“事件流通知”不是“消息队列”它没有ACK和重试概念业务侧要做幂等。实测下来Change Streams非常适合两个场景一个是“任务状态变化通知”另一个是“实时统计面板”。但注意如果变更频率极高每秒上万次Change Streams的推送压力也可能成为瓶颈这时候还是应该考虑引入真正的消息队列做削峰。3.2 基于唯一索引的分布式锁抢锁、续约、释放的一次完整实现任务调度经常需要“同一时间只有一个worker在跑某个定时批量任务”比如每天凌晨的全量数据清理。用数据库实现分布式锁MongoDB有一招非常好用利用唯一索引的约束力。建锁集合时给锁的名称建唯一索引db.locks.createIndex({ lockKey: 1 }, { unique: true });抢锁就是一次插入try { db.locks.insertOne({ lockKey: cron:daily_clean, owner: worker-01, acquiredAt: new Date(), expiresAt: new Date(Date.now() 30 * 1000) }); // 插入成功拿到锁 } catch (e) { if (e.code 11000) { // 唯一索引冲突锁被别人持有 } }续约就是把过期时间往后推同时校验owner是自己const result db.locks.findOneAndUpdate( { lockKey: cron:daily_clean, owner: worker-01, expiresAt: { $gt: new Date() } }, { $set: { expiresAt: new Date(Date.now() 30 * 1000) } } ); if (result) { // 续约成功 } else { // 锁已经过期或者不属于自己重新抢锁 }释放锁是删除同样要带owner条件db.locks.deleteOne({ lockKey: cron:daily_clean, owner: worker-01 });这套方案比Redis SETNX的好处是锁的元数据看得见摸得着可以排查是谁持锁、持锁多久了。缺点也很明显需要自己处理“持锁进程挂了锁没释放”的问题所以要带上expiresAt并养成续约的习惯。3.3 租约机制让卡死在running状态的任务自动复活分布式系统里worker宕机是常态。任务刚被领走、还没执行完worker进程突然挂了这一条任务如果没人管就永远卡在running状态。解决思路就是前面提到的leaseUntil字段。我把它叫做“租约”worker领取任务时设了一个保险丝——leaseUntil 当前时间 5分钟。如果worker在5分钟内没有处理完任务也没更新这个字段说明它可能出问题了。后台专门跑一个恢复进程或者一个定时任务扫描状态异常的running任务const expiredTasks db.tasks.find({ state: running, leaseUntil: { $lt: new Date() } }); for (const task of expiredTasks) { db.tasks.findOneAndUpdate( { _id: task._id, state: running, leaseUntil: task.leaseUntil }, { $set: { state: pending, workerId: null, leaseUntil: null, lastError: lease expired, retry after crash, } } ); }注意恢复的时候也要带条件更新防止万一是worker只是网络抖动、其实还在跑结果被重复调度。比较稳妥的做法是更新条件里带上leaseUntil: task.leaseUntil只有值没变过才允许重置。worker正常运行的时候对于耗时任务要主动续约。简单说就是每隔一两分钟把leaseUntil往后推一次证明“我还活着”。这个机制类比一下就是你租房子房东定期来查你有没有续租没续租就收房再租给别人。4. 高频翻车现场这些误区我基本都见过有的自己也踩过这一节是整篇的精华之一。下面这些坑我在大量代码评审和故障排查里反复见到每一个都值得单独拿出来说。4.1 误区一“查出来再更新”并发一上来就重复消费前面其实已经提到过这里再说透一点。用find查出一批pending任务再逐条update改成running两个操作之间存在时间窗口并发越高窗口越容易被撞上。我见过一个项目在压测到200并发的时候重复领取率直接飙升到30%以上。更麻烦的是这个问题不是必现的时有时无排查看不到明显报错最后是通过对比“领取记录”和“实际执行记录”才发现的。根治办法只有一条把“查询更新”合并为findOneAndUpdate这一个原子操作。用MongoDB的术语说就是对单文档的操作是原子的findAndModify系列命令把这个原子性用在了“先查再改”这个组合上。4.2 误区二把TTL索引当延迟队列用我在技术社区见过不止一个人问“能不能用TTL索引实现延迟任务”。先明确TTL索引是干什么的它是MongoDB用来自动清理过期数据的后台机制比如日志表保留30天到期自动删。它的工作方式是后台每60秒启动一次清理任务扫描所有TTL索引删除符合条件的文档。问题就在这里清理不是实时的默认60秒扫一次文档过期后最长可能还要等59秒才被删对延迟任务来说这个时间误差不可接受。它只是删除数据不会触发任何业务逻辑。你不能在文档删除的瞬间收到一个“任务到期了”的信号。如果用Change Streams监听delete事件技术上是能搭出一条链路但为了一个延迟任务引入这么一套不精确的机制完全是给自己挖坑。延迟任务的正确做法就是第2节说的scheduledAt字段精确到毫秒靠索引查询过滤完全可控。4.3 误区三对$natural排序的想当然有初学者写任务扫描代码时不指定排序指望MongoDB按插入顺序返回文档。MongoDB确实提供了一个$natural排序表示“按磁盘上的物理顺序返回”但文档在磁盘上的位置是会变的更新文档导致它变大超过原始分配的空间MongoDB会把文档移动到新的位置。compact命令或者副本集的同步逻辑也可能改变物理分布。所以$natural不等于插入顺序尤其在更新频繁的任务表里物理位置早就乱成一锅粥了。需要按顺序取任务就老老实实加一个createdAt或者用_id排序并建立对应的索引。4.4 误区四索引策略走两个极端一类人完全不建索引任务表数据到百万级后每次扫描都是全表扫CPU和IO飙升另一类人走另一极端给所有字段都建索引结果写入性能严重下滑。先说完全不建索引的后果。任务表每秒钟插入几千条扫描pending任务的时候如果没有合适的索引数据库要把整个集合的文档都读一遍才能过滤出符合条件的那几条。我帮人排查过一个案例单次任务领取耗时从几毫秒恶化到几百毫秒explain结果显示totalDocsExamined高达几十万选的索引是none。再看乱建索引的问题。每建一个索引写操作就要多维护一棵B树。任务表是写多读少的典型索引数量控制不住写入吞吐直接腰斩。正确的姿态是索引要为真实查询服务不为可能的查询服务。第2节给的那两个组合索引基本够覆盖90%以上的任务调度查询了。这是一张我常用的索引数量对照表数据量级索引数量建议说明百万以内2-3个主键之外建1-2个关键组合索引百万到千万3-5个根据实际慢查询日志逐步增加千万以上5-6个以内超过就要考虑分片或者调整查询逻辑4.5 误区五连接池和写入关注的隐形瓶颈任务调度的吞吐量上不去很多时候不是MongoDB本身慢而是客户端连接池配置不合理。MongoDB各语言驱动的默认连接池大小不一样Node.js驱动默认是100Java驱动默认也是100。如果你的worker是多线程密集领取任务每个线程持有一个连接超过池子大小之后新请求就要排队等空闲连接。我在压测时见过一个服务连接池默认设置并发从50提到200吞吐量不但没涨反而因为排队时间增加而降了30%。解决思路是分场景调整。领取任务用同步连接池批量写入用单独的异步批量通道两者互不干扰。maxPoolSize也不是越大越好池子太大数据库端连接数爆炸反而把MongoDB压垮。建议从100起步压测摸底找到自己的拐点。另一个容易忽视的是writeConcern。默认的写入关注是{ w: 1 }表示写入主节点即返回。如果业务要求更高持久性设置{ w: majority }表示大多数副本节点都写入成功才返回。代价是写入延迟升高因为要等待副本同步。任务调度场景我建议任务插入和状态更新这种核心写操作用w: majority结果回写这种允许丢失重建的数据用w: 1吞吐量能提升一大截。4.6 误区六忽视重复消费的幂等保护最后一个误区也最隐蔽很多人把“任务不重复”的希望寄托在存储层但分布式系统里没有任何存储方案能100%保证不重复。网络超时会让worker以为领取失败实际上数据库可能已经写入成功任务执行到一半网络分区恢复后系统重新调度业务逻辑可能已经执行了一部分。我的建议是存储层做到“尽量不重复”业务层必须做到“重复也无妨”。比如给任务加一个executionId每次领取任务时更新业务逻辑处理前检查是否已经处理过这个executionId或者保证业务逻辑的幂等性——重复跑结果一致。这是调度的最后一道防线。5. 吞吐量优化与容量规划从“能跑”到“扛得住压”任务调度的吞吐量优化是一个全局工程问题。索引和连接池调好之后还有几个点值得投入精力。5.1 批量写入和批量更新别再一条一条insert了任务量大的时候一条条插入简直是自残。比如仿真平台每天凌晨批量生成几百万条任务如果用循环单条insert光是网络往返就够让CPU原地起飞。正确姿势是insertManydb.tasks.insertMany( tasksArray, { ordered: false } );ordered: false很关键。默认ordered: true是一条失败就整体终止而false可以让MongoDB跳过失败的文档继续插入后面的任务。任务创建场景里某条数据格式不对根本不影响其他任务用false能最大化吞吐。批量更新结果也一样尽量把同类型的更新合并成一条bulkWrite而不是一个循环一个更新db.tasks.bulkWrite([ { updateOne: { filter: { _id: taskId1 }, update: { $set: { state: success, finishedAt: new Date(), result: result1 } } } }, { updateOne: { filter: { _id: taskId2 }, update: { $set: { state: failed, finishedAt: new Date(), lastError: timeout } } } } ]);注意bulkWrite要求所有操作在内存中构建批量太大的话要分片处理比如每批5000条。我实测过同样写10万条任务批量insertMany对比单条insert耗时能缩短到原来的十分之一左右。5.2 用explain和聚合管道定位慢操作别凭感觉优化吞吐量上不去先别急着加机器先用explain看查询计划。我排查慢查询时最常看几个指标db.tasks.find({ state: pending, scheduledAt: { $lte: new Date() } }).sort({ priority: -1 }).explain(executionStats);重点是这几个字段totalDocsExamined实际扫描的文档数。如果这个数字远大于返回条数说明索引没生效或者过滤条件太宽。totalKeysExamined扫描的索引条目数对应索引使用情况。executionTimeMillis总耗时。stage是否出现了COLLSCAN全集合扫描一旦出现且集合很大基本可以断定是慢查询。聚合管道在任务调度里也很有用。比如统计任务成功率、平均执行时长、失败原因分布db.tasks.aggregate([ { $match: { taskType: simulation, createdAt: { $gte: new Date(2024-05-01) } } }, { $group: { _id: $state, count: { $sum: 1 }, avgDurationMs: { $avg: { $subtract: [$finishedAt, $startedAt] } } } } ]);统计结果可以直接落到一个独立的task_metrics集合里前端展示实时报表时就不用反复扫描全表了。5.3 什么时候分片什么时候归档数据量上来后单集合的容量会成为新的天花板。我建议在数据量预估超过千万级且持续增长时认真考虑分片或者归档策略。MongoDB分片有几个坑要先避开分片键不能选自增字段比如ObjectId虽然单调增长但会导致写入总是落在同一个分片上完全失去分散效果。分片键必须是查询条件里高频出现的字段否则查询要广播到所有分片性能反而更差。任务表里_id哈希分片是一种选择适合按主键查询较多的场景如果更常按status createdAt查询可以考虑组合分片键但要先想清楚会不会造成写热点。我个人的习惯是给任务设置一个明确的数据生命周期比如“执行完成后保留最近7天数据7天以上的结果归档到另一个集合或者数据仓库”。这样主任务表的数据量始终控制在合理水平索引效率不会持续恶化查询和写入吞吐都更稳定。归档可以简单地用定时任务同步到task_archive集合然后从tasks集合删除原始文档。注意删除大批量数据时要用小批量多次删除或者按_id范围分批删避免一次性删除造成资源争用。6. 最后分享一点我在生产环境里沉淀下来的判断标准这篇文章写到这里核心内容讲得差不多了。最后不做什么大总结就分享几个我在多个调度系统项目里沉淀出来的判断标准和个人习惯仅供参考。第一不要一上来就上Kafka。如果你的调度任务量级在日千万以下MongoDB这套方案足够稳定还能省下一套中间件的运维成本。等数据量和复杂度真到了天花板再考虑引入消息队列也不迟而且到时候你的数据模型和状态机设计基本可以直接平移过去。第二状态机一定要简单。只有pending/running/success/failed/cancelled这几种状态宁可多搞几个字段记录细节也不要把状态拆出花来。状态越多并发边界越复杂排查问题越困难。第三每个任务都要有明确的归属者。workerId和leaseUntil这对组合是必须的否则出了问题你连“这条任务到底是谁在执行”都查不到。第四压测要早做且要模拟真实的并发场景。不要只测单条接口的延迟要测“多个worker同时抢任务”的竞争场景。这个场景才会暴露原子领取、连接池、索引设计这些真实问题。第五养成写幂等保护的习惯。这条放在最后因为它是整个系统的安全网。存储层能做到99.9%不重复但剩下的0.1%就靠业务层的幂等逻辑兜底。我见过太多调度系统在存储层做得很好却栽在一次网络抖动的重复执行上。MongoDB在任务调度场景里更像是一个“瑞士军刀”——队列、状态存储、结果存储、分布式锁、事件源能在一个系统里全部搞定。但它的边界也要清楚它擅长的是任务与状态的管理而不是极致的消息投递。用对了场景它能让你的系统简洁、可靠、吞吐可观用错了场景它也会让你在并发和一致性问题上焦头烂额。希望这篇文章能把你在调度系统存储选型上的思路理顺并帮你少踩几个我当年踩过的坑。
返回列表