""" 跨进程的定时任务互斥锁。 为什么需要这个 -------------- `gunicorn.conf.py` 配了 8 个 worker(`workers = min(cpu*2+1, 8)`),而 `run.py` 是在**模块级**启动 APScheduler 的 —— gunicorn 未开 preload_app,每个 worker 都会独立 import 一次 `run.py`,于是**每个 worker 各起一份调度器**,同一个 cron 任务在相同时刻被并发执行 8 次。 实测:容器里跑着 8 个 worker 进程,`create_app()` 与"调度器已启动"两条日志 的出现次数完全相等 —— 每个 worker 都配了一份。 后果:库存预警邮件每天实际会重复发 8 封。 APScheduler 自带的 `max_instances` 只在**单个调度器实例内**生效,管不了 跨进程;`replace_existing` 同理。这里用 PostgreSQL 的**会话级咨询锁** (`pg_try_advisory_lock`)做真正的跨进程互斥:抢到锁的那个 worker 才执行, 其余直接跳过。好处是不用引入 Redis 之类的新组件(compose 里也没有 redis, `redis_client` 恒为 None,相关装饰器全程 fail-open,指望不上)。 ★ 会话级锁绑定在**连接**上,必须在同一条连接上解锁 —— 所以这里显式取一条 专用连接,用完归还,不复用请求上下文里的 session。 """ import logging from contextlib import contextmanager from sqlalchemy import text from app.extensions import db logger = logging.getLogger(__name__) # 锁标识:PostgreSQL 咨询锁是 bigint,取固定值便于排查(不要用随机数)。 LOCK_INVENTORY_WARNING = 891001001 # 库存预警每日邮件 LOCK_DAILY_REPORT = 891001002 # MOM 系统日报 @contextmanager def advisory_lock(key): """ 尝试获取会话级咨询锁,产出「是否抢到」。 用法:: with advisory_lock(LOCK_DAILY_REPORT) as acquired: if not acquired: return # 别的 worker 正在跑,本轮跳过 ...实际干活... ★ 抢不到锁是**正常路径**,不是错误 —— 说明另一个 worker 正在执行同一任务。 调用方应当静默跳过,而不是报警。 ★ 拿锁失败(数据库不可用等)会让异常向上抛:宁可让调度器记一次失败, 也不要"没抢到锁"和"抢锁时数据库挂了"两种截然不同的情况被混为一谈。 """ conn = db.engine.connect() acquired = False try: acquired = bool( conn.execute(text("SELECT pg_try_advisory_lock(:k)"), {"k": key}).scalar() ) yield acquired finally: if acquired: try: conn.execute(text("SELECT pg_advisory_unlock(:k)"), {"k": key}) except Exception as e: # noqa: BLE001 # 解锁失败不应盖过业务异常;连接关闭时锁也会随之释放 logger.warning(f"[JobLock] 释放咨询锁 {key} 失败: {e}") conn.close()