PostgreSQL 原子领取任务:用行锁隔离并发发布器

PostgreSQL 原子领取任务:用行锁隔离并发发布器原创封面

先定义“领取”,不要从一条 SELECT 开始

任务领取的业务含义不是“某个进程看见了这行”,而是数据库已经把处理权在一个明确期限内交给某个工作者。记录至少需要待处理、执行中、已完成和需重试等状态,同时保存计划时间、领取者、租约截止时间、尝试次数与最后错误分类。只有这些字段共同变化,监控页面才能区分排队、真正运行和失去工作者的悬挂任务,也才能在进程重启后重建事实,而不是依靠内存队列猜测。

领取契约还要说明一次处理允许产生什么副作用。数据库内的纯计算可以放在同一事务完成;发送邮件、调用支付或生成大文件则不适合持锁等待。对于后一类任务,领取只授予执行资格,外部动作必须带稳定操作标识,完成后再用版本条件提交结果。这样即使工作者在外部系统成功后、本地确认前退出,接管者也会先查询或重放幂等结果,不会直接制造第二次副作用。

表结构和索引要服务于到期查询

常见候选条件是状态为待处理、run_at 不晚于数据库当前时间,再按 run_at 与稳定主键排序。对应索引应把低基数字段与时间、主键组成可用的访问路径,或使用只覆盖待处理状态的部分索引。是否选择部分索引要看状态分布和迁移习惯:它能显著缩小热队列索引,却要求查询谓词与索引条件保持一致。使用 EXPLAIN ANALYZE 时要装入接近生产的状态比例,空表上的顺序扫描没有判断价值。

任务正文和领取元数据最好分离考虑。巨大 JSON 载荷频繁更新会增加行宽、WAL 与 vacuum 压力,而领取只需要修改少量控制字段。可以让任务表保存载荷引用,或确保每次状态变更不会无谓重写大型字段。还要给 operation_key、业务聚合标识等真正要求唯一的字段建立约束,让数据库承担并发裁决;仅在应用中先查重再插入,会在两个连接同时通过检查时失效。

用短事务完成候选选择与条件更新

领取事务可以从按序候选中使用 FOR UPDATE SKIP LOCKED,让多个工作者跳过已经被其他事务锁住的行。锁定之后立即把状态改为执行中,写入 worker_id、lease_until 与递增版本,然后提交并返回本次领取的数据。选择和更新必须属于同一事务;若先 SELECT、提交,再逐行 UPDATE,其他连接会在空隙里得到相同候选,SKIP LOCKED 的保护也就完全消失。

批量大小应由单项耗时、租约长度和内存预算决定,而不是越大越好。一次锁住数千行后慢慢处理,会让后续工作者看不到这些任务,也会增加崩溃后的恢复延迟。更稳妥的做法是领取一个可在租约内完成的小批次,事务提交后才执行耗时工作。提交结果时再次匹配 id、version、worker_id 和执行中状态,影响行数为零就表示所有权已经变化,旧工作者不得覆盖接管者的新状态。

租约、心跳与接管需要同一个时钟事实源

租约到期判断应尽量使用 PostgreSQL 的当前时间,而不是混用多台应用服务器的本地时钟。租约长度至少覆盖正常处理的高分位耗时,并留出网络与调度抖动;过短会让仍在运行的任务被接管,过长则拖慢真正崩溃后的恢复。长任务可以心跳续租,但续租必须匹配原工作者和版本,并限制最大执行时长,避免故障代码无限占有任务。

接管流程不是把所有过期执行中任务直接改回待处理。它要先判断外部动作是否可能已经完成:若任务保存了幂等键或远端资源编号,应查询外部状态并据此提交完成、继续等待或安全重试。无法查询且副作用不可逆的任务需要进入人工对账队列。数据库能原子地转移本地所有权,却不能证明事务之外的系统没有执行,因此“租约过期等于未执行”是危险假设。

失败分类决定退避与终止

格式错误、缺少业务资源和违反不变量属于内容性失败,继续原样重试不会改善结果。发布器应在一次受控事务中把该任务放入暂停或死信状态,记录短小、脱敏且可行动的原因。连接超时、数据库重启、限流和暂时网络故障则属于系统性失败,此时保留可重试状态,按有上限的指数退避调整下次时间。把所有异常都归为同一 retry,会制造热点任务并掩盖真正需要修复的数据。

尝试次数不能单独决定最终失败,因为一次进程崩溃可能发生在远端已经成功之后。每次尝试都应携带同一业务操作键,并记录阶段:尚未调用、调用结果未知、已获得远端结果、正在提交本地确认。恢复逻辑依据阶段采取查询、重试或对账。错误日志不要保存完整载荷和凭据,可记录任务 id、版本、错误类别、依赖响应标识与下一次计划时间,使值班人员能复盘而不扩大敏感数据范围。

