Python脚本如何避免定时任务堆积

wen python案例 23

本文目录导读:

Python脚本如何避免定时任务堆积

  1. 使用单线程串行执行(最基础)
  2. ">使用APScheduler的"间隔"模式(推荐)
  3. 使用锁控制并发(APScheduler)
  4. 使用Celery的自动重试和限流
  5. 使用Redis分布式锁(多进程场景)
  6. 使用信号量控制并发数
  7. 最佳实践建议

在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())

最佳实践建议

  1. 优先使用APScheduler:它原生支持任务不重叠执行

  2. 设置合理的重试策略:使用 misfire_grace_time 控制跳过还是追赶

  3. 监控任务状态

    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)
  4. 使用分布式方案:当有多个服务实例时,用Redis/Celery保证全局只有一个实例执行

选择方案时考虑:

  • 单机简单场景:方案1或方案2
  • 需要跳过堆积:方案3(锁控制)
  • 分布式系统:方案4或方案5
  • 高并发控制:方案6(信号量)

抱歉,评论功能暂时关闭!