本文目录导读:

- 使用单线程串行执行(最基础)
- ">使用APScheduler的"间隔"模式(推荐)
- 使用锁控制并发(APScheduler)
- 使用Celery的自动重试和限流
- 使用Redis分布式锁(多进程场景)
- 使用信号量控制并发数
- 最佳实践建议
在Python定时任务中,任务堆积(任务还没执行完,下一个调度周期又到了)是个常见问题,以下是几种解决方案:
使用单线程串行执行(最基础)
import schedule
import time
from datetime import datetime
def my_task():
print(f"开始执行: {datetime.now()}")
time.sleep(10) # 模拟耗时任务
print(f"执行结束: {datetime.now()}")
# 每5秒执行一次,但如果任务执行10秒,不会启动新任务
schedule.every(5).seconds.do(my_task)
while True:
schedule.run_pending()
time.sleep(1) # 这里的sleep不会影响任务执行
缺点:即使任务超时,也不会跳过,只是延迟执行
使用APScheduler的"间隔"模式(推荐)
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.interval import IntervalTrigger
import time
from datetime import datetime
def task_with_cooldown():
print(f"执行任务: {datetime.now()}")
time.sleep(10) # 模拟耗时
scheduler = BackgroundScheduler()
# 关键:设置 interval 模式,任务执行完才开始计时
scheduler.add_job(
task_with_cooldown,
trigger=IntervalTrigger(seconds=5), # 任务结束后5秒再执行
id='my_task',
replace_existing=True
)
scheduler.start()
try:
while True:
time.sleep(1)
except KeyboardInterrupt:
scheduler.shutdown()
使用锁控制并发(APScheduler)
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.interval import IntervalTrigger
import threading
import time
class NonOverlappingTask:
def __init__(self):
self.lock = threading.Lock()
self.is_running = False
def execute(self):
if not self.lock.acquire(blocking=False): # 非阻塞获取锁
print("任务正在执行中,跳过本次调度")
return
try:
self.is_running = True
print("开始执行任务...")
time.sleep(10) # 模拟耗时
print("任务执行完成")
finally:
self.is_running = False
self.lock.release()
task = NonOverlappingTask()
scheduler = BackgroundScheduler()
# 即使调度的间隔是3秒,只要任务没完成就跳过
scheduler.add_job(
task.execute,
trigger='interval', # 固定间隔模式
seconds=3,
id='non_overlap_task',
misfire_grace_time=None # 不允许追赶执行
)
scheduler.start()
使用Celery的自动重试和限流
# tasks.py
from celery import Celery
from celery.utils.log import get_task_logger
app = Celery('tasks', broker='redis://localhost:6379/0')
logger = get_task_logger(__name__)
@app.task(
bind=True,
max_retries=3,
default_retry_delay=60,
rate_limit='1/m' # 每分钟最多1次
)
def my_task(self):
try:
# 业务逻辑
import time
time.sleep(5)
logger.info("任务执行成功")
except Exception as exc:
logger.error(f"任务失败: {exc}")
raise self.retry(exc=exc)
# 定时调用
@app.on_after_configure.connect
def setup_periodic_tasks(sender, **kwargs):
# 每10秒调度一次,但受rate_limit限制
sender.add_periodic_task(10.0, my_task.s(), name='every 10s')
使用Redis分布式锁(多进程场景)
import redis
import time
import uuid
from datetime import datetime
class DistributedTaskLock:
def __init__(self, redis_client, lock_key="task_lock"):
self.redis = redis_client
self.lock_key = lock_key
self.lock_value = str(uuid.uuid4())
self.lock_timeout = 30 # 锁超时时间(秒)
def acquire_lock(self):
"""尝试获取锁"""
return self.redis.set(
self.lock_key,
self.lock_value,
nx=True, # 只有key不存在时才能设置
ex=self.lock_timeout # 自动过期
)
def release_lock(self):
"""释放锁(只释放自己的锁)"""
# Lua脚本保证原子性
release_script = """
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
else
return 0
end
"""
self.redis.eval(release_script, 1, self.lock_key, self.lock_value)
def execute_with_lock(self, func, *args, **kwargs):
if self.acquire_lock():
try:
print(f"获取锁成功,执行任务: {datetime.now()}")
return func(*args, **kwargs)
finally:
self.release_lock()
else:
print(f"获取锁失败,跳过任务: {datetime.now()}")
return None
# 使用示例
def my_task():
time.sleep(10) # 模拟耗时任务
return "任务完成"
redis_client = redis.Redis(host='localhost', port=6379, db=0)
lock = DistributedTaskLock(redis_client)
# 每5秒尝试执行一次
while True:
lock.execute_with_lock(my_task)
time.sleep(5)
使用信号量控制并发数
import asyncio
import time
from datetime import datetime
class SemaphoreTask:
def __init__(self, max_concurrent=1):
self.semaphore = asyncio.Semaphore(max_concurrent)
self.task_counter = 0
async def execute_task(self, task_id):
async with self.semaphore:
print(f"任务{task_id}开始: {datetime.now()}")
await asyncio.sleep(5) # 模拟耗时
print(f"任务{task_id}结束: {datetime.now()}")
async def scheduler(self):
"""模拟定时调度"""
while True:
self.task_counter += 1
task_id = self.task_counter
# 非阻塞方式尝试执行
if self.semaphore.locked():
print(f"任务{task_id}被跳过(有任务在执行中)")
else:
asyncio.create_task(self.execute_task(task_id))
await asyncio.sleep(3) # 每3秒调度一次
# 运行
async def main():
scheduler = SemaphoreTask(max_concurrent=1)
await scheduler.scheduler()
# asyncio.run(main())
最佳实践建议
-
优先使用APScheduler:它原生支持任务不重叠执行
-
设置合理的重试策略:使用
misfire_grace_time控制跳过还是追赶 -
监控任务状态:
from apscheduler.events import EVENT_JOB_EXECUTED, EVENT_JOB_ERROR def job_listener(event): if event.exception: print(f"任务失败: {event.job_id}") else: print(f"任务成功: {event.job_id}") scheduler.add_listener(job_listener, EVENT_JOB_EXECUTED | EVENT_JOB_ERROR) -
使用分布式方案:当有多个服务实例时,用Redis/Celery保证全局只有一个实例执行
选择方案时考虑:
- 单机简单场景:方案1或方案2
- 需要跳过堆积:方案3(锁控制)
- 分布式系统:方案4或方案5
- 高并发控制:方案6(信号量)