公平性和吞吐量要用队列年龄观察

SKIP LOCKED 提升并发吞吐,但它不天然保证公平。如果队首任务频繁失败或持续被锁,简单按时间取小批次可能让某些租户长期等待。可以在排序中加入优先级与稳定主键,并为租户设置并发上限;更复杂的加权策略应先用真实流量验证。关键指标不是每分钟扫描多少行,而是最老可执行任务的等待时长、领取到开始的延迟、完成率和每类失败的积压。

工作者空轮询也会给数据库制造稳定负载。没有候选时可以退避,或通过通知仅作为唤醒提示,但最终仍以表中状态为事实。指标应分开记录候选查询耗时、锁跳过数量、领取批次大小、租约续期次数与接管次数。接管率突然升高通常意味着租约预算、下游延迟或进程稳定性出现变化,而不是简单增加更多工作者就能解决。

并发测试必须真的使用多个连接

单元测试可以验证状态转换和退避计算,却无法证明行锁行为。集成测试需要使用真实 PostgreSQL 和多个独立连接,同时在屏障处释放十个领取事务,然后断言每个任务只被一个工作者获得、没有遗漏且顺序满足约定。还应在锁定后提交前、提交后外部调用前、外部调用后确认前分别终止进程,验证租约接管与幂等路径。用串行 mock 返回固定结果无法覆盖这些竞争窗口。

测试时要保存任务初始快照,并在所有工作者结束后核对状态、版本、修订或审计数量。把一个工作者故意保持长事务,确认其他连接会跳过而非阻塞;让租约在可控时钟下到期,确认旧工作者的迟到提交影响零行。性能测试还要观察数据库连接数、锁等待、WAL 和表膨胀,任何重复副作用都应立刻判定失败,不能被较高吞吐量抵消。

上线时保留回退和清理路径

首次启用原子领取前,可以让新工作者只读取影子候选并记录它与旧队列的差异,再以少量租户灰度。迁移新增字段时先允许为空,回填历史任务,再逐步启用约束,避免长锁影响在线业务。扩容工作者要同步检查连接池预算,不能让每个进程默认连接数相乘后超过数据库容量。下线旧实现前,确认它不再修改相同状态字段,否则两个协议会互相破坏所有权。

完成任务的保留期和清理也属于队列设计。删除历史行应按主键小批次执行,并确保审计、对账和灾难恢复需要的证据已经转存;不要在高峰执行一个覆盖数月数据的巨大 DELETE。上线后定期演练数据库重启、工作者强杀和下游长时间不可用,检查积压能否在预算内恢复。只有在故障演练中仍保持一次所有权和可追踪副作用,领取协议才算真正闭环。

隔离级别、死锁与手工操作必须服从领取协议

PostgreSQL 默认 Read Committed 下,领取语句要把到期、状态、排序和行锁放进同一事务,并把 RETURNING 的行作为唯一候选来源。提高隔离级别不会代替租约与版本;Serializable 还可能产生需要整体重试的序列化失败。执行器只对尚未产生外部副作用的短事务有限重试,不能在远端结果未知后重跑整段业务。

任务若还要锁业务资源,所有工作者采用一致锁顺序以减少死锁。数据库中止事务后必须从新连接和新快照重新领取,不能继续使用已经失败的事务对象。死锁、连接断开和内容错误分别统计;把它们都写成“稍后重试”会让一条永远无效的任务反复占据队首,也让真正的数据库容量问题失去信号。

运维人员重置任务也要匹配当前状态、版本和租约,并填写原因审计。直接把全部 RUNNING 改回 PENDING 会抹掉存活工作者的所有权,迟到提交随后与接管者竞争。后台应提供暂停、确认远端结果和安全重新排队等受控命令,展示尝试次数、worker、租约与外部操作标识后再允许操作。

队列表权限只授予发布账户所需的查询、条件更新、修订和审计能力,公开应用不拥有任意批量领取权限。部署阶段在回滚事务中执行一次真实权限检查,证明服务账户能完成协议且不会提交测试行。若直到首个计划时间才发现权限不足,准点发布已经失去保障,故权限验证属于发布前门禁。

实现片段

SELECT id FROM jobs WHERE due_at <= now() FOR UPDATE SKIP LOCKED

一手参考资料

评论 · 0

还没有评论留下第一句经过思考的话。

游客评论需审核。注册后可直接公开,无需审核。

SHARE / 分享

分享这篇文章

WECHAT / 微信

用微信扫一扫

在手机微信中打开文章后,再从微信右上角分享给朋友或朋友圈。