本文目录导读:

我来介绍几种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())
推荐方案
- 简单场景:使用
asyncio.sleep()循环 - 生产环境:使用
apscheduler库,功能最完善 - cron表达式:使用
aiocron库 - 需要灵活控制:使用自定义调度器
选择哪种方案取决于你的具体需求,包括任务复杂性、并发要求、错误处理需求等。