资讯动态

Prefect 官方集成开发与发布指南:从包结构、工作流命令到基础设施装饰器的完整实践

发布时间:2026/9/12 7:30:01 来源:尧图企业网站定制
Prefect 官方集成开发与发布指南从包结构、工作流命令到基础设施装饰器的完整实践【免费下载链接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.项目地址: https://gitcode.com/GitHub_Trending/pr/prefect导读本文围绕 Prefect 仓库中 src/integrations/AGENTS.md 这份集成开发规范文档展开系统梳理 Prefect 官方集成Official Integrations的组织形态、关键契约、目录布局、日常开发与发布命令并深入解析集成专属的PrefectBaseSettings配置体系与基础设施装饰器 Bundle 步骤这一新形态 API。无论你是计划为 Prefect 贡献新集成、维护现有集成包还是想理解prefect-aws、prefect-gcp、prefect-redis等集成包内部如何运作本文都能让你快速建立完整的开发心智模型并直接复用到实际工作中。一、官方集成独立的 PyPI 包生态Prefect 的官方集成位于仓库 src/integrations/ 目录下目前包含 18 个面向外部服务和平台的扩展包prefect-aws、prefect-azure、prefect-bitbucket、prefect-dask、prefect-databricks、prefect-dbt、prefect-docker、prefect-email、prefect-gcp、prefect-github、prefect-gitlab、prefect-kubernetes、prefect-ray、prefect-redis、prefect-shell、prefect-slack、prefect-snowflake、prefect-sqlalchemy。与核心prefect包不同每个集成都是独立的 PyPI 包拥有自己的版本号、依赖声明和测试套件。这意味着集成包的发布节奏、依赖上限和破坏性变更都可以与核心库解耦用户按需安装即可pip install prefect-aws pip install prefect[aws] # 通过 extra 方式一并安装关键契约Key Contractssrc/integrations/AGENTS.md 明确规定了所有集成必须遵守的核心契约所有集成均处于 pre-1.0 阶段破坏性变更breaking changes通过提升 minor 版本号来表达而非 major 版本。例如prefect-dbt-0.7.x中引入 breaking change 时应发布0.8.0。新增集成必须先讨论贡献者在提交新集成的 PR 之前应当先开 issue 发起讨论一般情况下用户自定义的集成应该放在独立仓库中而不是并入本仓库。这条规则保证了官方集成目录的收敛性和维护成本可控。通过推送 tag 发布发布格式为prefect-name-semver例如prefect-dbt-0.7.20。justfile中的 unreleased-integrations 命令正是利用这一 tag 约定来统计自上次发布 tag 以来有提交的集成包它遍历src/integrations/prefect-*/目录用git tag -l ${pkg}-* --sort-v:refname | head -1找到每个包最近一次发布 tag再计算git log --oneline ${latest_tag}..HEAD -- $pkg_dir的提交数量从而列出所有未发布的变更。默认依赖 PyPI 上最新发布的prefect集成包始终面向核心库的正式发布版开发只有当你正在开发一个集成将要直接消费的核心接口时才允许使用核心 Prefect 的 editable 安装。凭据必须使用 Block绝不硬编码所有密钥、Token、连接字符串等敏感信息都必须通过 Prefect Block 管理严禁写死在 flow 代码或配置文件里。二、统一的集成包目录结构每个集成都遵循一致的结构参见 src/integrations/AGENTS.mdprefect-name/ ├── prefect_name/ # 源码blocks、tasks、workers ├── tests/ # 集成专属测试 ├── pyproject.toml # 包配置与依赖 ├── justfile # 任务运行命令 └── README.md以prefect-aws为例见 src/integrations/prefect-aws/其包根目录与 AGENTS.md 描述完全一致而prefect_aws/源码目录下则包含了更加细化的模块划分prefect_aws/ ├── __init__.py ├── credentials.py # AWS 凭据 Block ├── s3.py # S3 Block / 任务 ├── secrets_manager.py # Secrets Manager 集成 ├── batch.py # AWS Batch 任务 ├── glue_job.py # Glue Job 任务 ├── lambda_function.py # Lambda 任务 ├── decorators.py # 基础设施装饰器ecs ├── bundles/ # Bundle 上传/执行 CLI 步骤 │ ├── upload.py │ └── execute.py ├── experimental/ # 已废弃的向后兼容 shim ├── observers/ # ECS observer 等观测能力 ├── workers/ # ECSWorker 等 Worker 实现 ├── settings.py # 集成级 PrefectBaseSettings └── _cli/ # 面向基础设施的命令行工具这种目录划分并非prefect-aws独有——prefect-gcp、prefect-kubernetes等包同样采用workers/、decorators.py、bundles/的组织方式说明官方集成在演进过程中逐步沉淀出了一套统一的代码契约。集成级 justfile 与根目录 justfile 的分工每个集成目录下的justfile定义该包自身的常用命令。以 src/integrations/prefect-aws/justfile 为例just test # 运行该集成的全部测试uv run pytest just api-ref # 生成 API 参考文档api-ref命令的实现展示了文档生成的完整链路它会先回到仓库根目录git rev-parse --show-toplevel然后通过uvx --with-editable ./src/integrations/prefect-aws以可编辑方式注入该集成包调用mdxify工具以prefect_aws为根模块生成 API 参考文档输出到docs/integrations/prefect-aws/api-ref/。而仓库根目录的 justfile 则负责跨包的发布编排详见下一节两者的定位是包内日常开发与仓库级发布流程的分工。三、开发与发布常用命令速查日常开发命令所有日常命令都在集成目录内执行例如src/integrations/prefect-aws/uv run pytest # 运行该集成的全部测试 uv run pytest tests/ -k test_name # 按名称运行指定测试-k 支持模糊匹配 just api-ref # 重新生成 API 参考文档当需要运行一个依赖集成 extra 的脚本如复现脚本时从仓库根目录执行uv run --extra aws repros/1234.py # 以安装 prefect-aws 的方式运行脚本这里的--extra aws与pip install prefect[aws]语义一致uv run --extra name会在临时环境中按 pyproject 中声明的 extra 依赖解析并安装对应集成包随后运行指定脚本。发布命令发布相关命令从仓库根目录执行参见 justfilejust unreleased-integrations # 列出自上次发布 tag 以来有提交的集成 just prepare-integration-release pkg # 为某个集成生成发布说明如 prefect-awsjust prepare-integration-release pkg的内部流程见 justfile非常完整调用uv run scripts/prepare_integration_release_notes.py pkg生成发布说明草稿若该集成目录存在justfile且包含api-ref目标则自动重新生成 API 参考文档just --justfile $INTEGRATION_DIR/justfile --working-directory $INTEGRATION_DIR api-ref若设置了EDITOR环境变量自动打开生成的文件docs/v3/release-notes/integrations/pkg.mdx供人工审阅否则提示手动审阅提示后续步骤审阅发布说明、提交变更并开 PR。真正发布动作则遵循前文契约推送形如prefect-name-semver的 tag例如prefect-dbt-0.7.20CI 据此构建并发布到 PyPI。这也是unreleased-integrations命令能高效统计未发布变更的根本原因。四、基础设施装饰器与 Bundle 步骤新形态 API演进背景从部署到直接运行文档特别指出src/integrations/AGENTS.md对于支持直接在基础设施上运行 flow、无需创建 Deployment的集成官方引入了两套新形态 API基础设施装饰器Infrastructure Decorators位于包根目录的decorators.py例如from prefect_aws.decorators import ecsBundle 上传/执行 CLI 步骤位于bundles/子包例如prefect_aws.bundles.execute。同时集成中凡是存在 GA 路径即decorators.py/bundles/的experimental/子包都被标记为已废弃的向后兼容 shim它只是从 GA 路径重新导出符号并附带DeprecationWarning不要再向其中添加新代码。例外prefect-snowflake的experimental/workers/spcs.py是仍然活跃的 SPCS Worker 实现不在此废弃之列。装饰器的实现细节以prefect_aws.decorators.ecs为例src/integrations/prefect-aws/prefect_aws/decorators.py 中ecs装饰器的签名如下def ecs( work_pool: str, include_files: Sequence[str] | None None, launcher: BundleLauncher | None None, *, include_files_base_dir: Path | str | None None, **job_variables: Any, ) - Callable[[Flow[P, R]], InfrastructureBoundFlow[P, R]]:参数语义work_pool目标 ECS work pool 名称必填include_files可选的文件模式序列用于打包进 Bundle 并随 flow 一起上传到远程执行环境。模式相对于 flow 文件所在目录解析支持 glob 语法例如*.yaml、data/**/*.csvlauncher可选的 Bundle 上传与执行 Launcher 覆盖实现默认使用包内置实现include_files_base_dir解析include_files时的基准目录相对路径默认基于 flow 文件目录**job_variables传递给基础设施配置的额外作业变量。值得注意的实现细节_validate_include_files_syntaxdecorators.py会在装饰时刻而非运行时刻校验include_files的合法性——所有元素必须是字符串且不能为空或纯空白。这意味着语法错误可以在 import 阶段就被提前暴露而不是等到远程执行才失败体现了fail fast的工程理念。ecs返回InfrastructureBoundFlow其类型定义来自核心库的prefect.flowsfrom prefect.flows import InfrastructureBoundFlow, bind_flow_to_infrastructure说明该装饰器本质上是把 flow 与 ECS 基础设施绑定bind起来的语法糖。典型用法from prefect import flow from prefect_aws.decorators import ecs ecs(work_poolmy-ecs-pool, include_files[*.yaml, data/**/*.csv]) def my_flow(): ...Bundle 步骤upload与executeprefect_aws.bundles/目录下包含upload.py与execute.py两个模块bundles/分别对应 Bundle 的上传与执行两个 CLI 步骤。它们与核心库的 Bundle 机制prefect.bundles见 src/prefect/bundles/配合将 flow 及其依赖文件打成一个 Bundle上传到 S3 等远端存储再在 ECS 上拉取执行——这就是无 Deployment 直接运行 flow的完整链路。五、集成配置体系PrefectBaseSettings与环境变量前缀通用规则build_settings_config自动生成前缀需要运行时配置行为的集成会在包根目录的settings.py中定义PrefectBaseSettings子类核心基类定义见 src/prefect/settings/base.py。关键机制是调用build_settings_config((integrations, name, ...))它会根据传入的命名空间元组自动生成对应的环境变量前缀。例如 prefect-gcp 的 settings.py 中的class CloudRunV2WorkerSettings(PrefectBaseSettings): model_config build_settings_config( (integrations, gcp, cloud_run_v2, worker) )对应环境变量前缀为PREFECT_INTEGRATIONS_GCP_CLOUD_RUN_V2_WORKER_*。同样的规则适用于其它集成prefect-aws的 settings.py 使用(integrations, aws, ecs, worker)→PREFECT_INTEGRATIONS_AWS_ECS_WORKER_*以prefect-gcp为例其完整的设置层级为GcpSettings.cloud_run_v2.worker见 settings.py字段通过default_factory嵌套构造最终映射到环境变量PREFECT_INTEGRATIONS_GCP_CLOUD_RUN_V2_WORKER_*。用真实设置字段理解前缀规则以prefect-aws的 ECS Worker 设置为例settings.py字段与完整环境变量名的对应关系一目了然设置字段默认值环境变量语义create_task_run_max_attempts3PREFECT_INTEGRATIONS_AWS_ECS_WORKER_CREATE_TASK_RUN_MAX_ATTEMPTS创建 ECS task run 的最大重试次数应对集群扩容等瞬态资源约束create_task_run_min_delay_seconds1PREFECT_INTEGRATIONS_AWS_ECS_WORKER_CREATE_TASK_RUN_MIN_DELAY_SECONDS重试间最小固定延迟秒create_task_run_max_delay_jitter_seconds3PREFECT_INTEGRATIONS_AWS_ECS_WORKER_CREATE_TASK_RUN_MAX_DELAY_JITTER_SECONDS重试延迟的抖动上限秒register_task_definition_max_attempts3PREFECT_INTEGRATIONS_AWS_ECS_WORKER_REGISTER_TASK_DEFINITION_MAX_ATTEMPTS注册 task definition 的最大重试次数应对限流等瞬态错误register_task_definition_initial_delay_seconds1.0PREFECT_INTEGRATIONS_AWS_ECS_WORKER_REGISTER_TASK_DEFINITION_INITIAL_DELAY_SECONDS指数退避的初始延迟秒register_task_definition_max_delay_seconds10.0PREFECT_INTEGRATIONS_AWS_ECS_WORKER_REGISTER_TASK_DEFINITION_MAX_DELAY_SECONDS指数退避的最大延迟秒同样地prefect-gcp的 Cloud Run V2 Worker 设置settings.py对 job 创建、job 提交和只读 API 调用分别定义了*_max_attempts、*_initial_delay_seconds、*_max_delay_seconds三组参数采用指数退避 抖动exponential jitter backoff策略应对 Cloud Run API 的瞬态错误HTTP 429/500/503。例如export PREFECT_INTEGRATIONS_GCP_CLOUD_RUN_V2_WORKER_CREATE_JOB_MAX_ATTEMPTS5 export PREFECT_INTEGRATIONS_GCP_CLOUD_RUN_V2_WORKER_CREATE_JOB_INITIAL_DELAY_SECONDS2.0 export PREFECT_INTEGRATIONS_GCP_CLOUD_RUN_V2_WORKER_CREATE_JOB_MAX_DELAY_SECONDS30.0例外prefect-redis不使用PREFECT_INTEGRATIONS_REDIS_*prefect-redis是一个明确标注的例外src/integrations/AGENTS.md它不在包根目录维护settings.py而是把设置定义在各个模块内部并使用(redis, subsystem)命名空间。查看源码可以验证这一点messaging.py 中的RedisMessagingPublisherSettings使用(redis, messaging, publisher)→PREFECT_REDIS_MESSAGING_PUBLISHER_*字段包括batch_size默认 5、publish_every默认 10 秒、deduplicate_by同文件的RedisMessagingConsumerSettingsmessaging.py使用(redis, messaging, consumer)→PREFECT_REDIS_MESSAGING_CONSUMER_*字段包括block默认 1 秒、min_idle_time默认 5 秒、max_retries默认 3、trim_every默认 60 秒、trim_idle_threshold默认 5 分钟、should_process_pending_messages、starting_message_id、automatically_acknowledgecleanup_queue.py 中的RedisWorkerCleanupQueueSettings使用(redis, worker_cleanup_queue)→PREFECT_REDIS_WORKER_CLEANUP_QUEUE_*。因此不要指望prefect-redis支持PREFECT_INTEGRATIONS_REDIS_*形式的环境变量——在为其配置环境变量时必须使用上述PREFECT_REDIS_*前缀。与核心设置的区别独立的 Settings 层级与核心库prefect.settings不同src/prefect/settings/集成设置不会接入根Settings层级而是独立存在。访问方式不是通过全局配置对象而是直接实例化设置类from prefect_aws.settings import AwsSettings settings AwsSettings() print(settings.ecs.worker.create_task_run_max_attempts) # 默认 3在测试中通常使用mock.patch.dict(os.environ, {...})覆盖环境变量来验证配置行为这一测试模式在 tests/ 与各集成tests/目录中被广泛使用因为PrefectBaseSettings会从环境变量读取覆盖值。六、补充实践建议开发新集成前先在 issue 中发起讨论确认官方是否有意愿维护多数自定义集成应放置于独立仓库。版本管理牢记 pre-1.0 契约破坏性变更提升 minor 版本每次变更后通过just unreleased-integrations检查是否有未发布提交。发布前流程用just prepare-integration-release pkg生成发布说明并触发 API 文档重生成审阅docs/v3/release-notes/integrations/pkg.mdx后提交 PR正式发布以推送prefect-name-semvertag 为准。配置命名新集成如需运行时配置统一在包根目录settings.py用build_settings_config((integrations, name, ...))生成前缀除非你有充分的理由采用prefect-redis式的例外方案。敏感信息任何情况下凭据都通过 Block 管理参考 blocks/ 与各集成的credentials.py实现严禁硬编码。结语Prefect 官方集成生态以独立 PyPI 包 统一代码契约 版本化发布为骨架以decorators.py/bundles/的新形态 API 和PrefectBaseSettings配置体系为血肉构成了一个既可独立演进又与核心库深度协同的扩展体系。理解 src/integrations/AGENTS.md 这份规范是参与 Prefect 集成开发最直接的起点而本文结合prefect-aws、prefect-gcp、prefect-redis的真实源码对规范逐条印证希望能帮助你在此基础上快速上手无论是贡献官方集成还是构建自己的扩展。进一步阅读集成专属文档位于 docs/integrations/集成发布说明示例见 docs/v3/release-notes/。【免费下载链接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.项目地址: https://gitcode.com/GitHub_Trending/pr/prefect创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价