调度系统升级复盘:Crontab → Airflow → Prefect 的三步迭代
调度系统升级复盘:Crontab → Airflow → Prefect 的三步迭代
一、为什么我们要换调度系统?
今年初我们做了一次调度系统的全面升级。回顾这个过程,从最初写在服务器上的 Crontab,到后来部署 Apache Airflow,再到最近切换到 Prefect——每一步都不是"为了追新技术",而是被真实的业务痛点逼出来的。
Crontab 阶段的痛点很直接:几十个定时任务散落在4台服务器上,谁也不知道全貌。某天一台服务器磁盘满了,上面跑的5个ETL任务静默失败,下游报表数据全是空的,但没人收到任何通知。我们排查了2个小时才发现是磁盘满了导致任务没跑。
更头疼的是依赖管理。Crontab 只能设定时触发,但任务B依赖任务A的产出,A延迟了B怎么办?靠人工写touch /tmp/flag_file这种土办法,维护成本极高。
二、Crontab 阶段:简单但脆弱
Crontab 的优点只有一个:简单。写一行0 6 * * * /path/to/script.sh就完事了,不需要任何框架。
但缺点随着任务数量增长而爆发式暴露:
缺点一:无依赖管理
# crontab 配置示例 - 完全靠时间触发,没有依赖逻辑 # 每天凌晨6点跑用户数据清洗 0 6 * * * /opt/scripts/clean_user_data.sh # 每天凌晨7点跑用户画像计算(依赖上面的产出) 0 7 * * * /opt/scripts/calc_user_profile.sh # 每天凌晨8点跑用户画像导出到报表库(依赖上面的产出) 0 8 * * * /opt/scripts/export_user_profile.sh # 问题:如果第一个任务延迟或失败,后面的任务照跑不误 # 结果:画像计算拿到的数据是旧的或不完整的缺点二:无监控无告警
Crontab 任务失败了?没有人知道。除非你专门写日志检查脚本定期扫描,但这又是一个新的 Crontab 任务……
缺点三:无版本管理
改了一个任务的脚本,谁知道改了什么?没有代码仓库,没有变更记录,出了问题只能问"谁改了这个文件"。
缺点四:资源浪费
所有任务都在固定时间跑,不管数据是否到位、系统资源是否空闲。有些任务其实可以更早开始(数据已经准备好了),但 Crontab 不支持"数据就绪触发"。
这个阶段我们忍了8个月。忍不了的关键事件就是那次磁盘满导致5个任务静默失败,下游报表空了2小时。
三、Airflow 阶段:功能强大但运维沉重
痛定思痛后,我们决定迁移到 Airflow。Airflow 解决了 Crontab 的所有核心问题:依赖管理(DAG)、监控告警(UI + 回调)、版本管理(Git 集成)、数据就绪触发(Sensor)。
DAG 编写:依赖关系一目了然
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime, timedelta # 默认参数配置 default_args = { 'owner': 'data_team', 'depends_on_past': False, # 不依赖过去的运行 'retries': 3, # 失败重试3次 'retry_delay': timedelta(minutes=5), # 每次重试间隔5分钟 'email_on_failure': True, # 失败发邮件告警 'email': ['data-alert@company.com'] } # 定义DAG:用户数据处理流水线 with DAG( dag_id='user_data_pipeline', default_args=default_args, schedule_interval='0 6 * * *', # 每天凌晨6点调度 start_date=datetime(2025, 1, 1), catchup=False # 不补跑历史任务 ) as dag: # 任务1:清洗用户原始数据 clean_data = PythonOperator( task_id='clean_user_data', python_callable=clean_user_data_fn ) # 任务2:计算用户画像(依赖任务1) calc_profile = PythonOperator( task_id='calc_user_profile', python_callable=calc_user_profile_fn ) # 任务3:导出到报表库(依赖任务2) export_profile = PythonOperator( task_id='export_user_profile', python_callable=export_user_profile_fn ) # 定义依赖关系:clean → calc → export clean_data >> calc_profile >> export_profile依赖关系用>>语法表达,直观又清晰。任务1失败了,任务2和3自动等待或跳过,不会用脏数据继续跑。
Airflow 的运维痛点
但 Airflow 也有让人头疼的地方:
痛点一:部署复杂。Airflow 需要 PostgreSQL 做元数据库、Redis 做队列、Webserver + Scheduler + Worker 三组件协同。我们花了3周才把生产环境部署搞定,期间踩了无数坑。
痛点二:Scheduler 性能瓶颈。当DAG数量超过200个时,Scheduler 开始频繁超时,任务调度延迟从秒级变成了分钟级。我们不得不把 Scheduler 的解析间隔从30秒调到60秒,牺牲了实时性。
痛点三:DAG 解析开销。Airflow 每隔一段时间就要解析所有 DAG 文件,哪怕99%的DAG今天根本不会执行。这个开销随着DAG数量增长线性增加,很浪费。
痛点四:本地开发调试不便。在本地跑一个 Airflow DAG 需要启动整套组件,开发体验远不如直接跑 Python 脚本。
# Airflow 的Sensor等待数据就绪 - 功能强大但资源消耗高 from airflow.sensors.filesystem import FileSensor # 等待上游数据文件出现,最长等2小时 wait_for_data = FileSensor( task_id='wait_for_user_data', filepath='/data/raw/user_data_{{ ds }}.csv', poke_interval=60, # 每60秒检查一次 timeout=7200, # 最长等2小时 mode='poke' # poke模式会一直占用Worker slot )Sensor 的 poke 模式是个隐藏的坑:它会持续占用 Worker slot,几个长时间等待的 Sensor 就能堵住整个 Worker pool。改 mode='reschedule' 可以释放 slot,但调度延迟会变大。
四、Prefect 阶段:轻量、弹性、开发友好
今年5月我们开始试点 Prefect。选择 Prefect 的核心原因有三个:本地开发体验好、部署轻量、动态调度能力强。
本地开发:Python 函数就是任务
Prefect 的核心理念是"普通 Python 函数加上装饰器就是任务"。在本地开发时,你不需要启动任何服务,直接跑就行:
from prefect import flow, task, get_run_logger import pandas as pd @task(retries=3, retry_delay_seconds=300) def clean_user_data(date: str) -> pd.DataFrame: """清洗用户原始数据,失败自动重试3次""" logger = get_run_logger() raw_data = pd.read_csv(f'/data/raw/user_data_{date}.csv') # 去除重复记录 cleaned = raw_data.drop_duplicates(subset=['user_id']) # 填充缺失值 cleaned['age'] = cleaned['age'].fillna(cleaned['age'].median()) logger.info(f"清洗完成: {len(raw_data)} → {len(cleaned)} 条记录") return cleaned @task def calc_user_profile(data: pd.DataFrame) -> pd.DataFrame: """计算用户画像特征""" # 按用户分组,计算消费频次、客单价等画像指标 profile = data.groupby('user_id').agg({ 'order_amount': ['mean', 'sum', 'count'], 'category_id': 'nunique' }) profile.columns = ['avg_amount', 'total_amount', 'order_count', 'category_count'] return profile @task def export_user_profile(profile: pd.DataFrame, date: str) -> None: """导出用户画像到报表数据库""" logger = get_run_logger() profile.to_parquet(f'/data/output/user_profile_{date}.parquet') logger.info(f"导出完成: {len(profile)} 条画像记录") # 定义Flow:等价于Airflow的DAG @flow(name="user_data_pipeline") def user_data_pipeline(date: str): """用户数据处理流水线""" cleaned = clean_user_data(date) profile = calc_user_profile(cleaned) export_user_profile(profile, date) # 本地直接运行,无需启动任何服务 if __name__ == "__main__": user_data_pipeline("2026-07-01")对比 Airflow:不需要写 DAG 文件、不需要启动 Webserver/Scheduler/Worker、不需要连接 PostgreSQL。在本地就是普通 Python 代码,调试方便极了。
动态调度:运行时决定分支
Prefect 的动态调度能力是 Airflow 做不到的。Airflow 的 DAG 在执行前就固定了拓扑结构,而 Prefect 可以在运行时根据数据状态动态决定执行哪些分支:
from prefect import flow, task @task def check_data_freshness(date: str) -> bool: """检查数据新鲜度:是否在可接受范围内""" import os filepath = f'/data/raw/user_data_{date}.csv' if not os.path.exists(filepath): return False # 检查文件修改时间是否在今天 mtime = os.path.getmtime(filepath) return mtime > some_threshold @flow(name="adaptive_user_pipeline") def adaptive_user_pipeline(date: str): """自适应流水线:根据数据状态动态选择路径""" is_fresh = check_data_freshness(date) if is_fresh: # 数据新鲜,走快速路径 cleaned = clean_user_data(date) profile = calc_user_profile(cleaned) export_user_profile(profile, date) else: # 数据不新鲜,走补数据路径 backfill = backfill_user_data(date) # 从其他数据源补数据 cleaned = clean_user_data(date) profile = calc_user_profile(cleaned) export_user_profile(profile, date) notify_data_delay(date) # 通知数据延迟这种"运行时决定分支"的能力在 Airflow 里需要用 BranchOperator 实现,代码冗余且调试困难。
Prefect 的不足与应对
Prefect 也有缺点:
- 生态不如 Airflow 成熟:Airflow 有上百个 Operator 和 Provider,Prefect 的集成库还在增长中
- 社区规模较小:遇到问题找答案不如 Airflow 方便
- 企业版功能差异大:一些高级功能(如 RBAC、SSO)只在付费版里
我们的应对策略是:核心数据流水线用 Prefect,与外部系统集成的部分用 Prefect 调用 Airflow Operator 的底层 API,两套系统共存过渡。
五、总结
调度系统的三步迭代,每一步都是业务痛点驱动的:
- Crontab → Airflow:解决依赖管理、监控告警、版本控制。但部署重、调度器性能瓶颈
- Airflow → Prefect:解决开发体验、部署轻量、动态调度。但生态不够成熟
- 当前策略:核心流水线 Prefect + 集成场景 Airflow,两者共存过渡
几个关键复盘经验:
- 不要一步到位:从 Crontab 直接跳到 Prefect 会丢掉中间的经验积累,Airflow 阶段让我们学会了 DAG 设计和任务编排的核心方法论
- 调度器的选择要看团队规模:10人以下团队 Prefect 更合适,50人以上团队 Airflow 的成熟生态更有保障
- 本地开发体验是被低估的因素:Airflow 的本地调试痛苦间接导致了代码质量下降,Prefect 的"原生Python"体验让开发效率提升明显
- 过渡期并存是安全的做法:两套调度器同时运行,新任务用 Prefect,老任务保持 Airflow,逐步迁移不搞一刀切
调度系统的升级从来不是技术选型的问题,而是"现在这套系统到底解决不了什么问题"的答案在变。
