调度系统升级复盘: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(minutes5), # 每次重试间隔5分钟 email_on_failure: True, # 失败发邮件告警 email: [data-alertcompany.com] } # 定义DAG用户数据处理流水线 with DAG( dag_iduser_data_pipeline, default_argsdefault_args, schedule_interval0 6 * * *, # 每天凌晨6点调度 start_datedatetime(2025, 1, 1), catchupFalse # 不补跑历史任务 ) as dag: # 任务1清洗用户原始数据 clean_data PythonOperator( task_idclean_user_data, python_callableclean_user_data_fn ) # 任务2计算用户画像依赖任务1 calc_profile PythonOperator( task_idcalc_user_profile, python_callablecalc_user_profile_fn ) # 任务3导出到报表库依赖任务2 export_profile PythonOperator( task_idexport_user_profile, python_callableexport_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_idwait_for_user_data, filepath/data/raw/user_data_{{ ds }}.csv, poke_interval60, # 每60秒检查一次 timeout7200, # 最长等2小时 modepoke # poke模式会一直占用Worker slot )Sensor 的 poke 模式是个隐藏的坑它会持续占用 Worker slot几个长时间等待的 Sensor 就能堵住整个 Worker pool。改 modereschedule 可以释放 slot但调度延迟会变大。四、Prefect 阶段轻量、弹性、开发友好今年5月我们开始试点 Prefect。选择 Prefect 的核心原因有三个本地开发体验好、部署轻量、动态调度能力强。本地开发Python 函数就是任务Prefect 的核心理念是普通 Python 函数加上装饰器就是任务。在本地开发时你不需要启动任何服务直接跑就行from prefect import flow, task, get_run_logger import pandas as pd task(retries3, retry_delay_seconds300) 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(nameuser_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(nameadaptive_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 和 ProviderPrefect 的集成库还在增长中社区规模较小遇到问题找答案不如 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逐步迁移不搞一刀切调度系统的升级从来不是技术选型的问题而是现在这套系统到底解决不了什么问题的答案在变。
