欢迎光临
我们一直在努力

定时任务与工作流调度:APScheduler/Prefect 完整配置

文章目录

    • 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 的局限就会逐一暴露:

需求crontab 能做到吗APSchedulerPrefect
任务失败自动重试
任务间有依赖关系 ❌(需自己实现) ✅(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 擅长跨步骤、有依赖、需协作的工作流编排。

维度APSchedulerPrefect
依赖管理 ❌ 需自行实现 ✅ 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 后端集成时非常有用。

如果本文对你有帮助,欢迎点赞、关注!你的支持是我持续输出高质量技术内容最大的动力。

赞(0)
未经允许不得转载:171主机测评 » 定时任务与工作流调度:APScheduler/Prefect 完整配置
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址