如何避免两个EC2实例重复执行同一APScheduler定时任务
解决多EC2实例下APScheduler任务重复执行的问题
核心思路
多实例独立调度时,必须通过分布式锁来协调,确保同一时间只有一个实例能执行任务。任务执行前先尝试获取锁,成功则执行,失败直接跳过。
Redis分布式锁实现示例
Redis的SETNX(SET if Not eXists)命令支持原子性操作,适合做分布式锁。以下是修改后的代码:
首先安装依赖:
pip install apscheduler redis
修改后的任务代码:
from datetime import datetime from apscheduler.schedulers.blocking import BlockingScheduler import redis import time # 初始化Redis连接(根据你的实际配置修改) redis_client = redis.Redis(host='your-redis-host', port=6379, db=0, password='your-redis-password') sched = BlockingScheduler() def acquire_lock(lock_name, acquire_timeout=10, lock_timeout=7200): """获取分布式锁""" identifier = str(time.time()) end = time.time() + acquire_timeout while time.time() < end: # SETNX原子操作:不存在则设置,返回1表示成功获取锁 if redis_client.set(lock_name, identifier, nx=True, ex=lock_timeout): return identifier time.sleep(0.1) return None def release_lock(lock_name, identifier): """释放分布式锁""" pipe = redis_client.pipeline(True) while True: try: pipe.watch(lock_name) if pipe.get(lock_name).decode('utf-8') == identifier: pipe.multi() pipe.delete(lock_name) pipe.execute() return True pipe.unwatch() break except redis.exceptions.WatchError: pass return False @sched.scheduled_job('interval', id='my_job_id', hours=2) def job_function(): lock_name = "my_job_lock" # 获取锁:超时等待10秒,锁有效期2小时(需大于任务最长执行时间) lock_id = acquire_lock(lock_name, acquire_timeout=10, lock_timeout=7200) if not lock_id: print(f"{datetime.now()}: 未获取到锁,跳过任务执行") return try: print(f"{datetime.now()}: Hello World") # 这里编写你的实际任务逻辑 finally: # 确保锁被释放,避免死锁 release_lock(lock_name, lock_id) sched.start()
关键注意事项
- 锁有效期设置:必须大于任务的最长预估执行时间,防止任务执行过程中锁过期,导致其他实例重复执行。
- 异常处理:任务执行过程中如果抛出异常,通过
finally块确保锁被释放,避免死锁。 - Redis高可用:生产环境建议使用Redis集群或哨兵模式,避免单点故障导致锁机制失效。
其他可选方案
如果不想依赖Redis,也可以用AWS原生服务实现:
- DynamoDB条件写入:创建一张锁表,任务执行前尝试插入一条以任务ID为主键的唯一记录,插入成功则执行任务,完成后删除记录。
- S3对象锁:尝试创建一个唯一命名的S3对象,创建成功(不存在则创建)则执行任务,完成后删除对象。
内容的提问来源于stack exchange,提问作者Raghu
相关产品推荐
相关产品推荐

