Prefect 3 实战:告别 DAG 地狱,用纯 Python 编排生产级数据流水线
实战, 告别那种被称为 DAG 地狱的状况, 通过纯粹的方式去编排处于生产级别的数据流水线。以下这些词, 3, 工作流编排, 替代, 数据流水线, AI Agent编排, 适用于这样部分读者, 数据工程师、后端 以及全栈开发者, 还有正在评估编排工具的技术负责人。若是你曾经编写过数据流水线, 很大概率有过这般痛苦的经历: 为要运行一个ETL过程, 得先去学习一套DAG定义语法, 接着进行配置、撰写XCom来传递数据, 最终还得管理一整套相关的东西。代码没写出多少行, 运维的内心已然疲惫不堪了。思路很是直接的3: 你的函数, 那便是工作流。添加两个装饰器flow以及task, 它便自动获取到了状态追踪、重试、缓存、调度、可视化DAG还有生产级可观测性——然而你所编写的依旧是原生的结合我近期落地经验产生的这篇文章, 引领你从概念起始到进行部署, 将3切实运用起来, 并且探讨它与怎么进行选择、以及怎样给AI Agent充当“编排大脑”相关的情况。一、核心概念五个词讲清楚概念是什么类比Flow一个工作流的入口用 flow 装饰的普通函数主程序 main()Task工作流里的一个可追踪步骤用 task 装饰一个带重试/缓存的函数把 flow 固化成可调度、可远程触发的部署单元发布后的服务Work Pool /描述在哪跑的基础设施 真正拉取任务执行的进程队列 消费者Block /安全的配置 / 密钥存储UI 可编辑、类型校验环境变量 / 管理3 最关键的两个变化值得单独拎出来取代了2的Agent的Work Pool, 其声明基础设施类型 / /, 然后去轮询该池并拉起运行, 哪个与哪个环境绑定状态清晰显然明白, 不必再有“发生这个run却跑到错误的agent上”这类玄之又玄且难以捉摸的调试状况。原本独家为Cloud所享有的事件引擎, 如今进入了开源包, 现在可以使用“文件落到S3便即时促发事件发生”,“一小时之内出现三次结果失败的情况才进行示警通知”这类具备的规则。无需再写针对相关操作的轮询的脚本。在性能方面, 3 对客户端引擎进行了重新书写, 相较于 2, 运行时所产生的开销, 最高能够降低 90%至 98%, 在官方所处的分布式并行场景当中, 这种情况表现得格外显著。二、30 秒上手第一个 Flowpip install prefect prefect server start # 自带 SQLite一条命令起本地服务UI 在 http://localhost:4200from prefect import flow, task task def say_hello(name: str) - str: msg fHello, {name}! print(msg) return msg flow(namegreeting-flow) def greeting_flow(names: list[str]) - None: for name in names: say_hello(name) # 像普通函数一样调用但每一次都有状态/日志 if __name__ __main__: greeting_flow([Alice, Bob, Carol])通过直接运行flow.py之后, 打开UI, 进而就可以看见自动生成的依赖图, 及其每次调用时所产生的耗时、日志以及状态。需要注意的是: 未撰写任何DAG文件, 未配置YAML且未启动调度器, 仅仅如此便完成了编排。三、真实例子每日报表同步流水线下面有个例子, 它贴近生产, 而且集成了重试, 还含输入级缓存, 有参数化调度, 也有日志透出:from datetime import timedelta from prefect import flow, task from prefect.tasks import task_input_hash task(retries3, retry_delay_seconds5, cache_policytask_input_hash) def fetch_report(report_date: str) - dict: 拉取报表cache_policy 基于输入哈希 同一 report_date 一小时内重跑直接命中缓存不再真正请求。 return {date: report_date, rows: 1000} task(retries2) def transform(data: dict) - dict: data[clean] True return data task def load_to_warehouse(data: dict) - None: print(floaded {data[rows]} rows into warehouse) flow(namedaily-report-sync, log_printsTrue) def daily_report_sync(report_date: str 2026-08-26) - None: raw fetch_report(report_date) clean transform(raw) load_to_warehouse(clean) if __name__ __main__: daily_report_sync()几个关键设计点四、动态工作流 真正的杀手锏所涉及的DAG属于运行前所编译好的静态图, 而动态分支依据数据来决定运行哪些任务在书写方面极为别扭。相关的图乃是运行时凭借控制流自然生成的其中if/else、while以及动态fan-out均为原生语法。from prefect import flow, task task def process_file(path: str) - int: return len(path) flow def batch_process(files: list[str]) - list[int]: # 运行时才知道要处理哪些文件自动 fan-out 并行 return process_file.map(files)import httpx from prefect import flow, task task async def fetch_url(url: str) - str: async with httpx.AsyncClient() as client: return (await client.get(url)).text flow async def crawl(urls: list[str]) - list[str]: return await fetch_url.map(urls) # 原生 async直接 map这种范式呈现出“代码即工作流”的特点, 它对于那种依据运行时数据来决定走向的 AI / LLM 来说是特别友善的, 比如说, 会先使得模型去判断意图, 然后再据此决定调用哪一个下游 task。五、部署与执行模型数据不出你自己的基础设施采用混合执行方式时, 那你的flow代码一直会在自身的环境当中运行, 仅仅是去做调度以及具备可观测性的“控制管理层面”的工作。哪怕是处于Cloud这种环境下, 同样不存在从云端到你内部网络的进入连接情况, 密钥以及数据都不会脱离你的基础设施范围——这对于合规要求以及数据驻留而言, 是天然就呈现友好态势的。部署三步# 1. 构建 deployment指定入口、名称、目标 work pool prefect deployment build my_flow.py:my_flow --name prod --pool default-process # 2. 应用写入 server prefect deployment apply my_flow-deployment.yaml # 3. 起一个 worker 拉取该池的任务 prefect worker start --pool default-process调度所支持的有Cron, 还有/, 以及RRule这三种。下面呈现的这张图, 将从代码起始一直到运行的完整链路展示了出来:六、容错与幂等失败能回滚task级重试之外, 3提供了事务接口, 会将一组task包成事务, 失败时自动的那种将副作用进行回滚, 这对于“扣款 发邮件”“写库 调外部API”这类需要一致性的步骤而言是非常关键的。from prefect.transactions import transaction def create_user(name: str): with transaction() as txn: user_id db.insert(name) # 注册回滚钩子事务失败自动清理 txn.add_rollback(lambda: db.delete(user_id)) send_welcome_email(user_id)予以配合, 你的流水线既有能力“越过未发生改变的计算”, 同时还能够“在计算错误之时自行补救”。七、事件驱动让 自己响应世界事件, 也就是 Event, 属于 3 的一等公民。能被触发的有任何状态变化, 还有自定义事件以及外部事件, 并且无需去写轮询任务。典型场景是, S3桶出现新文件, 之后自动启动数据处理。时序是这样的如下:你同样能够依据“某成功”这一条件, 以及“一小时内失败N次”这种情况, 去触发通知, 暂停调度, 或者启动下游的flow。八、 vs 怎么选维度3首次发布20152018工作流定义静态 DAG 动态 flow / task纯 调度Cron 为主 事件Cron / 原生事件驱动执行模型从 Work Pool 拉取自托管复杂度高//DB//低单 代码位置计算在集群代码在你自己基础设施控制面分离学习曲线中等简单生态极大1000 / 成长中任务即 可直接用任意库适用场景成熟批处理 ETL、大团队、强治理、动态 / 事件驱动、快速迭代一句话决策一个是 2.0 开源, 另一个也是 2.0 开源, 不存在授权方面的成本问题, 差距在于范式以及运维应去承担的负担, 并非。九、Bonus给 AI Agent 当编排大脑这是我最为激动满怀的部分, 3 提供了 MCP, 它能够赋予编排环境直接接入 Code、Codex、CLI 这类 AI 助手的能力这些 AI 助手能够以只读的方式, 去查看 flow run、task run 以及日志, 还可以进行内置文档的检索。prefect mcp start # 启动 MCP server针对多步的Agent任务, 针对长链路的Agent任务, 针对需要重试与人工确认的Agent任务, 将“步骤编排”予以交付, 将“智能决策”交付给LLM, 这是一种相较于纯ReAct循环更为稳健的工程结构, 每一步都具备状态, 每一步都能够恢复, 每一步都具有可观测性, 而非在黑盒中盲目行事。十、每个 task 仅做一件事, 这是十条最佳实践之一, 因为粒度细才利于复用与缓存。对于 task, 要始终添加特定内容, 且瞬时故障不应依赖人工盯守。使用显式方式而非隐式缓存, 要明确界定“什么叫相同输入”。将 True 设置为让调试日志进入 UI。对 flow 进行参数化, 调度期间通过传参而非修改代码来实现。采用 Block 存储密钥, 切勿把 token 固定写在脚本里。先在本地让 start 成功运行, 之后再进行 build 并部署到生产环境。针对一个 Work Pool 此情况, 仅绑定一种基础设施, 要是出现混用状况, 那就拆分成多个池。关于事件驱动替代轮询方面, 能够触发的情况, 就不要定时进行空跑检查。对于长链路添加方面, 失败之时要有回滚操作, 以此保证幂等性。总结。3 的价值并非在于“又出现一个编排工具”, 而是在于它将编排的摩擦降至最低程度: 你所撰写的就是如此它所负责的乃是生产级的那一部分诸如状态、重试、缓存、调度、可视化这些方面。对于动态的、事件驱动的以及日益常见的 AI / LLM 工作流而言, 它的代码优先范式相较于静态 DAG 要顺手许多。若你的团队正因那运维成本而被劝退, 又或者你的那个怎样怎样的东西越来越似“凭借数据现编出来的图”, 3这事是值得尝试一下的——pip能弄起一个什么什么, 花十分钟就能让头一条生产级流水线运行通畅。参考的资料包含: 有着名为docs..io的官方文档、以3.0作为其发布说明的内容以及关于vs 2026的对比评测。