{T}

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.py

Python 管理 Cron

bash
pip install python-crontab
python
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_limittask_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/pypdfPyPDF2推荐 pypdf(4.x/5.x,PyPDF2 已停止维护)
Pillow(图像)旧版11.x

本文讲解的自动化办公流程(读写 Excel/Word/PDF/PPT)与核心 API 在最新版本中成立;注意 PyPDF2 已迁移至 pypdf。