Python脚本如何实现异步定时任务

wen python案例 25

本文目录导读:

Python脚本如何实现异步定时任务

  1. 使用 asyncio + 内置库
  2. 使用 apscheduler 库(推荐)
  3. 使用 aiocron 库
  4. 自定义灵活调度器
  5. 推荐方案

我来介绍几种Python实现异步定时任务的方法:

使用 asyncio + 内置库

基础方案:asyncio.sleep() 循环

import asyncio
import time
async def scheduled_task(name, interval, count):
    """定时执行的任务"""
    for i in range(count):
        print(f"[{time.strftime('%H:%M:%S')}] {name} - 第 {i+1} 次执行")
        await asyncio.sleep(interval)
    print(f"{name} 执行完毕")
async def main():
    # 创建多个定时任务
    tasks = [
        scheduled_task("任务A", 2, 5),  # 每2秒执行一次,共5次
        scheduled_task("任务B", 3, 3),  # 每3秒执行一次,共3次
    ]
    # 并发执行所有任务
    await asyncio.gather(*tasks)
# 运行
asyncio.run(main())

使用 asyncio 的循环和回调

import asyncio
import time
class AsyncScheduler:
    def __init__(self):
        self.tasks = []
    def add_task(self, interval, callback, *args, **kwargs):
        """添加定时任务"""
        async def task_wrapper():
            while True:
                await asyncio.sleep(interval)
                if callable(callback):
                    await callback(*args, **kwargs)
        self.tasks.append(task_wrapper())
    async def run(self):
        """运行所有定时任务"""
        await asyncio.gather(*self.tasks)
async def my_task(name):
    print(f"[{time.strftime('%H:%M:%S')}] {name} 执行了")
# 使用示例
async def main():
    scheduler = AsyncScheduler()
    # 添加定时任务
    scheduler.add_task(2, my_task, "任务A")  # 每2秒执行
    scheduler.add_task(5, my_task, "任务B")  # 每5秒执行
    # 运行调度器
    await scheduler.run()
# asyncio.run(main())

使用 apscheduler 库(推荐)

安装

pip install apscheduler

异步模式使用

from apscheduler.schedulers.asyncio import AsyncIOScheduler
from apscheduler.triggers.interval import IntervalTrigger
from apscheduler.triggers.cron import CronTrigger
import asyncio
import time
async def my_async_task(name):
    print(f"[{time.strftime('%H:%M:%S')}] 异步任务 {name} 执行")
    await asyncio.sleep(1)  # 模拟异步操作
    print(f"[{time.strftime('%H:%M:%S')}] 任务 {name} 完成")
async def main():
    # 创建调度器
    scheduler = AsyncIOScheduler()
    # 方法1:使用装饰器
    @scheduler.scheduled_job('interval', seconds=3, id='job1')
    async def job1():
        await my_async_task("定时任务1")
    # 方法2:手动添加任务
    scheduler.add_job(
        my_async_task,
        trigger=IntervalTrigger(seconds=5),
        args=["定时任务2"],
        id='job2',
        max_instances=3  # 最大并发实例数
    )
    # 方法3:使用cron表达式
    scheduler.add_job(
        my_async_task,
        trigger=CronTrigger(second='*/10'),  # 每10秒执行
        args=["定时任务3"],
        id='job3'
    )
    # 启动调度器
    scheduler.start()
    print("调度器已启动,按 Ctrl+C 停止")
    try:
        # 保持程序运行
        await asyncio.Event().wait()
    except (KeyboardInterrupt, SystemExit):
        print("正在关闭调度器...")
        scheduler.shutdown()
# 运行
# asyncio.run(main())

更完整的 APScheduler 示例

from apscheduler.schedulers.asyncio import AsyncIOScheduler
from apscheduler.triggers.interval import IntervalTrigger
from apscheduler.triggers.cron import CronTrigger
from apscheduler.triggers.date import DateTrigger
from datetime import datetime, timedelta
import asyncio
import logging
# 配置日志
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
class AsyncTaskScheduler:
    def __init__(self):
        self.scheduler = AsyncIOScheduler()
        self.tasks = {}
    def add_interval_task(self, task_id, func, seconds, args=None, kwargs=None):
        """添加间隔任务"""
        self.scheduler.add_job(
            func,
            trigger=IntervalTrigger(seconds=seconds),
            id=task_id,
            args=args or [],
            kwargs=kwargs or {},
            replace_existing=True
        )
        logger.info(f"添加间隔任务: {task_id}, 每 {seconds} 秒执行")
    def add_cron_task(self, task_id, func, cron_expr, args=None, kwargs=None):
        """添加cron定时任务"""
        self.scheduler.add_job(
            func,
            trigger=CronTrigger.from_crontab(cron_expr),
            id=task_id,
            args=args or [],
            kwargs=kwargs or {},
            replace_existing=True
        )
        logger.info(f"添加cron任务: {task_id}, 表达式: {cron_expr}")
    def add_one_time_task(self, task_id, func, run_date, args=None, kwargs=None):
        """添加一次性任务"""
        self.scheduler.add_job(
            func,
            trigger=DateTrigger(run_date=run_date),
            id=task_id,
            args=args or [],
            kwargs=kwargs or {},
            replace_existing=True
        )
        logger.info(f"添加一次性任务: {task_id}, 执行时间: {run_date}")
    def remove_task(self, task_id):
        """移除任务"""
        try:
            self.scheduler.remove_job(task_id)
            logger.info(f"移除任务: {task_id}")
        except Exception as e:
            logger.error(f"移除任务失败: {e}")
    def pause_task(self, task_id):
        """暂停任务"""
        self.scheduler.pause_job(task_id)
        logger.info(f"暂停任务: {task_id}")
    def resume_task(self, task_id):
        """恢复任务"""
        self.scheduler.resume_job(task_id)
        logger.info(f"恢复任务: {task_id}")
    async def start(self):
        """启动调度器"""
        self.scheduler.start()
        logger.info("调度器已启动")
        try:
            await asyncio.Event().wait()
        except (KeyboardInterrupt, SystemExit):
            logger.info("正在关闭调度器...")
        finally:
            self.scheduler.shutdown()
            logger.info("调度器已关闭")
