
1. 问题背景多进程部署下的定时任务重复执行1.1 为什么 FastAPI 项目要采用多进程部署FastAPI 本质上是一个 ASGI 框架核心运行机制是 event loop 驱动的异步并发。单个 uvicorn worker 跑一个 event loop对 I/O 密集型请求已经足够快但受 GIL全局解释器锁限制CPU 密集型操作没法利用多核。我在上一家公司维护一个内部数据服务接口主要负责查数据和触发数据同步单进程压测能到 2000 QPS但一遇到 CPU 计算密集的聚合任务就开始卡顿最终决定用uvicorn main:app --workers 4拉四个 worker 分摊流量。这个部署方式本身没有任何问题。它是生产环境的标准操作问题是很多人包括我在内都会忽略一个事实每个 worker 都是完全独立的操作系统进程各自拥有独立的 event loop 和独立的内存空间。这个看似理所当然的事实在定时任务场景下会演变成一个隐蔽的坑。如果你问我什么时候应该考虑多进程部署我的经验是看两个指标一是 CPU 密集型逻辑占比是不是高到单核跑不满流量二是单进程 event loop 的阻塞事件是不是已经开始影响尾延迟。只有这两条里至少占一条才值得上多 worker。否则单纯为并发加进程反而会增加内存占用和排查问题的难度。1.2 定时任务重复执行的典型现象我们用 APScheduler 挂了一个每 30 秒执行一次的数据同步任务部署到测试环境后立刻发现现象很诡异任务日志里同一时刻出现 4 条一模一样的执行记录间隔完全一致数据库里同一批次的数据被处理了 4 次下游接口收到了 4 份重复的推送通知。最迷惑人的一点在于这个现象在本地uvicorn main:app单进程运行的时候完全不会出现。只有加上--workers 4才会出现而且不是偶发是 100 % 复现。如果你从来没处理过多进程调度问题第一反应肯定是找代码的 bug是不是任务注册了多次是不是replace_existing参数没写对我当初就把这些全都检查了一遍确认没有问题最后才把怀疑点放到 worker 进程数量上。1.3 复现场景一个最简单的错误示范下面这段代码几乎可以完整复现问题。它表面上看起来无懈可击该加的配置都加了但放到多进程环境里就出问题from contextlib import asynccontextmanager from fastapi import FastAPI from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.interval import IntervalTrigger scheduler AsyncIOScheduler() async def sync_data(): print(task running...) asynccontextmanager async def lifespan(app: FastAPI): scheduler.add_job( sync_data, IntervalTrigger(seconds30), idsync_data, replace_existingTrue, max_instances1, coalesceTrue, ) scheduler.start() yield scheduler.shutdown() app FastAPI(lifespanlifespan)max_instances1、replace_existingTrue、coalesceTrue全写了但它们在多进程环境下一点用都没有。原因后面细讲这里先记住一个结论APScheduler 的所有实例限制参数都只对当前进程内的 scheduler 实例生效管不到其他进程里的 scheduler。2. 为什么每个 worker 都会执行任务根因分析2.1 FastAPI 生命周期钩子与进程模型要理解这个问题得先搞清楚 uvicorn 多 worker 的启动流程。使用--workers 4时uvicorn 主进程会 fork 出 4 个子进程每个子进程做的事情完全一样导入app对象执行 lifespan 启动逻辑监听同一个端口各自处理进来的请求。问题就出在每个子进程都会执行 lifespan 启动逻辑这句话上。我们项目里scheduler AsyncIOScheduler()是模块级的全局变量lifespan里调用scheduler.start()。于是每个子进程都创建了一个属于自己的AsyncIOScheduler都往自己的调度循环里注册了sync_data这个 job。4 个进程4 个调度器4 份 job每个周期当然会执行 4 次。注意这里不是 多个 worker 争抢同一个 job而是 每个 worker 各有一个独立的 job 实例。这两者的区别很关键前者可以通过协调机制解决后者必须先明白所有 worker 都在独立运行才会想到用跨进程的互斥手段。验证方法很简单在启动代码里加一行print(fprocess {os.getpid()} scheduler started)多进程启动后日志里会出现 4 个不同的 PID每个 PID 都打印了一次。2.2 进程、线程、协程的锁为什么都解决不了踩坑之后团队里有同事提出一个想法用asyncio.Lock包住任务执行函数让同时只有一个协程能进去跑。这个方案当然不行。asyncio.Lock是协程级别的锁作用范围仅限于你当前进程的 event loop。而多进程部署下每个 worker 是独立的系统进程内存地址空间完全隔离进程 A 里的 Lock 对象在进程 B 里根本不存在也就不可能产生任何互斥效果。threading.Lock是线程级别的只能在当前进程内对线程互斥跨进程同样无效。Python 的 multiprocessing 模块提供跨进程锁基于信号量或共享内存但它只适用于同一台机器上的进程。如果你的服务做了多机水平扩展4 台机器各跑两个 worker这种锁也管不住跨机器的重复执行。我后来在给团队分享时用了这个类比进程就像一户户独立的人家灯泡、水龙头各用各的。线程是同一户人家里的不同房间协程是同一个房间里轮流干活的人。你不可能用 A 家的门锁去锁 B 家的大门。分布式锁的逻辑就是在小区门口装一个公共闸机谁要出小区都得先过闸机闸机会记录当前是谁在通行。2.3 两个根治方向拆分调度进程 vs 分布式互斥理解了根因解决方案其实就两条路。第一条路是让调度器只存在于一个进程里。具体做法是把 scheduler 从 FastAPI 应用的 lifespan 里摘出来单独写一个脚本scheduler_main.py用python scheduler_main.py单独起一个进程跑调度逻辑业务服务只负责处理请求。这种方式架构上最干净但代价是你需要自己保证调度进程的高可用。如果它挂了整个定时任务系统就停了得有额外的守护机制比如 systemd 拉起、supervisor 监控或者干脆上 K8s 的 Deployment 单副本。第二条路是保留每个 worker 里的 scheduler但在任务真正执行前加一道分布式锁。所有 worker 执行任务前先去抢锁只有抢到锁的那个 worker 才继续跑其他的直接跳到下一个周期。这条路对现有代码侵入小而且天然支持多台机器横向扩展不需要额外维护一个调度进程的高可用逻辑。我最后选了第二条。原因很简单我们的定时任务本来就有比较完整的失败告警和重跑机制偶尔因为 Redis 抖动漏跑一期可以接受但整个调度进程挂掉需要额外运维投入团队资源不够。如果你对任务执行率要求极高且运维能力到位第一条路其实更可靠。3. Redis 分布式锁让任务只跑一次3.1 核心原理SET NX EX 原子操作分布式锁要解决的核心问题只有一个多个进程如何对一个共享资源进行互斥访问。Redis 分布式锁就是利用 Redis 单线程命令处理的特性把判断 设置两个动作合并成一条原子命令SET scheduler_lock:sync_data owner_id NX EX 60这条命令拆开看NX表示 key 不存在时才设置成功存在则返回空EX 60表示设置 60 秒过期时间。多个进程同时执行这条命令Redis 单线程模型保证它们被逐个处理最终只有一个进程能拿到成功返回这个进程就相当于拿到了锁。执行完任务之后调用DEL把 key 删掉锁就释放了。为什么要加过期时间因为如果持有锁的进程突然崩溃OOM、kill -9、宿主机宕机它永远来不及执行DEL如果锁没有过期时间其他进程就永远无法接手这个任务。EX是分布式锁最基本的兜底设计目的是防止锁永久不释放这种比重复执行更糟糕的情况。3.2 用 redis-py 实现一个最小可用的分布式锁项目里我们用redis 4.6的redis.asyncio客户端。这里强调一下为什么必须用异步客户端FastAPI 是异步应用所有请求都跑在 event loop 上如果任务里用同步 Redis 客户端去获取锁一个SET命令阻塞几十毫秒整个服务的并发能力就会被拖垮。下面是当时实现的锁类去掉了一些日志和埋点核心逻辑保持原样import asyncio import os import time import uuid from redis import asyncio as aioredis class DistributedLock: 基于 Redis SET NX EX 的分布式锁用于跨进程互斥 def __init__(self, redis: aioredis.Redis, lock_name: str, ttl: int 60): self.redis redis self.lock_key fscheduler_lock:{lock_name} self.owner_id f{os.getpid()}-{uuid.uuid4().hex} self.ttl ttl async def acquire(self, timeout: float 0.0) - bool: deadline time.time() timeout while True: ok await self.redis.set( self.lock_key, self.owner_id, nxTrue, exself.ttl, ) if ok: return True if time.time() deadline: return False await asyncio.sleep(0.2) async def release(self) - None: # 用 Lua 脚本保证只删除自己持有的锁避免误删其他 worker 的锁 script if redis.call(get, KEYS[1]) ARGV[1] then return redis.call(del, KEYS[1]) else return 0 end await self.redis.eval(script, 1, self.lock_key, self.owner_id)几个细节值得单独说明。第一锁的 value 用进程 PID uuid作为持有者标识。这个标识的意义在于释放锁的时候可以确认这把锁确实是我持有的避免误删。如果只用固定字符串比如1就无法做这个判断。第二release用了 Lua 脚本。这是安全释放锁的经典做法不是直接DEL。为什么要这么麻烦设想这个场景任务 A 持有锁运行运行时间太长锁在 TTL 到期后自动过期了任务 B 此时抢到锁开始运行任务 A 终于结束执行DEL把锁删掉。结果就是任务 B 的锁被 A 误删任务 B 还没跑完任务 C 就能再次拿到锁三个任务同时执行。用 Lua 脚本先GET判断 value 是否等于自己的owner_id等于才删可以完全杜绝这种连锁误删问题。3.3 把锁嵌入任务执行流程有了锁类之后任务函数本身不需要改业务逻辑。我写了一个通用包装函数把所有定时任务统一用这个包装器跑async def run_with_lock(func, lock_name: str, ttl: int 60): lock DistributedLock(redis_client, lock_name, ttl) if not await lock.acquire(timeout0): logger.info(f[{lock_name}] 未获取到锁本轮跳过) return try: await func() finally: await lock.release()这里timeout0非常关键拿不到锁直接跳过而不是傻等。原因在于调度器每个周期都会触发一次任务等待没有意义拿不到锁就说明这个周期已经有别的 worker 在执行了直接放弃即可。如果把 timeout 设为很大比如 30 秒多个 worker 会同时挂起等待浪费资源不说锁释放的一瞬间还可能发生惊群效应多个 worker 同时抢锁好在 SET NX 是原子命令最终只有一个能抢到但白白消耗了一堆连接和内存。3.4 锁的 TTL 应该怎么设置TTL 是分布式锁最容易翻车的参数没有之一。设太短任务没执行完锁就过期其他 worker 趁机抢锁重复执行的问题又回来了设太长如果拿到锁的 worker 真的进程崩溃Redis 里这个 key 要等很久才被自然清理这段时间任务会一直处于无人执行状态等于定时任务停了。我的经验是先统计任务的历史执行耗时取 p95 或 p99 的耗时再乘以 2 到 3 的系数。拿我们的sync_data任务举例正常耗时 5 秒最慢一次跑到过 20 秒下游接口慢于是 TTL 设置为 60 秒。这样既不会频繁过期也不会在异常退出后阻塞太久。如果你拿不准某个任务的耗时分布可以先不接锁纯日志跑两三天看历史数据再定 TTL。3.5 还有哪些更复杂的方案Redlock 与开源锁库严格来说我上面讲的属于单节点 Redis 分布式锁在单节点 Redis 场景下够用。如果你追求更高可靠性也就是服务 Redis 主从切换时锁不失效业界有 Redlock 算法多节点仲裁和现成开源库比如redis-lock、aioredlock。但 Redlock 在分布式系统界一直有争议很多架构师认为它在极端网络分区场景下也存在安全窗口。我的建议很务实用于定时任务去重这种场景单节点 Redis 锁的可靠性已经足够。如果你真的需要一个极端可靠的调度系统正确的方向不是升级锁而是改用任务队列加幂等消费设计把任务不丢失、不重复执行的语义交给消息中间件来保证比在锁层面较劲更实际。顺带提一句后续如果项目复杂度上来了建议调研一下 Celery Beat 或者 APScheduler 结合scheduler-lock这类成熟方案。生产环境的稳定性往往不取决于某一层技术的先进程度而是取决于每一层是否都保持简单可靠。4. 实操过程与完整代码4.1 项目环境准备演示项目的依赖和版本如下都是当时实际使用的版本依赖版本说明Python3.11原生支持 asyncio类型提示完善fastapi0.104使用 lifespan 上下文管理uvicorn[standard]0.24支持--workers多进程启动apscheduler3.10异步调度器 AsyncIOSchedulerredis4.6提供了 redis.asyncio 异步客户端本地需要有一个 Redis 服务。没有现成 Redis 的可以直接用 docker 起一个最小实例docker run -d -p 6379:6379 redis:7-alpine4.2 完整可运行示例这是一个可以直接跑起来的完整项目我把所有代码都放在main.py里方便演示。实际项目建议拆模块锁相关代码放core/lock.py任务注册放jobs.py主应用只保留生命周期装配。import asyncio import logging import os import time import uuid from contextlib import asynccontextmanager from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.interval import IntervalTrigger from fastapi import FastAPI from redis import asyncio as aioredis logging.basicConfig( levellogging.INFO, format%(asctime)s [%(levelname)s] [%(processName)s] %(message)s, ) logger logging.getLogger(scheduler-demo) # ---------- 分布式锁 ---------- class DistributedLock: 基于 Redis SET NX EX 的分布式锁 def __init__(self, redis: aioredis.Redis, lock_name: str, ttl: int 60): self.redis redis self.lock_key fscheduler_lock:{lock_name} self.owner_id f{os.getpid()}-{uuid.uuid4().hex} self.ttl ttl async def acquire(self, timeout: float 0.0) - bool: deadline time.time() timeout while True: ok await self.redis.set( self.lock_key, self.owner_id, nxTrue, exself.ttl ) if ok: return True if time.time() deadline: return False await asyncio.sleep(0.2) async def release(self) - None: # 只删除自己持有的锁避免误删 script if redis.call(get, KEYS[1]) ARGV[1] then return redis.call(del, KEYS[1]) else return 0 end await self.redis.eval(script, 1, self.lock_key, self.owner_id) # ---------- 全局对象 ---------- redis_client aioredis.from_url(redis://127.0.0.1:6379/0, decode_responsesTrue) scheduler AsyncIOScheduler(timezoneAsia/Shanghai) # ---------- 定时任务 ---------- async def sync_data(): lock DistributedLock(redis_client, sync_data, ttl60) if not await lock.acquire(timeout0): logger.info([sync_data] 未获取到锁本轮跳过) return try: logger.info([sync_data] 任务开始执行) # 模拟业务逻辑耗时 await asyncio.sleep(5) logger.info([sync_data] 任务执行完成) finally: await lock.release() # ---------- FastAPI 生命周期 ---------- asynccontextmanager async def lifespan(app: FastAPI): scheduler.add_job( sync_data, IntervalTrigger(seconds30), idsync_data_job, replace_existingTrue, max_instances1, coalesceTrue, misfire_grace_time30, ) scheduler.start() logger.info(fprocess {os.getpid()} scheduler started) yield scheduler.shutdown() await redis_client.aclose() app FastAPI(title多进程定时任务去重演示, lifespanlifespan) app.get(/healthz) async def healthz(): return {status: ok}代码里的每个对象都值得解释清楚。AsyncIOScheduler的运行基于当前进程的 event loop所以多 worker 下每个 worker 都有独立调度循环。IntervalTrigger(seconds30)表示每 30 秒触发一次。misfire_grace_time30表示任务错过了预定执行时间但在 30 秒内补跑都是合法的这个参数在任务队列积压时避免大量补跑我建议保留。max_instances1和coalesceTrue要写上虽然跨进程管不住但能在单进程内防止任务重叠问题是有益无害的。4.3 多进程启动与验证命令本地验证用 4 worker 模式启动uvicorn main:app --host 0.0.0.0 --port 8000 --workers 4 --log-level info如果用 gunicorn注意指定 UvicornWorker 的 worker classgunicorn main:app -w 4 -k uvicorn.workers.UvicornWorker -b 0.0.0.0:8000启动后观察日志理想情况下每 30 秒只打印一次任务开始执行另外 3 个 worker 打印的是未获取到锁本轮跳过。我当时的实际运行日志长这样简化版[process worker1] scheduler started [process worker2] scheduler started [process worker3] scheduler started [process worker4] scheduler started [process worker3] [sync_data] 任务开始执行 [process worker1] [sync_data] 未获取到锁本轮跳过 [process worker2] [sync_data] 未获取到锁本轮跳过 [process worker4] [sync_data] 未获取到锁本轮跳过 [process worker3] [sync_data] 任务执行完成这组日志就是项目验收的最直接证据。如果任务开始执行只出现一次说明分布式锁生效了。4.4 锁的流转过程验证除了看日志还可以直接观察 Redis 里的锁 key 状态。任务运行期间执行redis-cli keys scheduler_lock:* redis-cli get scheduler_lock:sync_data任务运行时能拿到 keyvalue 是某个 worker 的 PID 加 uuid任务结束后 key 被删除下一轮周期立刻又有新的 value。如果任务运行中 key 突然消失说明 TTL 设置得太短锁提前过期了这是最需要警惕的信号。我为了给团队演示锁的真实争抢过程还做过一个实验把IntervalTrigger改成每 2 秒触发一次任务里sleep改成 10 秒。这样一来2 秒下锁根本来不及释放下一轮就该触发了日志里会持续出现其他 worker 的跳过记录。这个实验直观展示了锁被持有期间其他 worker 只能干等或者跳过的状态建议你在自己的环境里也跑一下对理解分布式锁的互斥语义非常有帮助。5. 常见问题与避坑经验5.1 锁的 TTL 设置不当导致任务重叠这是我踩过的最大一个坑。最初我把 TTL 设成 30 秒结果某个任务上传文件到对象存储某次网络抖动跑了 45 秒锁在第 30 秒过期另一个 worker 抢到锁开始执行两个 worker 同时处理同一个文件最终数据异常。排查了很久才意识到是 TTL 不够。解决思路有两个一是把 TTL 调大比如任务最长耗时的 3 倍二是给锁加续期机制也就是后台起一个协程定期把 TTL 重置直到任务结束。续期机制的实现代码会多一些我当时为了快速解决直接采用了统计历史耗时 3 倍系数的方案。如果你的任务执行时间波动极大建议早日上续期别在 TTL 上赌运气。5.2 Redis 不可用时任务执行策略如何选择分布式锁依赖 RedisRedis 挂掉之后 acquire 会直接失败任务会跳过不会重复执行但任务也彻底不跑了。这可能是你不想看到的。可以增加降级策略如果 Redis 连不上在本进程内用asyncio.Lock加threading.Lock双锁保证单机多个 worker 之间不重复但多机部署无法保证。我们内部小项目直接选择了Redis 挂掉任务暂停的策略因为核心数据同步有离线兜底跑批定时任务只是提高时效性暂停可接受。如果你对任务不可执行零容忍建议把任务消费模式改成消息队列驱动这是更彻底的解耦方案。5.3 锁误删问题比重复执行更隐蔽上面 Lua 脚本已经规避了锁误删但我还是要单独强调一次不要因为DEL简单就省掉先判断再删除这个步骤。锁误删的典型事故链是任务 A 持有锁锁过期任务 B 抢到锁任务 A 释放时把 B 的锁删了任务 C 趁机拿到锁三个任务同时执行数据全乱。这个场景里每个环节看起来都是正常操作但组合在一起就是灾难。用 Lua 脚本判断持有者是把整个链路锁死的唯一办法。注意release里的 Lua 脚本判断的是 value 等于自己的owner_id才删除但这只保证持有者身份正确。如果想彻底避免过期锁的问题还需要配合续期机制保证任务执行期间锁不提前释放。5.4 本地单进程开发时锁会影响调试吗本地开发一般uvicorn main:app单进程运行锁同样会走 Redis。如果没有本地 Redis开发环境就起不了服务。我建议加一个显式开关比如用环境变量ENABLE_DISTRIBUTED_LOCK控制开发环境关闭生产环境打开。这个开关有两个细节要注意第一默认值必须设为不跳过锁。一旦默认跳过很可能有人忘记打开把 bug 带到生产。第二跳过锁的日志要打印明显告警。我吃过这个亏有同事本地开着开发模式测试没意识到生产配置里这个开关根本没生效结果又重复执行了一次我排查了半小时才发现是开关配置问题。显式日志可以让这类问题在第一时间暴露。5.5 与任务队列方案Celery Beat、xxl-job 等对比如果你的定时任务本身比较复杂需要重试、失败补偿、批量执行、可视化监控建议直接评估 Celery Beat、xxl-job 这类独立的调度框架。Celery 的语义是任务分发worker 之间天然只有一个进程消费同一个 task配合 Redis broker 也天然解决了持久化和重试的问题但代价是引入更重的依赖和运维组件。我做个简单的选型对比方便你按场景对号入座方案优点缺点适用场景APScheduler 分布式锁轻量、侵入小、无额外组件需自己保证锁与幂等无可视化轻量周期任务、内部数据同步独立调度进程逻辑简单、无锁依赖需额外守护调度进程高可用团队运维能力较强Celery Beat worker成熟可靠、带重试与监控依赖组件多、配置复杂复杂任务流、需要精准控制频率xxl-job 等外部调度平台可视化、分布式管理需要部署额外服务、引入平台依赖中大型团队统一任务治理结论是对于单机或小规模集群里的轻量定时任务APScheduler 加分布式锁是性价比最高的方案。等任务量和复杂度上来之后再考虑迁移到独立调度平台那个时候你的心跳监控、告警、日志体系已经构建好迁移成本反而低。5.6 定时任务本身要做幂等设计最后这条最重要放在最后压轴。锁只是保证同一时刻只有一个进程在跑幂等是保证即使重复执行也不会产生脏数据两者是完全不同层面的设计千万不要混为一谈。我前段时间写了一个 job功能是同步用户余额到报表表。加锁之前重复执行会重复累加余额加锁之后表面上没有重复了但我心里清楚锁不是万能的——如果某天 Redis 集群故障、或者某次发布新代码改了 TTL 导致锁短暂失效任务重复执行金额就会翻倍。这种问题靠事后排查是发现不了的唯一靠谱的防线是任务自身的幂等数据库里加唯一约束业务上使用先删后插或幂等键设计保证重复触发时数据依然正确。如果你现在只学到了分布式锁的用法那这个问题只解决了一半。另一半永远在业务层不管外层机制多强核心数据的写入逻辑都必须能容忍重复调用。这不仅是技术习惯更是一种工程防御思维。最后再说两个实操心得踩完这个坑之后我最大的体会是多进程定时任务重复执行这个问题本质上不是 FastAPI 或 APScheduler 的 bug而是每个进程各自初始化调度器这种部署模型与全局唯一调度需求之间的天然矛盾。解决矛盾的方式可以有很多种但核心原则只有一条就是让调度和执行的互斥发生在进程之外、机器之外。最后分享两个小技巧。第一上线验证时不要只看业务日志直接对比数据库落库数据的条号或者唯一 ID日志可能被框架缓存误导数据是骗不了人的。第二如果你们用 git 做 CI/CD建议把单 worker 与多 worker 的对比测试写成一个小小的冒烟测试脚本每次发布自动跑一遍确认锁在目标 worker 数下依然生效。我在项目里就是这么做的虽然多花了一点时间但后面几次上线都因为这一步排查掉了不少潜在问题。