如果你写过稍微复杂一点的Python数据处理脚本,大概率遇到过这样的场景:任务A跑完才能跑任务B,任务B失败后要手动重跑,每天凌晨还要靠crontab定时触发。脚本一多,依赖关系就像一团乱麻,出了问题排查起来非常痛苦。Prefect就是为解决这类问题而生的现代化工作流编排框架,它用纯Python的方式定义工作流,配合自带的Prefect服务器,可以轻松实现任务调度、依赖管理、失败重试和运行监控。本文带你从零认识Prefect,并动手搭建一套可用的数据管道。

Prefect是什么?它和其他编排工具有什么区别
Prefect是一个开源的Python工作流编排框架,它的核心理念是让你的代码尽量保持普通Python代码的样子,只需要加上装饰器,就能获得完整的编排能力。相比之下,Airflow需要学习DAG的概念和一套特定的编写规范,代码风格受框架约束较大;而Prefect几乎零侵入,老的脚本加上一个装饰器就能变成可管理的工作流,迁移成本非常低。
Prefect 2.x之后的版本在架构上做了很大简化。它引入了动态工作流的概念,工作流的依赖关系不需要提前静态声明,而是在运行时动态确定,这让分支、循环、映射等逻辑写起来和写普通Python没太大区别。同时Prefect原生支持失败重试、结果缓存、并发限制等生产环境常用特性,官方还提供Prefect Cloud云服务,但对于数据敏感的团队,自建的Prefect服务器往往是更合适的选择。
核心概念:Flow与Task到底怎么理解
用Prefect之前必须搞清楚两个最基本的单位:Flow和Task。Task是工作流中最小的执行单元,比如下载一个文件、清洗一批数据、写入一张数据库表,都可以定义为一个Task。在代码里,Task就是一个加了@task装饰器的函数。Flow则是把多个Task组织起来的容器,用@flow装饰器标记,它定义了任务的执行顺序和依赖关系。
举个直观的例子:一个每日ETL管道,可以拆成三个Task,分别是拉取API数据、清洗转换数据、写入数据库,然后由一个Flow来统一调度这三个Task。Flow内部就是普通的函数调用,你写clean_data(raw),Prefect就自动识别出写入任务依赖于清洗任务的结果,从而保证执行顺序。这种写法的好处是逻辑一目了然,几乎不需要额外学习编排语法。
除了Flow和Task,还有几个概念值得了解。Deployment(部署)负责把工作流注册到服务器上,让工作流可以被调度和远程触发;Work Pool(工作池)决定了工作流实际在哪个环境中执行;Work Queue则负责任务的分发排队。初次接触可以先忽略后面几个,把Flow和Task玩熟了再深入。
安装与本地服务器搭建步骤
Prefect的安装非常简单,建议先用虚拟环境隔离依赖。执行pip install prefect即可完成安装,安装完成后通过prefect version确认版本。Prefect对Python版本有一定要求,建议使用Python 3.9以上的环境。
接下来启动本地Prefect服务器,只需要一条命令:prefect server start。首次启动时Prefect会在本地初始化一个SQLite数据库存储元数据,然后启动API服务和Web UI。默认情况下,Web UI的访问地址是http://127.0.0.1:4200,打开浏览器就能看到控制面板。如果想让局域网内其他机器访问,可以加--host 0.0.0.0参数。
服务器启动后,还需要让你的命令行环境指向这个服务器地址,执行prefect config set PREFECT_API_URL=http://127.0.0.1:4200/api。这样之后所有工作流的运行记录都会上报到本地服务器,而不是使用临时的本地存储。完成这两步,一套完整的编排环境就算搭好了。
编写第一个工作流:从单任务到任务编排
下面写一个最简单的工作流感受一下Prefect的用法:
from prefect import flow, task
@task(retries=2, retry_delay_seconds=5)
def fetch_data():
print("模拟拉取数据")
return [1, 2, 3]
@task
def process_data(data):
return [x * 10 for x in data]
@flow(name="我的第一个数据管道")
def my_pipeline():
raw = fetch_data()
result = process_data(raw)
print(result)
if __name__ == "__main__":
my_pipeline()
直接运行这个脚本,任务就会执行,同时打开Web UI会发现这次运行已经被完整记录下来,包括每个Task的耗时、日志和状态。注意fetch_data任务设置了retries参数,表示失败后自动重试两次,每次间隔5秒,这在处理网络请求这类不稳定操作时非常实用。
如果希望任务并行执行,可以把任务调用改成.submit()方式并配合.result()获取结果,例如future = fetch_data.submit()再future.result()。Prefect会自动管理并发,配合并发限制参数还能控制资源占用,避免下游数据库被并发任务压垮。
定时调度:让管道每天自动运行
本地手动跑通之后,下一步是让工作流定时执行。在Prefect中这通过Deployment实现。为Flow添加调度配置,可以在代码中使用.serve()方法,一行代码即可完成部署加调度:
if __name__ == "__main__":
my_pipeline.serve(name="每日ETL", cron="0 2 * * *")
这里的cron表达式0 2 * * *表示每天凌晨两点执行一次,语法与Linux crontab一致。运行脚本后,Prefect会自动创建Deployment并注册到服务器上,即使你关闭了这个脚本,调度依然由服务器负责,到了时间就会把任务派发给工作池执行。前提是你要用prefect worker start --pool default-agent-pool启动一个worker进程来实际执行任务。
除了cron,Prefect还支持interval间隔调度和自定义的RRule规则,也支持传入时区参数,避免夏令时之类的时区陷阱。对于数据处理场景,还可以开启防重叠机制,确保上一次还没跑完时不会重复触发新的运行。
监控、日志与失败处理实践
编排工具的价值很大程度体现在可观测性上。Prefect的Web UI提供了直观的运行视图,Flow Run页面能看到每次运行的状态、时间轴和任务层级,点击任意Task还能查看完整日志。Prefect自动捕获了print输出和logging记录,所以老脚本几乎不需要改造就有完整的日志体系。
失败处理方面,建议在实践中遵循几个原则:对外部依赖(如API、数据库)的调用一定要配置重试参数;对关键失败要主动抛出异常让Flow标记为失败状态,而不是静默吞掉;对可以容忍失败的任务,可以设置continue_on_failure让管道继续往下走。Prefect还支持配置自动化规则(Automation),比如在Flow失败时发送webhook通知到企业微信或钉钉,让团队第一时间感知问题。
另一个实用技巧是使用缓存策略。对于代价高昂的计算步骤,可以给Task加上cache_key_fn和cache_expiration,让相同输入在有效期内直接复用上次结果,重复运行管道时能节省大量时间。
总结与进阶方向
整体来看,Prefect用极低的侵入性补齐了Python数据管道在调度、依赖、重试和监控上的短板。搭建路径也很清晰:安装Prefect、启动本地服务器、编写Flow和Task、创建Deployment配置调度、启动worker执行任务,五步就能跑起一套生产可用的编排系统。
当你熟悉基础用法后,可以继续探索几个方向:使用Docker Work Pool把任务隔离到容器中运行,保证环境一致性;用Secret Block管理数据库密码和API密钥;通过子流程(Subflow)拆分复杂管道,让逻辑更清晰;以及在多台机器上部署worker实现分布式执行。掌握这些能力之后,Prefect足以支撑从个人脚本到企业级数据平台的各种编排需求。