ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

如何把 Prefect 从 2.x 升级到 3.0 并处理 Pydantic 2 等破坏性变更

如何把 Prefect 从 2.x 升级到 3.0 并处理 Pydantic 2 等破坏性变更 如何把 Prefect 从 2.x 升级到 3.0 并处理 Pydantic 2 等破坏性变更【免费下载链接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.项目地址: https://gitcode.com/GitHub_Trending/pr/prefect如果你的 Python 项目还在使用prefect2.x 版本的包这篇文档对应的官方升级指南升级指南给出了一条明确的迁移路径升级prefect包、自托管时升级数据库、再逐个处理 3.0 的破坏性变更Pydantic 2、流程最终状态语义、任务缓存、unmapped()可变对象等。适用前提是 Python 3.10 或更新版本Prefect 作为 Python 包发布要求 Python 3.10大部分用户只需少量甚至零代码改动。官方说明Prefect 2.0指prefect包的 2.x 系列Prefect 3.0仅指 3.x 系列两者都与 Prefect Cloud 商业产品本身无严格绑定。升级前建议先阅读 Whats new in Prefect 3.0 了解新增能力events automations、事务式缓存层、任务嵌套执行等。升级前的准备暂停所有 deployment 的调度。官方明确建议由于 Pydantic 对 datetime 的处理差异会影响调度器的幂等逻辑存在调度器在升级后第一轮循环中重复触发 run 的小概率风险因此升级前应暂停调度。确认你使用哪些集成extra。如果你安装过prefect[aws]这类集成包升级prefect时要一并升级对应集成否则可能出现版本错配。确认自托管 server 的情况。如果你自托管 Prefect server除了升级客户端包还需要升级数据库并保持客户端版本与 server 兼容通常客户端版本不大于 server 版本详见版本兼容规则。执行升级pip install -U prefect自托管 server 的用户接着运行数据库升级命令prefect server database upgrade使用了集成或 extra 时同时升级对应的集成包例如pip install -U prefect[aws]验证升级结果安装文档Install Prefect给出的确认命令是prefect version文档中的示例输出如下注意这是文档示例具体数值以你环境的实际输出为准重点看Version与Pydantic versionVersion: 3.4.24 API version: 0.8.4 Python version: 3.12.8 Git commit: 2428894e Built: Mon, Oct 13, 2025 07:16 PM OS/Arch: darwin/arm64 Profile: local Server type: server Pydantic version: 2.11.7 Server: Database: sqlite SQLite version: 3.47.1升级后版本应为 3.x且Pydantic version为 2.x。若你自托管 server恢复之前暂停的 deployment 调度。处理 Pydantic 2Prefect 3.0 基于 Pydantic 2.0 构建。Prefect 自身的对象会自动升级需要你自己处理的是如果 flow 参数或自定义 block 使用了自定义 Pydantic 模型必须让它们兼容 Pydantic 2.0。仅在你自己代码内部使用、不与 Prefect 直接交互的 Pydantic 1.0 模型可以继续使用。官方建议参考 Pydantic 的迁移指南Pydantic 官方文档中的 migration 页面查看 1.x 到 2.0 需要做的具体改动。检查模块重命名与移除的导入路径部分较少使用的模块被重命名、重组或移除。旧的导入路径仍会保留6 个月的支持但会发出 deprecation 警告。完整的受影响路径列表在仓库的 migration.py即官方文档所称的 deprecation code中其中MOVED_IN_V3列出旧路径 → 新路径的映射例如prefect.client:get_client→prefect.client.orchestration:get_clientprefect.engine:pause_flow_run/resume_flow_run/suspend_flow_run→prefect.flow_runs下对应函数prefect.deployments:load_flow_from_flow_run→prefect.flows:load_flow_from_flow_run而REMOVED_IN_V3列出已移除对象及替代方案例如prefect.agent:PrefectAgent改用 workers、prefect.filesystems:S3改用prefect_aws.S3Bucket等。按警告提示把代码改成新导入路径即可这是有 6 个月缓冲的迁移项不必在升级当天全部改完但应尽早规划。处理行为层面的破坏性变更以下是升级指南中标注的行为变化按你是否受影响逐一检查同步 flow 中调用 async taskPrefect 2.0 允许在同步 flow 中调用原生asynctask非常规 Python 用法3.0 移除了该行为。如果你的 flow 依赖这种模式必须把 flow 改成异步或改用支持异步执行的 task runner。Flow 最终状态的判定规则变了2.0 中只要有任一 task run 失败flow run 就会被标记为 failed。3.0 中 flow run 的最终状态只由两点决定flow 函数的return值字面量视为成功显式返回的State即为最终状态返回State可迭代对象时全部Completed才算Completed任一Failed则 flow 为Failed。flow 函数是否让异常raise出去异常向上传播导致Failed被raise_on_failureFalse压制的异常不影响 flow 状态。也就是说任务失败不再自动导致 flow run 失败。官方给出的三种让关键任务失败即 flow 失败的方案# 方案一让 task 异常向上传播不使用 raise_on_failureFalse from prefect import flow, task task def failing_task(): raise ValueError(Task failed) flow def my_flow(): failing_task() # Exception propagates, causing flow failure try: my_flow() except ValueError as e: print(fFlow failed: {e}) # Output: Flow failed: Task failed# 方案二return_stateTrue 并显式检查 task 状态 from prefect import flow, task from prefect.states import Failed task def failing_task(): raise ValueError(Task failed) flow def my_flow(): state failing_task(return_stateTrue) if state.is_failed(): raise ValueError(state.result()) return Flow completed successfully# 方案三try/except 捕获并返回失败状态 from prefect import flow, task from prefect.states import Failed task def failing_task(): raise ValueError(Task failed) flow def my_flow(): try: failing_task() except ValueError: return Failed(messageFlow failed due to task failure) return Flow completed successfully print(my_flow()) # Output: Failed(messageFlow failed due to task failure)任务自动缓存3.0 引入了幂等引擎默认情况下同一 flow run 内用相同输入被多次调用的 task 会被自动缓存。如果你的 task 依赖副作用每次调用都必须真正执行给 task 传cache_policyNone关闭缓存。unmapped()传入可变对象2.0 中用unmapped()给 mapped task 传 list/dict 等可变对象时每个 task 实例看起来拿到的是独立副本3.0 中该对象是按引用共享的任一 task 修改它其他并发 task 都会看到变化可能产生竞态。官方推荐的做法是在 task 内部先做深拷贝再修改from prefect import flow, task, unmapped from copy import deepcopy task def process_item(item: str, shared_state: dict): # Create a copy for this task local_state deepcopy(shared_state) local_state[current_item] item print(fProcessing {local_state[current_item]}) flow def my_flow(): shared_state {current_item: None} process_item.map( item[A, B, C], shared_stateunmapped(shared_state) )如果需要在多个 task 间共享只读配置改用不可变类型字符串、元组、frozen dataclass更安全。调度参数schedule改为schedules如果你调用Flow.serve或Flow.deploy时用到过schedule参数会报TypeError: Flow.deploy() got an unexpected keyword argument schedule。3.0 中该参数改名为schedules列表且调度对象来自prefect.schedules# Prefect 2.0 写法 from datetime import timedelta from prefect import flow from prefect.client.schemas.schedules import IntervalSchedule flow def my_flow(): pass my_flow.serve( namemy-flow, scheduleIntervalSchedule(intervaltimedelta(minutes1)) )# Prefect 3.0 写法 from datetime import timedelta from prefect import flow from prefect.schedules import Interval flow def my_flow(): pass my_flow.serve( namemy-flow, schedules[Interval(timedelta(minutes1))] )升级后常见报错排查官方把两类高频报错单独列为常见坑gotchasAttributeError: coroutine object has no attribute some attribute在异步 task/flow 上下文中没有await异步方法就会触发。修复方式是await该方法或传_syncTrue。例如Secret.load在异步上下文中是异步方法# 错误未 await my_secret Secret.load(my-secret) print(my_secret.get()) # AttributeError: coroutine object has no attribute get # 正确await my_secret await Secret.load(my-secret) # 也可以显式同步 my_secret Secret.load(my-secret, _syncTrue)RuntimeWarning: coroutine some_async_callable was never awaited与上面是同一类问题——存在未被 await 的协程同样用await或_syncTrue修复。TypeError: object some Type cant be used in await expression对非协程对象使用了await。注意PrefectFuture在 3.x 中是标准同步接口my_task.submit(...)始终是同步调用不应awaittask async def my_task(): pass flow async def my_flow(): # future await my_task.submit() # TypeError future my_task.submit() # This will work可选分支从 agents 迁移到 workers如果你还在用早期 Prefect 2.0 的 agents3.0 中 workers 已是标准REMOVED_IN_V3中prefect.agent:PrefectAgent的提示信息即指向 workers。仓库内有专门的迁移文档 upgrade-agents-to-workers 描述 agents → workers 的升级路径不属于当前升级主路径但使用 agents 的读者应把它列入迁移清单。版本兼容性限制自托管场景下官方建议所有客户端版本与 server 相同或更旧新版本客户端可能依赖旧 server 尚不支持的能力例如 3.1.0 客户端可用于 3.4.0 server反之不行。Prefect Cloud 则兼容所有客户端版本。弃用功能的保留期为至少 3 个 minor 版本或 6 个月取较长者所以旧导入路径的 6 个月缓冲期是明确的迁移窗口。Prefect 的版本号不完全遵循语义化版本minor 版本升级也可能包含不向后兼容的变更因此跨版本升级时应持续阅读 release notes官方在 versioning 中提示了这一点。完成上述步骤后你的代码库就完成了一次 2.x → 3.0 的升级prefect version显示 3.x 与 Pydantic 2.x、自托管数据库已执行prefect server database upgrade、代码不再使用旧导入路径和schedule参数、且你按 flow 的实际错误处理需求选择了最终状态方案。【免费下载链接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.项目地址: https://gitcode.com/GitHub_Trending/pr/prefect创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表