# 使用示例
async def data_collection_task(source_name):
    """模拟数据采集任务"""
    logger.info(f"开始从 {source_name} 采集数据...")
    await asyncio.sleep(2)  # 模拟采集过程
    logger.info(f"从 {source_name} 采集数据完成")
async def report_generation_task():
    """模拟报表生成任务"""
    logger.info("开始生成报表...")
    await asyncio.sleep(3)
    logger.info("报表生成完成")
async def main():
    scheduler = AsyncTaskScheduler()
    # 添加各种定时任务
    scheduler.add_interval_task(
        "data_collector",
        data_collection_task,
        seconds=10,
        args=["API数据源"]
    )
    scheduler.add_cron_task(
        "daily_report",
        report_generation_task,
        cron_expr="0 9 * * *"  # 每天9点执行
    )
    # 添加一次性任务(5分钟后执行)
    run_time = datetime.now() + timedelta(minutes=5)
    scheduler.add_one_time_task(
        "one_time_job",
        data_collection_task,
        run_date=run_time,
        args=["一次性数据源"]
    )
    # 启动调度器
    await scheduler.start()
if __name__ == "__main__":
    asyncio.run(main())

使用 aiocron 库

安装

pip install aiocron

使用示例

import aiocron
import asyncio
import time
@aiocron.crontab('*/5 * * * * *')  # 每5秒执行
async def cron_job1():
    print(f"[{time.strftime('%H:%M:%S')}] cron任务1执行")
@aiocron.crontab('*/10 * * * * *')  # 每10秒执行
async def cron_job2():
    print(f"[{time.strftime('%H:%M:%S')}] cron任务2执行")
async def main():
    print("aiocron调度器启动")
    # 保持运行
    while True:
        await asyncio.sleep(1)
# 运行
# asyncio.run(main())

自定义灵活调度器

import asyncio
import time
from typing import Callable, Any, Optional
from dataclasses import dataclass
from datetime import datetime
@dataclass
class Task:
    """任务数据类"""
    name: str
    func: Callable
    interval: float
    args: tuple = ()
    kwargs: dict = None
    last_run: Optional[float] = None
    running: bool = False
class FlexibleAsyncScheduler:
    def __init__(self):
        self.tasks = {}
        self._running = False
    def add_task(self, name: str, func: Callable, interval: float, 
                 *args, **kwargs):
        """添加定时任务"""
        self.tasks[name] = Task(
            name=name,
            func=func,
            interval=interval,
            args=args,
            kwargs=kwargs
        )
    def remove_task(self, name: str):
        """移除任务"""
        if name in self.tasks:
            del self.tasks[name]
    async def _run_task(self, task: Task):
        """运行单个任务"""
        if task.running:
            return  # 如果任务正在运行,跳过
        task.running = True
        try:
            await task.func(*task.args, **(task.kwargs or {}))
        except Exception as e:
            print(f"任务 {task.name} 执行失败: {e}")
        finally:
            task.running = False
            task.last_run = time.time()
    async def _scheduler_loop(self):
        """调度器主循环"""
        while self._running:
            current_time = time.time()
            for task in self.tasks.values():
                # 检查是否需要执行
                if task.last_run is None:
                    should_run = True
                else:
                    should_run = (current_time - task.last_run) >= task.interval
                if should_run:
                    # 异步执行任务
                    asyncio.create_task(self._run_task(task))
            # 每秒检查一次
            await asyncio.sleep(1)
    async def start(self):
        """启动调度器"""
        self._running = True
        await self._scheduler_loop()
    def stop(self):
        """停止调度器"""
        self._running = False
# 使用示例
async def monitor_cpu():
    print(f"[{time.strftime('%H:%M:%S')}] 监控CPU使用率...")
    await asyncio.sleep(0.5)
async def check_memory():
    print(f"[{time.strftime('%H:%M:%S')}] 检查内存使用...")
    await asyncio.sleep(0.5)
async def main():
    scheduler = FlexibleAsyncScheduler()
    # 添加定时任务
    scheduler.add_task("cpu_monitor", monitor_cpu, interval=3)
    scheduler.add_task("memory_check", check_memory, interval=5)
    print("自定义调度器启动,按 Ctrl+C 停止")
    try:
        await scheduler.start()
    except KeyboardInterrupt:
        scheduler.stop()
        print("调度器已停止")
if __name__ == "__main__":
    asyncio.run(main())

推荐方案

  1. 简单场景:使用 asyncio.sleep() 循环
  2. 生产环境:使用 apscheduler 库,功能最完善
  3. cron表达式:使用 aiocron
  4. 需要灵活控制:使用自定义调度器

选择哪种方案取决于你的具体需求,包括任务复杂性、并发要求、错误处理需求等。

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