Python 定时任务与调度指南
概述
定时任务是自动化办公的核心——让程序在指定时间自动执行,无需人工干预。从简单的定时备份到复杂的分布式任务调度,Python 生态提供了多层次的解决方案。
图表渲染中…
| 方案 | 复杂度 | 持久化 | 分布式 | 适用场景 |
|---|---|---|---|---|
| schedule | ⭐ | ❌ | ❌ | 简单定时脚本 |
| APScheduler | ⭐⭐⭐ | ✅ | ⚠️ 有限 | Web 应用、中等规模 |
| Celery | ⭐⭐⭐⭐⭐ | ✅ | ✅ | 大规模分布式 |
| cron | ⭐⭐ | ✅ | ❌ | Linux 系统级 |
| Windows 任务计划 | ⭐⭐ | ✅ | ❌ | Windows 系统级 |
schedule:轻量级调度
bash
pip install schedule基础用法
python
import schedule
import time
def job():
print('执行任务...')
# 固定间隔
schedule.every(10).minutes.do(job)
schedule.every(1).hour.do(job)
schedule.every(1).day.do(job)
schedule.every(1).week.do(job)
# 指定时间
schedule.every().day.at('09:00').do(job)
schedule.every().monday.at('09:30').do(job)
schedule.every().wednesday.at('18:00').do(job)
# 直到某个时间
schedule.every().day.until('23:59').do(job)
# 运行调度器
while True:
schedule.run_pending()
time.sleep(1)高级用法
python
import schedule
import time
import logging
from datetime import datetime
# 带参数的任务
def send_report(recipient, subject):
print(f'发送报告给 {recipient}: {subject}')
schedule.every().day.at('09:00').do(send_report, recipient='admin@company.com', subject='日报')
# 取消任务
job_obj = schedule.every(5).minutes.do(job)
schedule.cancel_job(job_obj)
# 清除所有任务
schedule.clear()
# 装饰器语法
@schedule.repeat(schedule.every(1).hour)
def hourly_cleanup():
print('每小时清理...')
# 任务标签(分组管理)
schedule.every(5).minutes.do(job).tag('maintenance')
schedule.every(1).hour.do(job).tag('reports')
# 按标签取消
schedule.clear('maintenance')
# 获取待执行任务
for job in schedule.get_jobs():
print(f'任务: {job.job_func.__name__}, 下次执行: {job.next_run}')
# 错误处理
def safe_job(job_func):
"""装饰器:捕获任务异常,防止调度器崩溃"""
def wrapper(*args, **kwargs):
try:
return job_func(*args, **kwargs)
except Exception as e:
logging.error(f'任务 {job_func.__name__} 执行失败: {e}')
return wrapper
@safe_job
def risky_job():
1 / 0 # 不会导致调度器崩溃
schedule.every(1).minute.do(risky_job)优雅退出
python
import schedule
import signal
import sys
running = True
def signal_handler(signum, frame):
global running
print('\n收到退出信号,正在优雅停止...')
running = False
signal.signal(signal.SIGINT, signal_handler)
signal.signal(signal.SIGTERM, signal_handler)
schedule.every(10).seconds.do(job)
while running:
schedule.run_pending()
time.sleep(1)
print('调度器已停止')APScheduler:企业级调度
bash
pip install apscheduler三大核心概念
图表渲染中…
基础用法
python
from apscheduler.schedulers.blocking import BlockingScheduler
from apscheduler.schedulers.background import BackgroundScheduler
from datetime import datetime
# 阻塞式调度器(适合独立脚本)
scheduler = BlockingScheduler()
# 后台调度器(适合 Web 应用)
# scheduler = BackgroundScheduler()
# Date 触发器(一次性任务)
from apscheduler.triggers.date import DateTrigger
scheduler.add_job(job, trigger='date', run_date='2026-06-07 09:00:00')
# Interval 触发器(固定间隔)
scheduler.add_job(job, 'interval', minutes=30)
scheduler.add_job(job, 'interval', hours=1, start_date='2026-06-07 08:00', end_date='2026-06-30')
# Cron 触发器(Cron 表达式)
scheduler.add_job(job, 'cron', hour=9, minute=0) # 每天 9:00
scheduler.add_job(job, 'cron', day_of_week='mon-fri', hour=9) # 工作日 9:00
scheduler.add_job(job, 'cron', day=1, hour=0) # 每月1号 0:00
scheduler.add_job(job, 'cron', hour='9,18') # 每天 9:00 和 18:00
scheduler.add_job(job, 'cron', minute='*/15') # 每15分钟
scheduler.add_job(job, 'cron', day_of_week='mon', hour=9, minute=30) # 每周一 9:30
# 带参数
scheduler.add_job(send_report, 'cron', hour=9, args=['admin@company.com'], kwargs={'subject': '日报'})
# 任务管理
job = scheduler.add_job(job, 'interval', minutes=10, id='my_job', name='我的任务')
scheduler.remove_job('my_job')
scheduler.pause_job('my_job')
scheduler.resume_job('my_job')
scheduler.modify_job('my_job', minutes=5) # 修改间隔
# 启动
scheduler.start()持久化存储
python
from apscheduler.schedulers.blocking import BlockingScheduler
from apscheduler.jobstores.sqlalchemy import SQLAlchemyJobStore
from apscheduler.executors.pool import ThreadPoolExecutor, ProcessPoolExecutor
# 配置持久化
jobstores = {
'default': SQLAlchemyJobStore(url='sqlite:///jobs.sqlite'),
# 也可以用 MySQL/PostgreSQL
# 'default': SQLAlchemyJobStore(url='mysql+pymysql://user:pass@localhost/scheduler'),
}
executors = {
'default': ThreadPoolExecutor(20),
'processpool': ProcessPoolExecutor(5),
}
job_defaults = {
'coalesce': True, # 合并错过的执行
'max_instances': 3, # 同一任务最大并发实例数
'misfire_grace_time': 300, # 错过执行的容忍时间(秒)
}
scheduler = BlockingScheduler(
jobstores=jobstores,
executors=executors,
job_defaults=job_defaults,
)
# 添加任务时指定执行器
scheduler.add_job(cpu_intensive_task, 'cron', hour=2,
executor='processpool')事件监听
python
from apscheduler.events import EVENT_JOB_EXECUTED, EVENT_JOB_ERROR, EVENT_JOB_MISSED
def job_executed(event):
print(f'任务 {event.job_id} 执行成功')
def job_error(event):
print(f'任务 {event.job_id} 执行失败: {event.exception}')
def job_missed(event):
print(f'任务 {event.job_id} 错过执行')
scheduler.add_listener(job_executed, EVENT_JOB_EXECUTED)
scheduler.add_listener(job_error, EVENT_JOB_ERROR)
scheduler.add_listener(job_missed, EVENT_JOB_MISSED)Flask/Django 集成
python
# Flask 集成
from flask import Flask
from apscheduler.schedulers.background import BackgroundScheduler
app = Flask(__name__)
scheduler = BackgroundScheduler()
def scheduled_task():
with app.app_context():
# 在 Flask 应用上下文中执行
pass
scheduler.add_job(scheduled_task, 'cron', hour=9)
scheduler.start()
# 确保应用退出时关闭调度器
import atexit
atexit.register(lambda: scheduler.shutdown())
@app.route('/jobs')
def list_jobs():
jobs = []
for job in scheduler.get_jobs():
jobs.append({
'id': job.id,
'name': job.name,
'next_run': str(job.next_run_time),
})
return {'jobs': jobs}
# Django 集成类似,使用 django-apscheduler 包Celery:分布式任务队列
bash
pip install celery redis
# 需要安装 Redis 作为 Broker基础架构
图表渲染中…
配置与使用
python
# celery_app.py
from celery import Celery
from celery.schedules import crontab
app = Celery(
'my_tasks',
broker='redis://localhost:6379/0',
backend='redis://localhost:6379/1',
)
# 配置
app.conf.update(
task_serializer='json',
accept_content=['json'],
result_serializer='json',
timezone='Asia/Shanghai',
enable_utc=True,
task_track_started=True,
task_time_limit=3600, # 硬超时 1 小时
task_soft_time_limit=3000, # 软超时 50 分钟
worker_max_tasks_per_child=1000, # 每个 worker 处理 N 个任务后重启
worker_prefetch_multiplier=1, # 预取倍数
)
# 定义任务
@app.task(bind=True, max_retries=3)
def send_email_task(self, recipient, subject, body):
"""异步发送邮件任务"""
try:
# 发送邮件逻辑...
return {'status': 'success', 'recipient': recipient}
except Exception as exc:
raise self.retry(exc=exc, countdown=60) # 60 秒后重试
@app.task
def generate_report_task(report_type, date_range):
"""异步生成报表任务"""
# 报表生成逻辑...
return {'report_url': 'https://...', 'status': 'completed'}
# 定时任务(Celery Beat)
app.conf.beat_schedule = {
'daily-report': {
'task': 'celery_app.generate_report_task',
'schedule': crontab(hour=9, minute=0),
'args': ('daily', 'today'),
},
'hourly-healthcheck': {
'task': 'celery_app.healthcheck',
'schedule': crontab(minute=0), # 每小时整点
},
'weekly-backup': {
'task': 'celery_app.backup_database',
'schedule': crontab(day_of_week=0, hour=2), # 每周日凌晨2点
},
}
@app.task
def healthcheck():
"""健康检查任务"""
import psutil
return {
'cpu': psutil.cpu_percent(),
'memory': psutil.virtual_memory().percent,
'disk': psutil.disk_usage('/').percent,
}
@app.task
def backup_database():
"""数据库备份任务"""
# 备份逻辑...
return {'status': 'completed', 'file': 'backup_20260606.sql.gz'}调用与监控
python
# 调用异步任务
from celery_app import send_email_task, generate_report_task
# 方式一:delay(简单调用)
result = send_email_task.delay('admin@company.com', '测试', '正文')
# 方式二:apply_async(高级调用)
result = send_email_task.apply_async(
args=['admin@company.com', '测试', '正文'],
countdown=60, # 60 秒后执行
expires=3600, # 1 小时后过期
retry=True,
priority=5, # 优先级(0-9,0 最高)
queue='emails', # 指定队列
)
# 获取结果
print(result.id) # 任务 ID
print(result.status) # PENDING / STARTED / SUCCESS / FAILURE
print(result.get(timeout=30)) # 阻塞等待结果
print(result.ready()) # 是否完成
# 批量任务
from celery import group
job = group([
send_email_task.s(email, '周报', content)
for email in ['a@co.com', 'b@co.com', 'c@co.com']
])
result = job.apply_async()
# 链式任务
from celery import chain
workflow = chain(
generate_report_task.s('monthly', '2026-05'),
send_email_task.s('admin@co.com', '月报'),
)
workflow.apply_async()
# 启动 Worker
# celery -A celery_app worker --loglevel=info --concurrency=4
# 启动 Beat(定时调度)
# celery -A celery_app beat --loglevel=info
# 同时启动 Worker + Beat
# celery -A celery_app worker -B --loglevel=info系统级调度
Cron(Linux/macOS)
bash
# 编辑 crontab
crontab -e
# Cron 表达式格式:
# ┌──────── 分钟 (0-59)
# │ ┌────── 小时 (0-23)
# │ │ ┌──── 日 (1-31)
# │ │ │ ┌── 月 (1-12)
# │ │ │ │ ┌ 星期 (0-6, 0=周日)
# * * * * * command
# 示例
# 每天凌晨2点备份数据库
0 2 * * * /usr/bin/python3 /opt/scripts/backup.py >> /var/log/backup.log 2>&1
# 工作日9点发送日报
0 9 * * 1-5 /usr/bin/python3 /opt/scripts/daily_report.py
# 每5分钟检查服务状态
*/5 * * * * /usr/bin/python3 /opt/scripts/healthcheck.py
# 每月1号清理过期文件
0 0 1 * * /usr/bin/python3 /opt/scripts/cleanup.pyPython 管理 Cron
bash
pip install python-crontabpython
from crontab import CronTab
# 使用当前用户的 crontab
cron = CronTab(user=True)
# 添加任务
job = cron.new(command='/usr/bin/python3 /opt/scripts/backup.py')
job.minute.on(0)
job.hour.on(2)
# 等价于: 0 2 * * *
# 使用 cron 表达式
job2 = cron.new(command='/usr/bin/python3 /opt/scripts/report.py', comment='daily_report')
job2.setall('0 9 * * 1-5')
# 列出所有任务
for job in cron:
print(job)
# 按注释查找并删除
for job in cron.find_comment('daily_report'):
cron.remove(job)
cron.write() # 保存Windows 任务计划
python
# 使用 subprocess 调用 schtasks
import subprocess
def create_windows_task(task_name, script_path, schedule='daily', time='09:00'):
"""创建 Windows 计划任务"""
cmd = [
'schtasks', '/create',
'/tn', task_name,
'/tr', f'python {script_path}',
'/sc', schedule, # daily, weekly, monthly
'/st', time,
'/f', # 覆盖已存在的任务
]
subprocess.run(cmd, check=True)
def delete_windows_task(task_name):
"""删除 Windows 计划任务"""
subprocess.run(['schtasks', '/delete', '/tn', task_name, '/f'], check=True)
def list_windows_tasks():
"""列出所有计划任务"""
result = subprocess.run(['schtasks', '/query', '/fo', 'CSV'],
capture_output=True, text=True)
return result.stdout实战案例:定时巡检与告警系统
python
"""
定时巡检与告警系统
功能:自动巡检服务器状态,异常时多渠道告警
支持:APScheduler 调度 + 企业微信/钉钉/邮件通知
"""
import psutil
import platform
import smtplib
import json
import logging
from datetime import datetime
from dataclasses import dataclass, asdict
from pathlib import Path
from apscheduler.schedulers.blocking import BlockingScheduler
from apscheduler.triggers.cron import CronTrigger
from apscheduler.triggers.interval import IntervalTrigger
# 配置日志
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s [%(levelname)s] %(message)s',
handlers=[
logging.FileHandler('patrol.log', encoding='utf-8'),
logging.StreamHandler(),
]
)
@dataclass
class AlertConfig:
"""告警配置"""
cpu_threshold: float = 80.0 # CPU 使用率阈值
memory_threshold: float = 85.0 # 内存使用率阈值
disk_threshold: float = 90.0 # 磁盘使用率阈值
wecom_key: str = '' # 企业微信 Webhook Key
dingtalk_url: str = '' # 钉钉 Webhook URL
smtp_host: str = ''
smtp_port: int = 465
smtp_user: str = ''
smtp_password: str = ''
alert_email: str = ''
@dataclass
class PatrolResult:
"""巡检结果"""
timestamp: str
hostname: str
cpu_percent: float
memory_percent: float
disk_percent: float
disk_free_gb: float
load_avg: tuple
network_connections: int
alerts: list
class PatrolSystem:
"""巡检系统"""
def __init__(self, config: AlertConfig):
self.config = config
self.history: list = []
self.scheduler = BlockingScheduler()
def check_system(self) -> PatrolResult:
"""执行系统检查"""
cpu = psutil.cpu_percent(interval=1)
mem = psutil.virtual_memory()
disk = psutil.disk_usage('/')
load = psutil.getloadavg() if hasattr(psutil, 'getloadavg') else (0, 0, 0)
net_conns = len(psutil.net_connections())
alerts = []
if cpu > self.config.cpu_threshold:
alerts.append(f'🔴 CPU 使用率 {cpu:.1f}% 超过阈值 {self.config.cpu_threshold}%')
if mem.percent > self.config.memory_threshold:
alerts.append(f'🔴 内存使用率 {mem.percent:.1f}% 超过阈值 {self.config.memory_threshold}%')
if disk.percent > self.config.disk_threshold:
alerts.append(f'🔴 磁盘使用率 {disk.percent:.1f}% 超过阈值 {self.config.disk_threshold}%')
result = PatrolResult(
timestamp=datetime.now().isoformat(),
hostname=platform.node(),
cpu_percent=cpu,
memory_percent=mem.percent,
disk_percent=disk.percent,
disk_free_gb=disk.free / (1024**3),
load_avg=load,
network_connections=net_conns,
alerts=alerts,
)
self.history.append(asdict(result))
return result
def send_alert(self, result: PatrolResult):
"""发送告警通知"""
if not result.alerts:
return
# 构建告警内容
content = f'''## 🚨 服务器告警
**主机**: {result.hostname}
**时间**: {result.timestamp}
| 指标 | 数值 | 状态 |
|------|------|------|
| CPU | {result.cpu_percent:.1f}% | {'🔴' if result.cpu_percent > self.config.cpu_threshold else '🟢'} |
| 内存 | {result.memory_percent:.1f}% | {'🔴' if result.memory_percent > self.config.memory_threshold else '🟢'} |
| 磁盘 | {result.disk_percent:.1f}% | {'🔴' if result.disk_percent > self.config.disk_threshold else '🟢'} |
| 磁盘剩余 | {result.disk_free_gb:.1f} GB | — |
**告警详情**:
'''
for alert in result.alerts:
content += f'- {alert}\n'
# 企业微信通知
if self.config.wecom_key:
self._send_wecom(content)
# 钉钉通知
if self.config.dingtalk_url:
self._send_dingtalk(content)
# 邮件通知
if self.config.alert_email and self.config.smtp_host:
self._send_email(content)
def _send_wecom(self, content):
"""企业微信通知"""
import requests
url = f'https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key={self.config.wecom_key}'
payload = {
'msgtype': 'markdown',
'markdown': {'content': content},
}
try:
requests.post(url, json=payload, timeout=5)
logging.info('企业微信通知已发送')
except Exception as e:
logging.error(f'企业微信通知失败: {e}')
def _send_dingtalk(self, content):
"""钉钉通知"""
import requests
payload = {
'msgtype': 'markdown',
'markdown': {
'title': '服务器告警',
'text': content,
},
}
try:
requests.post(self.config.dingtalk_url, json=payload, timeout=5)
logging.info('钉钉通知已发送')
except Exception as e:
logging.error(f'钉钉通知失败: {e}')
def _send_email(self, content):
"""邮件通知"""
from email.mime.text import MIMEText
msg = MIMEText(content, 'html', 'utf-8')
msg['Subject'] = f'🚨 服务器告警 - {platform.node()}'
msg['From'] = self.config.smtp_user
msg['To'] = self.config.alert_email
try:
with smtplib.SMTP_SSL(self.config.smtp_host, self.config.smtp_port) as server:
server.login(self.config.smtp_user, self.config.smtp_password)
server.send_message(msg)
logging.info('邮件通知已发送')
except Exception as e:
logging.error(f'邮件通知失败: {e}')
def patrol_job(self):
"""巡检任务(被调度器调用)"""
logging.info('开始巡检...')
result = self.check_system()
# 记录状态
status = '正常' if not result.alerts else '告警'
logging.info(f'巡检完成: CPU={result.cpu_percent:.1f}%, '
f'内存={result.memory_percent:.1f}%, '
f'磁盘={result.disk_percent:.1f}% — {status}')
# 有告警时发送通知
if result.alerts:
self.send_alert(result)
def daily_report_job(self):
"""每日巡检报告"""
if not self.history:
return
today = datetime.now().strftime('%Y-%m-%d')
today_records = [r for r in self.history if r['timestamp'].startswith(today)]
if not today_records:
return
avg_cpu = sum(r['cpu_percent'] for r in today_records) / len(today_records)
avg_mem = sum(r['memory_percent'] for r in today_records) / len(today_records)
max_cpu = max(r['cpu_percent'] for r in today_records)
alert_count = sum(1 for r in today_records if r['alerts'])
report = f'''## 📊 每日巡检报告
**日期**: {today}
**主机**: {platform.node()}
| 指标 | 平均值 | 最大值 |
|------|--------|--------|
| CPU | {avg_cpu:.1f}% | {max_cpu:.1f}% |
| 内存 | {avg_mem:.1f}% | — |
| 告警次数 | {alert_count} | — |
'''
logging.info(f'每日报告生成: CPU均值={avg_cpu:.1f}%, 告警={alert_count}次')
def save_history(self, filepath='patrol_history.json'):
"""保存历史记录"""
Path(filepath).write_text(
json.dumps(self.history[-1000:], ensure_ascii=False, indent=2),
encoding='utf-8'
)
def run(self):
"""启动巡检系统"""
# 每 5 分钟巡检一次
self.scheduler.add_job(
self.patrol_job,
IntervalTrigger(minutes=5),
id='patrol',
name='系统巡检',
)
# 每天 9:00 发送日报
self.scheduler.add_job(
self.daily_report_job,
CronTrigger(hour=9, minute=0),
id='daily_report',
name='每日报告',
)
# 每天凌晨清理历史(保留最近 7 天)
self.scheduler.add_job(
lambda: self.history.__delitem__(slice(0, -1008)), # 7天 × 24小时 × 6次/小时
CronTrigger(hour=3, minute=0),
id='cleanup',
name='历史清理',
)
logging.info('巡检系统启动')
try:
self.scheduler.start()
except (KeyboardInterrupt, SystemExit):
self.save_history()
logging.info('巡检系统已停止')
# 使用
config = AlertConfig(
cpu_threshold=80,
memory_threshold=85,
disk_threshold=90,
wecom_key='YOUR_WECOM_KEY',
smtp_host='smtp.qq.com',
smtp_port=465,
smtp_user='your_email@qq.com',
smtp_password='your_auth_code',
alert_email='admin@company.com',
)
system = PatrolSystem(config)
system.run()常见陷阱
| 陷阱 | 说明 | 正确做法 |
|---|---|---|
| schedule 阻塞主线程 | while True 循环占用主线程 | 使用 APScheduler 后台调度或独立进程 |
| 任务异常导致调度停止 | 未捕获的任务异常 | 使用 try/except 包裹任务或 APScheduler 错误监听 |
| 时区问题 | UTC vs 本地时间 | 显式设置 timezone='Asia/Shanghai' |
| 持久化丢失 | 内存存储重启后任务丢失 | 使用 SQLAlchemyJobStore |
| Celery Worker 挂起 | 任务超时无响应 | 设置 task_time_limit 和 task_soft_time_limit |
| Cron 表达式错误 | 分钟/小时/日/月/周顺序混淆 | 使用在线工具验证(crontab.guru) |
| 重复部署导致任务重复 | 多实例运行同一调度器 | 使用分布式锁或 APScheduler 的 coalesce |
schedule 的 time.sleep 太短 | CPU 空转 | 至少 sleep(1),推荐 10-60 秒 |
延伸阅读
版本差异(自动化办公库 → 当前稳定版)
| 库 | 本文编写时 | 当前稳定版 |
|---|---|---|
openpyxl(Excel) | 旧版 | 3.1.x |
python-docx(Word) | 旧版 | 1.1.x |
python-pptx(PPT) | 旧版 | 1.0.x |
reportlab(PDF) | 旧版 | 4.x |
PyPDF2/pypdf | PyPDF2 | 推荐 pypdf(4.x/5.x,PyPDF2 已停止维护) |
Pillow(图像) | 旧版 | 11.x |
本文讲解的自动化办公流程(读写 Excel/Word/PDF/PPT)与核心 API 在最新版本中成立;注意 PyPDF2 已迁移至 pypdf。