多 Worker 安全抢任务:租约、围栏令牌与迟到写入

多 Worker 安全抢任务:租约、围栏令牌与迟到写入
从 outbox.pending 到可竞争的任务上一篇durable_store.py产出了outbox表、稳定command_id以及pending、mark_done。单 worker 可以扫描待处理行多 worker 却会同时读到同一条记录并执行。远端幂等能降低重复副作用但不能代替本地协调昂贵调用仍会重复限流额度会浪费非幂等旧系统更危险。本篇复用 outbox 行作为任务增加租约所有者、租约期限和单调递增的围栏令牌。租约与永久锁不同。worker 崩溃后不会释放普通应用锁而租约到期后其他 worker 可以接管。代价是原持有者可能只是长时间暂停并未死亡它恢复时会与新持有者并行。因此只有“我曾拿到租约”不够每次完成都必须验证令牌仍是最新。令牌像号码牌后领取者数字更大存储拒绝旧号码的完成请求。原子领取与条件完成保存为lease_store.py。SQLite 支持UPDATE ... RETURNING这里先选候选 ID再用带条件的更新领取。BEGIN IMMEDIATE在写事务开始时取得保留锁使两个领取者不能同时通过临界区。时间由数据库的unixepoch()产生避免不同 worker 系统时钟漂移参与比较。importsqlite3fromdurable_storeimportDurableStore LEASE_COLUMNS[ALTER TABLE outbox ADD COLUMN lease_owner TEXT,ALTER TABLE outbox ADD COLUMN lease_until INTEGER,ALTER TABLE outbox ADD COLUMN fence INTEGER NOT NULL DEFAULT 0,]classLeaseStore(DurableStore):def__init__(self,path:str):super().__init__(path)withself.connect()asdb:forstatementinLEASE_COLUMNS:try:db.execute(statement)exceptsqlite3.OperationalErroraserror:ifduplicate columnnotinstr(error):raisedefclaim(self,owner:str,ttl_seconds:int30):withself.connect()asdb:db.execute(BEGIN IMMEDIATE)rowdb.execute(SELECT command_id FROM outbox WHERE processed_at IS NULL AND (lease_until IS NULL OR lease_until unixepoch()) ORDER BY event_seq LIMIT 1).fetchone()ifrowisNone:returnNoneclaimeddb.execute(UPDATE outbox SET lease_owner?, lease_untilunixepoch()?, fencefence1 WHERE command_id? RETURNING *,(owner,ttl_seconds,row[command_id]),).fetchone()returnclaimeddefcomplete(self,command_id:str,owner:str,fence:int)-bool:withself.connect()asdb:cursordb.execute(UPDATE outbox SET processed_atCURRENT_TIMESTAMP WHERE command_id? AND lease_owner? AND fence? AND processed_at IS NULL,(command_id,owner,fence),)returncursor.rowcount1运行输出模块定义成功无标准输出租约期限不参与complete条件是一个有意选择如果租约虽过期但尚无人接管旧持有者完成仍可接受一旦新人领取fence增加旧持有者必然失败。也可以要求lease_until now但会让刚好跨过期限的成功结果被丢弃造成更多重试。真正保护并发所有权的是围栏令牌而不是对毫秒边界的迷信。模拟暂停、接管和迟到完成保存为demo_105.py。第一个 worker 领取一个零秒租约第二个立即接管。旧 worker 随后尝试完成条件更新影响零行新 worker 使用更大令牌完成成功。示例不依赖睡眠因此输出稳定、测试快速。importtempfilefrompathlibimportPathfromlease_storeimportLeaseStorefromminiflowimportCommandwithtempfile.TemporaryDirectory()asdirectory:storeLeaseStore(str(Path(directory)/flow.db))store.append_with_commands(trip-001,0,trip_requested,{},[Command(lock_budget,{trip_id:trip-001})],)oldstore.claim(worker-old,ttl_seconds0)assertoldisnotNoneprint(old fence:,old[fence])newerstore.claim(worker-new,ttl_seconds30)assertnewerisnotNoneprint(new fence:,newer[fence])old_okstore.complete(old[command_id],worker-old,old[fence])new_okstore.complete(newer[command_id],worker-new,newer[fence])print(old completion:,old_ok)print(new completion:,new_ok)print(pending:,len(store.pending()))assert(old_ok,new_ok)(False,True)运行输出old fence: 1 new fence: 2 old completion: False new completion: True pending: 0围栏必须延伸到真正资源这是租约最容易被讲浅的地方本地数据库拒绝旧 worker 的“完成标记”不代表外部副作用没发生。若两个 worker 都向不支持幂等或围栏的设备发命令旧 worker 仍可能在外部覆盖新结果。严格方案是让资源端也保存最大 fence只接受更大的令牌若资源端只支持幂等键则同一command_id至少可以折叠重复两者都不支持时只能通过单写代理把危险资源纳入可控边界。非平凡踩坑是自动续租线程。Python 进程发生长时间 GC、宿主机挂起或网络分区时续租可能停止而业务线程仍在执行。恢复后它不能因为“续租线程又正常了”就继续提交必须重新读取租约并验证 fence。活动执行时间可能超过 TTL 时应在安全检查点续租对不可中断的长调用TTL 要覆盖合理上界同时接受故障恢复变慢的权衡。另一个坑是公平性。始终ORDER BY event_seq会让某个快速失败的老任务占据扫描前列。真实查询应加入available_at与尝试次数只选到期任务并对工作流做限流避免一个大客户耗尽全部 worker。领取事务要短只更新所有权不在事务里调用网络。SQLite 同时只有一个写者若持锁执行 HTTP整个引擎都会排队。本篇交接本篇产出的lease_store.py提供claim(owner, ttl_seconds)和带fence的complete上一篇的 outbox 现在能被多个 worker 安全领取。下一篇将复用lease_until所体现的“数据库时间”原则但不再用短期租约等待业务日期我们会建立持久化 timer 表用确定性的 timer ID 安排退避和数天后的唤醒重启期间也不会丢失。 觉得有用就点个赞 收藏方便回头查阅有疑问直接在评论区留言我看到都会回。 文章里的代码都能直接跑。想要可直接 clone 的完整工程 配套部署脚本 / 踩坑清单评论一声或发邮件到cj2664qq.com我免费发你。如果你正好在做类似系统、或有工程化难题想找人做也欢迎邮件聊一句——我按实际情况评估能落地的就接单或出方案。评论和邮件都能直接找到我不用跳别的平台。