文章目录
-
- 1 定时任务的工程化需求
- 2 APScheduler 四组件模型
- 3 触发器实战
-
- 3.1 DateTrigger:一次性任务
- 3.2 IntervalTrigger:周期任务
- 3.3 CronTrigger:cron 表达式
- 4 执行器选择与混用
- 5 任务持久化:JobStore
- 6 任务生命周期监听与告警
- 7 Prefect 入门:从定时调度到工作流编排
- 8 Prefect 的任务编排与 DAG
- 9 APScheduler vs Prefect 选型决策
- 10 实战:每日数据采集 Pipeline(Prefect DAG)
- 11 总结
1 定时任务的工程化需求
一个定时任务从"能跑"到"生产可用"之间,隔着重重工程障碍。以每日凌晨爬取竞品数据的脚本为例——crontab -e 中写一行 0 2 * * * python crawl.py,谁都能在 5 分钟内搞定。但当场景变得复杂时,crontab 的局限就会逐一暴露:
| 任务失败自动重试 | ❌ | ✅ | ✅ |
| 任务间有依赖关系 | ❌ | ❌(需自己实现) | ✅(DAG 原生支持) |
| 任务执行失败后告警 | ❌ | ✅(需写 listener) | ✅(内置通知) |
| 服务重启后恢复未完成任务 | ❌ | ✅(JobStore 持久化) | ✅ |
| 同一个任务多个并发实例 | ❌ | ⚠️(需配置) | ✅(并发控制) |
| 任务执行历史查询 | ❌ | ❌ | ✅(UI) |
| 手动触发任意任务 | ❌ | ❌ | ✅ |
crontab 的模型是"无状态 shell 命令",而生产调度需要的是"有状态、有监控、可恢复的工作流"。这正是 APScheduler 和 Prefect 解决的问题。
2 APScheduler 四组件模型
APScheduler 的设计哲学是四组件解耦:
#mermaid-svg-YrAoyAKW2pLytv7m{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-YrAoyAKW2pLytv7m .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-YrAoyAKW2pLytv7m .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-YrAoyAKW2pLytv7m .error-icon{fill:#552222;}#mermaid-svg-YrAoyAKW2pLytv7m .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-YrAoyAKW2pLytv7m .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-YrAoyAKW2pLytv7m .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-YrAoyAKW2pLytv7m .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-YrAoyAKW2pLytv7m .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-YrAoyAKW2pLytv7m .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-YrAoyAKW2pLytv7m .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-YrAoyAKW2pLytv7m .marker{fill:#333333;stroke:#333333;}#mermaid-svg-YrAoyAKW2pLytv7m .marker.cross{stroke:#333333;}#mermaid-svg-YrAoyAKW2pLytv7m svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-YrAoyAKW2pLytv7m p{margin:0;}#mermaid-svg-YrAoyAKW2pLytv7m .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-YrAoyAKW2pLytv7m .cluster-label text{fill:#333;}#mermaid-svg-YrAoyAKW2pLytv7m .cluster-label span{color:#333;}#mermaid-svg-YrAoyAKW2pLytv7m .cluster-label span p{background-color:transparent;}#mermaid-svg-YrAoyAKW2pLytv7m .label text,#mermaid-svg-YrAoyAKW2pLytv7m span{fill:#333;color:#333;}#mermaid-svg-YrAoyAKW2pLytv7m .node rect,#mermaid-svg-YrAoyAKW2pLytv7m .node circle,#mermaid-svg-YrAoyAKW2pLytv7m .node ellipse,#mermaid-svg-YrAoyAKW2pLytv7m .node polygon,#mermaid-svg-YrAoyAKW2pLytv7m .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-YrAoyAKW2pLytv7m .rough-node .label text,#mermaid-svg-YrAoyAKW2pLytv7m .node .label text,#mermaid-svg-YrAoyAKW2pLytv7m .image-shape .label,#mermaid-svg-YrAoyAKW2pLytv7m .icon-shape .label{text-anchor:middle;}#mermaid-svg-YrAoyAKW2pLytv7m .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-YrAoyAKW2pLytv7m .rough-node .label,#mermaid-svg-YrAoyAKW2pLytv7m .node .label,#mermaid-svg-YrAoyAKW2pLytv7m .image-shape .label,#mermaid-svg-YrAoyAKW2pLytv7m .icon-shape .label{text-align:center;}#mermaid-svg-YrAoyAKW2pLytv7m .node.clickable{cursor:pointer;}#mermaid-svg-YrAoyAKW2pLytv7m .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-YrAoyAKW2pLytv7m .arrowheadPath{fill:#333333;}#mermaid-svg-YrAoyAKW2pLytv7m .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-YrAoyAKW2pLytv7m .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-YrAoyAKW2pLytv7m .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-YrAoyAKW2pLytv7m .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-YrAoyAKW2pLytv7m .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-YrAoyAKW2pLytv7m .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-YrAoyAKW2pLytv7m .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-YrAoyAKW2pLytv7m .cluster text{fill:#333;}#mermaid-svg-YrAoyAKW2pLytv7m .cluster span{color:#333;}#mermaid-svg-YrAoyAKW2pLytv7m div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-YrAoyAKW2pLytv7m .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-YrAoyAKW2pLytv7m rect.text{fill:none;stroke-width:0;}#mermaid-svg-YrAoyAKW2pLytv7m .icon-shape,#mermaid-svg-YrAoyAKW2pLytv7m .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-YrAoyAKW2pLytv7m .icon-shape p,#mermaid-svg-YrAoyAKW2pLytv7m .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-YrAoyAKW2pLytv7m .icon-shape .label rect,#mermaid-svg-YrAoyAKW2pLytv7m .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-YrAoyAKW2pLytv7m .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-YrAoyAKW2pLytv7m .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-YrAoyAKW2pLytv7m :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
Scheduler(调度器)
触发时机
任务持久化
执行任务
Executor(执行器)
AsyncioExecutor 异步
ThreadPoolExecutor 线程
ProcessPoolExecutor 进程
JobStore(状态存储)
MemoryJobStore 内存
SQLAlchemyJobStore 数据库
RedisJobStore Redis
Trigger(触发器)
DateTrigger 一次性
IntervalTrigger 周期
CronTrigger cron表达式
协调触发器、存储器、执行器的核心引擎
- Scheduler 是中央协调器,持有所有组件的引用,决定何时将哪个任务交给哪个执行器。
- Trigger 只负责计算"下一次触发时间",与任务本身解耦。同一个任务可以绑定不同的 Trigger。
- JobStore 管理任务的持久化——重启后任务不丢失。默认的 MemoryJobStore 在进程退出后清零。
- Executor 管理实际执行线程/进程。ThreadPoolExecutor 适合 I/O 密集型任务,ProcessPoolExecutor 适合 CPU 密集型。
3 触发器实战
3.1 DateTrigger:一次性任务
from datetime import datetime, timedelta
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.date import DateTrigger
def delayed_report():
print(f"报告生成于: {datetime.now()}")
scheduler = BackgroundScheduler()
scheduler.add_job(
delayed_report,
trigger=DateTrigger(run_date=datetime.now() + timedelta(minutes=5)),
id="one_time_report"
)
scheduler.start()
一次性任务在执行后自动从调度器移除。常见场景包括:延迟执行的导出任务、缓存预热、限时的资源抢占。
3.2 IntervalTrigger:周期任务
from apscheduler.triggers.interval import IntervalTrigger
scheduler.add_job(
sync_weather_data,
trigger=IntervalTrigger(hours=1, minutes=30), # 每1.5小时
id="weather_sync",
replace_existing=True # 重复 ID 的任务会被替换
)
IntervalTrigger 的 jitter 参数可以在固定间隔上叠加随机偏移,用于分散同一时刻大量定时任务的并发压力:
scheduler.add_job(
send_notification,
trigger=IntervalTrigger(hours=1, jitter=60), # 每次实际执行时间偏移 ±60 秒
)
3.3 CronTrigger:cron 表达式
CronTrigger 支持 6 字段格式(秒 分 时 日 月 星期),与 Linux crontab 的 5 字段格式有本质区别:
from apscheduler.triggers.cron import CronTrigger
# 每周一至周五 8:30 执行
scheduler.add_job(
morning_meeting_reminder,
trigger=CronTrigger(day_of_week="mon-fri", hour=8, minute=30),
)
# 每月最后一天 23:00 执行
scheduler.add_job(
monthly_close,
trigger=CronTrigger(day="last", hour=23, minute=0),
)
# 每年 1 月 1 日 00:00 执行
scheduler.add_job(
new_year_task,
trigger=CronTrigger(month=1, day=1, hour=0, minute=0),
)
# 每 15 分钟执行(与 IntervalTrigger 的效果等价)
scheduler.add_job(
health_check,
trigger=CronTrigger(minute="*/15"),
)
最容易踩的坑:Linux crontab 是 5 字段(月 日 时 分 周),而 APScheduler CronTrigger 默认是 6 字段(秒 分 时 日 月 周)。如果从 crontab 迁移而来而忽略秒字段,会发现任务总是提前或延迟 1 分钟触发。解决方案是明确使用 second=0 或 CronTrigger.from_crontab("*/15 * * * *") 从 5 字段兼容写法启动。
4 执行器选择与混用
三种执行器的选型逻辑非常清晰:
from apscheduler.executors.asyncio import AsyncIOExecutor
from apscheduler.executors.threadpool import ThreadPoolExecutor
from apscheduler.executors.processpool import ProcessPoolExecutor
# I/O 密集型:网络爬虫、API 调用、文件读写
# 推荐 ThreadPoolExecutor,线程池避免 GIL 阻塞 I/O 等待
executors = {
"default": ThreadPoolExecutor(20),
}
# CPU 密集型:数据加密、图像处理、机器学习推理
# ProcessPoolExecutor 将任务分配到子进程,绕过 GIL
executors = {
"cpu": ProcessPoolExecutor(4),
}
# 异步任务:已是 asyncio 生态
executors = {
"async": AsyncIOExecutor(),
}
配置 job_defaults 防止并发冲突:
scheduler = BackgroundScheduler(
executors=executors,
job_defaults={
"coalesce": True, # 堆积的多次执行合并为一次
"max_instances": 1, # 同一任务最多同时运行 1 个实例,防止重复
"misfire_grace_time": 60, # 任务错过触发时间后 60 秒内仍可补执行
},
)
5 任务持久化:JobStore
SQLAlchemyJobStore 是生产环境的标配——将任务元数据存储在数据库中,服务重启后自动恢复调度状态:
from apscheduler.jobstores.sqlalchemy import SQLAlchemyJobStore
from apscheduler.executors.pool import ThreadPoolExecutor, ProcessPoolExecutor
jobstores = {
"default": SQLAlchemyJobStore(url="sqlite:///scheduler.db"),
}
executors = {
"default": ThreadPoolExecutor(20),
"processpool": ProcessPoolExecutor(4),
}
job_defaults = {
"coalesce": True,
"max_instances": 3,
}
scheduler = BackgroundScheduler(
jobstores=jobstores,
executors=executors,
job_defaults=job_defaults,
)
注意:SQLAlchemyJobStore 存储的是"任务定义"(哪个函数、触发条件是什么),而非任务执行结果。任务输出仍然需要自行处理(写入数据库、发送消息等)。
6 任务生命周期监听与告警
add_listener() 接收事件代码并触发回调:
import httpx
def notify_on_event(event):
"""事件监听器:根据任务状态发送告警"""
job = scheduler.get_job(event.job_id)
job_name = job.name if job else "unknown"
if event.code == EVENT_JOB_EXECUTED:
print(f"[OK] {job_name} 执行成功")
elif event.code == EVENT_JOB_ERROR:
print(f"[ERROR] {job_name} 执行失败: {event.exception}")
# 飞书 Webhook 告警
httpx.post(
"https://open.feishu.cn/open-apis/bot/v2/hook/xxx",
json={
"msg_type": "text",
"content": {"text": f"⛔ 任务失败: {job_name}\\n异常: {event.exception}"}
}
)
elif event.code == EVENT_JOB_MISSED:
print(f"[WARN] {job_name} 错过了触发时间")
scheduler.add_listener(
notify_on_event,
EVENT_JOB_EXECUTED | EVENT_JOB_ERROR | EVENT_JOB_MISSED,
)
7 Prefect 入门:从定时调度到工作流编排
Prefect 的核心抽象是两个装饰器:
from prefect import flow, task
@task(retries=3, retry_delay_seconds=60, cache_key_fn=lambda ctx, fn, *args, **kwargs: str(args[0]))
def fetch_weather(city: str):
"""带自动重试和结果缓存的任务单元"""
response = httpx.get(f"https://api.weather.com/v3/wx/{city}")
return response.json()
@task
def save_to_database(data: dict):
"""写入数据库,失败不重试(幂等性无法保证)"""
db.insert("weather_log", data)
@flow(name="daily-weather-pipeline", log_prints=True)
def daily_pipeline():
"""工作流入口:组合多个任务"""
cities = ["beijing", "shanghai", "guangzhou"]
results = fetch_weather.map(cities) # 并行执行多个城市
save_to_database.map(results) # 顺序写入数据库
- @task 封装独立的计算单元,带自动重试(retries)、结果缓存(cache_key_fn)、超时控制(timeout_seconds)。
- @flow 是工作流的顶层容器,定义任务之间的执行关系。flow.run() 是同步入口,flow.serve() 是持久化服务。
- .map() 对列表中的每个元素并行执行同一个任务——类似 asyncio.gather() 但声明式、无需手动管理并发。
8 Prefect 的任务编排与 DAG
Prefect 的任务依赖通过数据流隐式定义——上游任务的返回值自动作为下游任务的输入:
@task
def extract() –> list[dict]:
return [{"id": 1}, {"id": 2}, {"id": 3}]
@task
def transform(record: dict) –> dict:
record["score"] = record["id"] * 10
return record
@task
def load(record: dict):
db.insert("processed", record)
@flow
def etl_flow():
records = extract() # 返回 list[dict]
transformed = transform.map(records) # 并行处理每个 record
load.map(transformed) # 顺序加载
如果任务间有明确的顺序要求(而非并行),使用 wait_for:
from prefect import wait_for
@flow
def sequential_flow():
task_a = task_a.submit()
task_b = task_b.submit(wait_for=[task_a])
task_c = task_c.submit(wait_for=[task_b])
9 APScheduler vs Prefect 选型决策
#mermaid-svg-a25x081V0m5CBtoe{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-a25x081V0m5CBtoe .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-a25x081V0m5CBtoe .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-a25x081V0m5CBtoe .error-icon{fill:#552222;}#mermaid-svg-a25x081V0m5CBtoe .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-a25x081V0m5CBtoe .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-a25x081V0m5CBtoe .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-a25x081V0m5CBtoe .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-a25x081V0m5CBtoe .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-a25x081V0m5CBtoe .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-a25x081V0m5CBtoe .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-a25x081V0m5CBtoe .marker{fill:#333333;stroke:#333333;}#mermaid-svg-a25x081V0m5CBtoe .marker.cross{stroke:#333333;}#mermaid-svg-a25x081V0m5CBtoe svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-a25x081V0m5CBtoe p{margin:0;}#mermaid-svg-a25x081V0m5CBtoe .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-a25x081V0m5CBtoe .cluster-label text{fill:#333;}#mermaid-svg-a25x081V0m5CBtoe .cluster-label span{color:#333;}#mermaid-svg-a25x081V0m5CBtoe .cluster-label span p{background-color:transparent;}#mermaid-svg-a25x081V0m5CBtoe .label text,#mermaid-svg-a25x081V0m5CBtoe span{fill:#333;color:#333;}#mermaid-svg-a25x081V0m5CBtoe .node rect,#mermaid-svg-a25x081V0m5CBtoe .node circle,#mermaid-svg-a25x081V0m5CBtoe .node ellipse,#mermaid-svg-a25x081V0m5CBtoe .node polygon,#mermaid-svg-a25x081V0m5CBtoe .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-a25x081V0m5CBtoe .rough-node .label text,#mermaid-svg-a25x081V0m5CBtoe .node .label text,#mermaid-svg-a25x081V0m5CBtoe .image-shape .label,#mermaid-svg-a25x081V0m5CBtoe .icon-shape .label{text-anchor:middle;}#mermaid-svg-a25x081V0m5CBtoe .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-a25x081V0m5CBtoe .rough-node .label,#mermaid-svg-a25x081V0m5CBtoe .node .label,#mermaid-svg-a25x081V0m5CBtoe .image-shape .label,#mermaid-svg-a25x081V0m5CBtoe .icon-shape .label{text-align:center;}#mermaid-svg-a25x081V0m5CBtoe .node.clickable{cursor:pointer;}#mermaid-svg-a25x081V0m5CBtoe .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-a25x081V0m5CBtoe .arrowheadPath{fill:#333333;}#mermaid-svg-a25x081V0m5CBtoe .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-a25x081V0m5CBtoe .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-a25x081V0m5CBtoe .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-a25x081V0m5CBtoe .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-a25x081V0m5CBtoe .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-a25x081V0m5CBtoe .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-a25x081V0m5CBtoe .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-a25x081V0m5CBtoe .cluster text{fill:#333;}#mermaid-svg-a25x081V0m5CBtoe .cluster span{color:#333;}#mermaid-svg-a25x081V0m5CBtoe div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-a25x081V0m5CBtoe .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-a25x081V0m5CBtoe rect.text{fill:none;stroke-width:0;}#mermaid-svg-a25x081V0m5CBtoe .icon-shape,#mermaid-svg-a25x081V0m5CBtoe .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-a25x081V0m5CBtoe .icon-shape p,#mermaid-svg-a25x081V0m5CBtoe .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-a25x081V0m5CBtoe .icon-shape .label rect,#mermaid-svg-a25x081V0m5CBtoe .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-a25x081V0m5CBtoe .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-a25x081V0m5CBtoe .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-a25x081V0m5CBtoe :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
否
是
否, 简单脚本
是
多步骤 DAG, 需要 UI
中等, 需持久化和告警
需重试/监控
需要定时执行任务?
任务间是否有依赖?
需要持久化?
Prefect ✅
crontab 或 APScheduler
复杂度有多高?
APScheduler + SQLAlchemyJobStore ✅
决策建议:
- crontab:临时脚本、一次性数据处理、不需要监控。
- APScheduler:单服务内定时任务、需要持久化(重启恢复)、需要基本告警。依赖 APScheduler 的"轻量"。
- Prefect:多步骤 DAG、团队协作(需要 UI 和版本管理)、需要手动触发任意任务、实验追踪。
10 实战:每日数据采集 Pipeline(Prefect DAG)
以下是一个完整的数据采集工作流,涵盖爬虫→清洗→入库→日报生成→推送四个阶段:
from prefect import flow, task
from datetime import datetime
import httpx
import json
import sqlite3
# ── 阶段 1:爬虫采集(每小时增量)───────────────────────────────
@task(retries=2, retry_delay_seconds=30)
def crawl_products() –> list[dict]:
"""从电商平台 API 抓取商品列表"""
response = httpx.get(
"https://api.example.com/products",
headers={"Authorization": "Bearer xxx"},
timeout=30.0,
)
response.raise_for_status()
return response.json()["products"]
# ── 阶段 2:数据清洗 ──────────────────────────────────────────
@task
def clean_products(products: list[dict]) –> list[dict]:
"""过滤无效数据,标准化字段"""
cleaned = []
for p in products:
if not p.get("price") or p["price"] <= 0:
continue
cleaned.append({
"product_id": p["id"],
"name": p["name"].strip(),
"price": float(p["price"]),
"category": p.get("category", "unknown"),
"crawled_at": datetime.now().isoformat(),
})
return cleaned
# ── 阶段 3:写入数据库 ────────────────────────────────────────
@task
def save_to_db(products: list[dict]):
"""写入 SQLite,每日覆盖"""
conn = sqlite3.connect("daily_products.db")
conn.execute("CREATE TABLE IF NOT EXISTS products (…)")
for p in products:
conn.execute(
"INSERT OR REPLACE INTO products VALUES (?, ?, ?, ?, ?)",
(p["product_id"], p["name"], p["price"], p["category"], p["crawled_at"])
)
conn.commit()
conn.close()
return len(products)
# ── 阶段 4:生成日报并推送 ────────────────────────────────────
@task
def generate_report(count: int):
"""生成本日采集报告并推送到钉钉"""
report = f"✅ 数据采集完成\\n今日新增: {count} 条\\n时间: {datetime.now()}"
httpx.post(
"https://oapi.dingtalk.com/robot/send?access_token=xxx",
json={"msgtype": "text", "text": {"content": report}},
)
# ── 主工作流 ─────────────────────────────────────────────────
@flow(
name="daily-product-pipeline",
schedule=CronTrigger(day_of_week="mon-fri", hour=2, minute=0),
log_prints=True,
)
def daily_pipeline():
"""每日凌晨 2 点执行"""
products = crawl_products()
cleaned = clean_products(products)
count = save_to_db(cleaned)
generate_report(count)
# ── 部署 ─────────────────────────────────────────────────────
if __name__ == "__main__":
# 本地测试:直接运行
daily_pipeline()
# 生产部署:
# $ prefect deploy daily_pipeline.py:daily_pipeline
# $ prefect worker start –pool "default-pool"
#mermaid-svg-MFhHTJfzzq7DjRzp{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-MFhHTJfzzq7DjRzp .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-MFhHTJfzzq7DjRzp .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-MFhHTJfzzq7DjRzp .error-icon{fill:#552222;}#mermaid-svg-MFhHTJfzzq7DjRzp .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-MFhHTJfzzq7DjRzp .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-MFhHTJfzzq7DjRzp .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-MFhHTJfzzq7DjRzp .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-MFhHTJfzzq7DjRzp .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-MFhHTJfzzq7DjRzp .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-MFhHTJfzzq7DjRzp .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-MFhHTJfzzq7DjRzp .marker{fill:#333333;stroke:#333333;}#mermaid-svg-MFhHTJfzzq7DjRzp .marker.cross{stroke:#333333;}#mermaid-svg-MFhHTJfzzq7DjRzp svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-MFhHTJfzzq7DjRzp p{margin:0;}#mermaid-svg-MFhHTJfzzq7DjRzp .mermaid-main-font{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}#mermaid-svg-MFhHTJfzzq7DjRzp .exclude-range{fill:#eeeeee;}#mermaid-svg-MFhHTJfzzq7DjRzp .section{stroke:none;opacity:0.2;}#mermaid-svg-MFhHTJfzzq7DjRzp .section0{fill:rgba(102, 102, 255, 0.49);}#mermaid-svg-MFhHTJfzzq7DjRzp .section2{fill:#fff400;}#mermaid-svg-MFhHTJfzzq7DjRzp .section1,#mermaid-svg-MFhHTJfzzq7DjRzp .section3{fill:white;opacity:0.2;}#mermaid-svg-MFhHTJfzzq7DjRzp .sectionTitle0{fill:#333;}#mermaid-svg-MFhHTJfzzq7DjRzp .sectionTitle1{fill:#333;}#mermaid-svg-MFhHTJfzzq7DjRzp .sectionTitle2{fill:#333;}#mermaid-svg-MFhHTJfzzq7DjRzp .sectionTitle3{fill:#333;}#mermaid-svg-MFhHTJfzzq7DjRzp .sectionTitle{text-anchor:start;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}#mermaid-svg-MFhHTJfzzq7DjRzp .grid .tick{stroke:lightgrey;opacity:0.8;shape-rendering:crispEdges;}#mermaid-svg-MFhHTJfzzq7DjRzp .grid .tick text{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;fill:#333;}#mermaid-svg-MFhHTJfzzq7DjRzp .grid path{stroke-width:0;}#mermaid-svg-MFhHTJfzzq7DjRzp .today{fill:none;stroke:red;stroke-width:2px;}#mermaid-svg-MFhHTJfzzq7DjRzp .task{stroke-width:2;}#mermaid-svg-MFhHTJfzzq7DjRzp .taskText{text-anchor:middle;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}#mermaid-svg-MFhHTJfzzq7DjRzp .taskTextOutsideRight{fill:black;text-anchor:start;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}#mermaid-svg-MFhHTJfzzq7DjRzp .taskTextOutsideLeft{fill:black;text-anchor:end;}#mermaid-svg-MFhHTJfzzq7DjRzp .task.clickable{cursor:pointer;}#mermaid-svg-MFhHTJfzzq7DjRzp .taskText.clickable{cursor:pointer;fill:#003163!important;font-weight:bold;}#mermaid-svg-MFhHTJfzzq7DjRzp .taskTextOutsideLeft.clickable{cursor:pointer;fill:#003163!important;font-weight:bold;}#mermaid-svg-MFhHTJfzzq7DjRzp .taskTextOutsideRight.clickable{cursor:pointer;fill:#003163!important;font-weight:bold;}#mermaid-svg-MFhHTJfzzq7DjRzp .taskText0,#mermaid-svg-MFhHTJfzzq7DjRzp .taskText1,#mermaid-svg-MFhHTJfzzq7DjRzp .taskText2,#mermaid-svg-MFhHTJfzzq7DjRzp .taskText3{fill:white;}#mermaid-svg-MFhHTJfzzq7DjRzp .task0,#mermaid-svg-MFhHTJfzzq7DjRzp .task1,#mermaid-svg-MFhHTJfzzq7DjRzp .task2,#mermaid-svg-MFhHTJfzzq7DjRzp .task3{fill:#8a90dd;stroke:#534fbc;}#mermaid-svg-MFhHTJfzzq7DjRzp .taskTextOutside0,#mermaid-svg-MFhHTJfzzq7DjRzp .taskTextOutside2{fill:black;}#mermaid-svg-MFhHTJfzzq7DjRzp .taskTextOutside1,#mermaid-svg-MFhHTJfzzq7DjRzp .taskTextOutside3{fill:black;}#mermaid-svg-MFhHTJfzzq7DjRzp .active0,#mermaid-svg-MFhHTJfzzq7DjRzp .active1,#mermaid-svg-MFhHTJfzzq7DjRzp .active2,#mermaid-svg-MFhHTJfzzq7DjRzp .active3{fill:#bfc7ff;stroke:#534fbc;}#mermaid-svg-MFhHTJfzzq7DjRzp .activeText0,#mermaid-svg-MFhHTJfzzq7DjRzp .activeText1,#mermaid-svg-MFhHTJfzzq7DjRzp .activeText2,#mermaid-svg-MFhHTJfzzq7DjRzp .activeText3{fill:black!important;}#mermaid-svg-MFhHTJfzzq7DjRzp .done0,#mermaid-svg-MFhHTJfzzq7DjRzp .done1,#mermaid-svg-MFhHTJfzzq7DjRzp .done2,#mermaid-svg-MFhHTJfzzq7DjRzp .done3{stroke:grey;fill:lightgrey;stroke-width:2;}#mermaid-svg-MFhHTJfzzq7DjRzp .doneText0,#mermaid-svg-MFhHTJfzzq7DjRzp .doneText1,#mermaid-svg-MFhHTJfzzq7DjRzp .doneText2,#mermaid-svg-MFhHTJfzzq7DjRzp .doneText3{fill:black!important;}#mermaid-svg-MFhHTJfzzq7DjRzp .doneText0.taskTextOutsideLeft,#mermaid-svg-MFhHTJfzzq7DjRzp .doneText0.taskTextOutsideRight,#mermaid-svg-MFhHTJfzzq7DjRzp .doneText1.taskTextOutsideLeft,#mermaid-svg-MFhHTJfzzq7DjRzp .doneText1.taskTextOutsideRight,#mermaid-svg-MFhHTJfzzq7DjRzp .doneText2.taskTextOutsideLeft,#mermaid-svg-MFhHTJfzzq7DjRzp .doneText2.taskTextOutsideRight,#mermaid-svg-MFhHTJfzzq7DjRzp .doneText3.taskTextOutsideLeft,#mermaid-svg-MFhHTJfzzq7DjRzp .doneText3.taskTextOutsideRight{fill:black!important;}#mermaid-svg-MFhHTJfzzq7DjRzp .crit0,#mermaid-svg-MFhHTJfzzq7DjRzp .crit1,#mermaid-svg-MFhHTJfzzq7DjRzp .crit2,#mermaid-svg-MFhHTJfzzq7DjRzp .crit3{stroke:#ff8888;fill:red;stroke-width:2;}#mermaid-svg-MFhHTJfzzq7DjRzp .activeCrit0,#mermaid-svg-MFhHTJfzzq7DjRzp .activeCrit1,#mermaid-svg-MFhHTJfzzq7DjRzp .activeCrit2,#mermaid-svg-MFhHTJfzzq7DjRzp .activeCrit3{stroke:#ff8888;fill:#bfc7ff;stroke-width:2;}#mermaid-svg-MFhHTJfzzq7DjRzp .doneCrit0,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCrit1,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCrit2,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCrit3{stroke:#ff8888;fill:lightgrey;stroke-width:2;cursor:pointer;shape-rendering:crispEdges;}#mermaid-svg-MFhHTJfzzq7DjRzp .milestone{transform:rotate(45deg) scale(0.8,0.8);}#mermaid-svg-MFhHTJfzzq7DjRzp .milestoneText{font-style:italic;}#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText0,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText1,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText2,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText3{fill:black!important;}#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText0.taskTextOutsideLeft,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText0.taskTextOutsideRight,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText1.taskTextOutsideLeft,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText1.taskTextOutsideRight,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText2.taskTextOutsideLeft,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText2.taskTextOutsideRight,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText3.taskTextOutsideLeft,#mermaid-svg-MFhHTJfzzq7DjRzp .doneCritText3.taskTextOutsideRight{fill:black!important;}#mermaid-svg-MFhHTJfzzq7DjRzp .vert{stroke:navy;}#mermaid-svg-MFhHTJfzzq7DjRzp .vertText{font-size:15px;text-anchor:middle;fill:navy!important;}#mermaid-svg-MFhHTJfzzq7DjRzp .activeCritText0,#mermaid-svg-MFhHTJfzzq7DjRzp .activeCritText1,#mermaid-svg-MFhHTJfzzq7DjRzp .activeCritText2,#mermaid-svg-MFhHTJfzzq7DjRzp .activeCritText3{fill:black!important;}#mermaid-svg-MFhHTJfzzq7DjRzp .titleText{text-anchor:middle;font-size:18px;fill:#333;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}#mermaid-svg-MFhHTJfzzq7DjRzp :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
08:00
08:00
08:00
08:00
08:00
08:00
08:00
08:00
08:00
crawl_products (并发API请求)
clean_products (数据标准化)
save_to_db (写入SQLite)
generate_report (钉钉通知)
爬虫
清洗
入库
推送
每日数据采集 Pipeline 时间线
11 总结
APScheduler 和 Prefect 解决的是两类问题:APScheduler 擅长单服务内的轻量定时调度,Prefect 擅长跨步骤、有依赖、需协作的工作流编排。
| 依赖管理 | ❌ 需自行实现 | ✅ DAG 原生支持 |
| 持久化 | ✅ SQLAlchemyJobStore | ✅ 内置 |
| UI 界面 | ❌ | ✅ Prefect Cloud/Server |
| 部署复杂度 | 低(单进程) | 中(Server + Worker) |
| 适合场景 | 简单定时任务 | 复杂工作流 |
两者并非互斥——在 Prefect flow 中,可以通过 PythonShellTask 调用 APScheduler 管理的外部脚本,实现渐进式迁移。
扩展阅读
- APScheduler 官方文档:https://apscheduler.readthedocs.io
- Prefect 文档:https://docs.prefect.io
- 定时任务的对称面是手动触发——Prefect 支持通过 API 和 UI 手动运行任意 flow,这在与 Web 后端集成时非常有用。
如果本文对你有帮助,欢迎点赞、关注!你的支持是我持续输出高质量技术内容最大的动力。


