Prefect Runner 模块化架构深度解析:薄门面 + 单一职责服务如何驱动 Flow Run 全生命周期

发布时间:2026/9/12 5:21:57
Prefect Runner 模块化架构深度解析:薄门面 + 单一职责服务如何驱动 Flow Run 全生命周期 Prefect Runner 模块化架构深度解析薄门面 单一职责服务如何驱动 Flow Run 全生命周期【免费下载链接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.项目地址: https://gitcode.com/GitHub_Trending/pr/prefect导读Runner 是 Prefect 中负责管理所有 Deployment 执行的运行时组件通过flow.serve或serve工具创建 Deployment 后Runner 负责轮询调度、提交运行、跟踪子进程、处理取消与崩溃并向 API 汇报状态。本文以仓库内 src/prefect/runner/AGENTS.md 为骨架结合src/prefect/runner/下各抽取类的源码实现剖析 Runner 从 1757 行单体类重构为薄门面 单一职责服务后的内部架构你将掌握各核心类的职责边界、服务启停的 LIFO 依赖顺序、基于 TCP 回环通道的取消意图协商机制、状态转换拆分逻辑以及存储拉取与 Git 仓库的安全校验细节。Runner 的定位与重构背景Runner 负责管理所有 Deployment 的执行。通过flow.serve或serve工具创建 Deployment 时Runner 还会轮询调度运行。runner.py顶部的 docstring 给出了典型用法见 src/prefect/runner/runner.pyimport time from prefect import flow, serve flow def slow_flow(sleep: int 60): Sleepy flow - sleeps the provided amount of time (in seconds). time.sleep(sleep) flow def fast_flow(): Fastest flow this side of the Mississippi. return if __name__ __main__: slow_deploy slow_flow.to_deployment(namesleeper, interval45) fast_deploy fast_flow.to_deployment(namefast) # serve generates a Runner instance serve(slow_deploy, fast_deploy)重构前src/prefect/runner/runner.py是一个 1757 行的单体类将三类截然不同的操作纠缠在一起轮询模式start()轮询调度运行并通过_get_and_submit_flow_runs提交、单次运行模式execute_flow_run()通过子进程直接运行或执行python -m prefect.engine、Bundle 模式execute_bundle()在multiprocessing.SpawnProcess中执行序列化 Bundle。单体难以隔离测试每个单元测试都要启动完整 Runner、难以推理三种执行模式共享状态、难以安全扩展。完整的设计动因与目标架构记录在 plans/completed/2026-02-18-runner-refactor.md 中。重构后的指导原则被写进AGENTS.md的第一句话Runner 是薄门面thin facade新行为属于抽取出的单一职责类而不是 Runner 本身。架构总览门面委托给哪些类Runnerrunner.py将工作委托给下表这些类每个类在 src/prefect/runner/ 目录下有独立文件类文件职责FlowRunExecutor_flow_run_executor.py单次运行生命周期提交submit→ 启动start→ 等待wait→ 崩溃/钩子crashed/hooksProcessManager_process_manager.py进程映射、PID 追踪、SIGTERM→SIGKILL 逐级终止StateProposer_state_proposer.py所有 API 状态转换提议CancellationManager_cancellation_manager.py控制通道信号 → kill → 钩子 → 状态 → 事件的取消序列CancelFinalizer_cancel_finalizer.pykill 后持久化 Cancelled 状态状态无法确认时回退为 CrashedControlChannel_control_channel.pyRunner 侧 TCP 回环 IPC用于在 kill 前向子进程传递取消意图HookRunner_hook_runner.pyon_cancellation / on_crashed 钩子执行EventEmitter_event_emitter.py通过 EventsClient 发射事件任何 WebSocket 连接失败时降级为 NullEventsClientLimitManager_limit_manager.py并发限制DeploymentRegistry_deployment_registry.pyDeployment/flow/storage/bundle 映射ScheduledRunPoller_scheduled_run_poller.py轮询循环、运行发现、调度ProcessStarter协议_flow_run_executor.py启动进程的策略接口FlowRunExecutorContext_flow_run_executor.py在 Runner 之外worker、CLI、bundle进行一次性执行的异步上下文管理器DirectSubprocessStarter_starter_direct.py通过run_flow_in_subprocess直接运行 Flow 对象EngineCommandStarter_starter_engine.py生成python -m prefect.engine子进程WorkspaceResolvingEngineCommandStarter_workspace_starter.py启动受管 supervisor准备 workspace、启动所选引擎命令并为钩子记录准备好的运行时BundleExecutionStarter_starter_bundle.py在 SpawnProcess 中执行序列化 Bundle其中FlowRunExecutor是运行生命周期的核心每个运行构建一次、submit()完成后即弃其生命周期为见 src/prefect/runner/_flow_run_executor.pypropose_submitting—— 服务器拒绝则提前返回已取消预检cancelling/cancelled—— 打日志并提前返回通过 starter 启动进程 ——task_status.started(handle)提前通知调用方将句柄加入process_manager阻塞直到进程退出starter.start 在发出 started 后阻塞快照终态 attempt 证据然后移除进程注册引擎回执receipt视为已处理否则解释非零退出码仅当 Runner 提议 Crashed 时运行崩溃钩子。核心契约新行为放哪、遗留方法怎么处理AGENTS.md用两段硬性契约约束演进方向。第一段新行为必须落入抽取类而不是门面。如果修复或新增功能进程生命周期 →ProcessManager或FlowRunExecutor状态转换 →StateProposer关闭/崩溃处理 →FlowRunExecutor.submit()取消 →CancellationManager钩子 →HookRunner第二段Runner 上的遗留方法仅为向后兼容而存在禁止添加新行为。以下方法均已找到替代实现遗留方法已被替换为_submit_run_and_capture_errors()FlowRunExecutor.submit()_run_process()ProcessStarter各实现_flow_run_process_map字典ProcessManager_kill_process()ProcessManager.kill()_run_on_crashed_hooks()/_run_on_cancellation_hooks()HookRunnerexecute_flow_run()已废弃2026 年 3 月改用FlowRunExecutorContextWorkspaceResolvingEngineCommandStarterexecute_bundle()已废弃2026 年 3 月改用prefect.bundles.execute中的execute_bundle()reschedule_current_flow_runs()已废弃2026 年 3 月SIGTERM 重调度改由 CLI 执行路径内联处理这些方法会在内部调用方迁移完成后移除。当前的迁移状态截至本文写作时直接 ProcessWorker 运行已使用FlowRunExecutorContext生成的命令使用WorkspaceResolvingEngineCommandStarter显式配置的命令使用EngineCommandStarter只有 ad hoc bundle 路径仍在调用已废弃的Runner.execute_bundle()。生命周期行为必须留在抽取类中不要将直接 worker 执行绕回 Runner。服务生命周期AsyncExitStack 的 LIFO 顺序Runner.__aenter__期间服务按以下顺序进入退出时严格逆序这是一条硬约束弄错会导致关闭时出现ClosedResourceErrorclient—— 最后退出所有服务都需要它process_manager—— 第 5 个退出运行结束后 kill 剩余进程limit_manager—— 第 4 个退出event_emitter—— 第 3 个退出在 client 关闭前 flush 事件runs_task_group—— 第 2 个退出等待在途运行完成cancelling_observer—— 最先退出在任务结束前停止取消检测。新增服务时必须仔细插入到这条序列中。同样的 6 步依赖顺序也以注释形式固化在FlowRunExecutor的类 docstring 中src/prefect/runner/_flow_run_executor.py其中第 6 位是cancellation_manager最先退出避免对部分拆除状态做取消检测。取消机制Attempt Control Session 的意图与回执ControlChannel_control_channel.py与子进程侧的prefect._internal.control_listener构成一个带认证、按 attempt 作用域的 TCP 回环会话。v1 对端保留单字节控制协议v2 对端在带内协商回执支持可以选择用一条结构化 Engine Outcome Receipt 结束会话。第一条有效回执或被确认的控制意图胜出ControlChannel.get_conclusion()暴露这条不可变的证据unregister()在清理注册的同时返回它。关键的安全与一致性设计第一条通过认证的连接消费该 attempt 的 token重放与连接替换都会被拒绝线缆常量、回执类型与共享意图字节映射都位于prefect._internal.attempt_control被 runner.py 以from prefect._internal.attempt_control import AttemptConclusion引用通道不替代平台的 kill 信号取消时 Runner 仍会走ProcessManager.kill()发送SIGTERMWindows 上为CTRL_BREAK_EVENT宽限期后升级为SIGKILL。通道额外提供的价值是预植入意图等子进程的SIGTERM处理器运行时control_listener.get_intent()已返回预植入的意图引擎的except TerminationSignal块因此可以分发到on_cancellation而非on_crashedreschedule/relinquish则不动状态。取消如何使用它3 步Runner 在通道上发送cancel意图并等待最多 1 秒的子进程ba确认_DEFAULT_ACK_TIMEOUT 1.0见 _control_channel.py若子进程确认意图在SIGTERM到达前已在子进程内提交引擎的except TerminationSignal块就能分发到on_cancellation钩子而不是on_crashedRunner 随后无论确认与否都继续执行ProcessManager.kill()。POSIX 与 Windows 的差异POSIXack 仅表示SIGTERM 桥已就绪、意图已植入。Runner 真实的SIGTERM仍是唯一能中断阻塞代码的触发源ack 后立即 killWindowsack 表示子进程已排队_thread.interrupt_main(SIGTERM)。Runner 给子进程 30 秒宽限期自行退出grace_seconds 30.0见 _cancellation_manager.py超时才回退到外部 kill。失败模式若子进程从未连接或在 1 秒内未确认signal()返回False。CancellationManager回退到普通优雅 kill引擎将终止视为崩溃。而prefect flow-run executesupervisor 则直接调用ProcessManager.kill(forceTrue)SIGKILL、无宽限因为未确认的引擎会提议Crashed从而撤销reschedule/relinquish。CancellationManager.cancel()的完整序列为 kill → 钩子 → 状态 → 事件见 _cancellation_manager.pyProcessLookupError进程已消失被视为预期继续其它 kill 异常中止序列钩子失败记为 ERROR 但状态与事件照常执行钩子运行时使用合成的Cancelling状态而非flow_run.state快照以规避并发取消路径观察者事件重复、websocket 与轮询竞争导致的钩子静默跳过。扩展意图当前意图为cancel、reschedule与relinquish。新增意图需同时更新prefect._internal.attempt_control中的共享字节映射与引擎的TerminationSignal分发。EventEmitter 的 WebSocket 降级策略事件是非关键遥测Runner 只发射prefect.runner.cancelled-flow-run因此事件 WebSocket 不可达时绝不能拖垮 flow run。EventEmitter.__aenter___event_emitter.py捕获连接失败静默将失败客户端替换为NullEventsClient。被捕获的失败分为两类_NONFATAL_CONNECTION_EXCEPTIONS见 src/prefect/runner/_event_emitter.py永久性拒绝websockets.exceptions.InvalidStatusHTTP 4xx——例如配置了PREFECT_SERVER_API_AUTH_STRING的服务器上≤3.6.13 的客户端连接 ≥3.6.14 的服务器瞬态失败来自prefect.events.clients的RETRYABLE_EXCEPTIONSConnectionClosed、TimeoutError、OSError且已超出事件客户端自身的重连次数。降级时记录一条WARNING。一个容易踩坑的细节如果__aenter__抛出异常不能在原始客户端上调用__aexit__它从未成功进入替换的NullEventsClient才会被正常进入/退出。ProcessStarter 策略模式执行模式可插拔每种执行模式都有一个ProcessStarter实现。ProcessStarter是结构化类型协议见 _flow_run_executor.py任何具有合规start()签名的对象都算数无需继承测试替身可以直接用AsyncMock或普通异步函数。其契约是在阻塞之前调用task_status.started(handle)让调用方在进程完成前就能拿到ProcessHandle阻塞直到进程退出进程退出后返回None。要新增一种执行模式就实现该协议并把实例注入FlowRunExecutor而不是给 Runner 添加新的代码路径。当前实现覆盖四种模式Starter文件启动方式DirectSubprocessStarter_starter_direct.pyrun_flow_in_subprocess直接运行 Flow 对象EngineCommandStarter_starter_engine.pypython -m prefect.engine子进程WorkspaceResolvingEngineCommandStarter_workspace_starter.py受管 supervisor先准备 workspace再启动所选引擎命令并为钩子记录准备好的运行时BundleExecutionStarter_starter_bundle.py在SpawnProcess中执行序列化 Bundle状态转换拆分ScheduledRunPoller 与 FlowRunExecutor 各管一段ScheduledRunPoller现在先调用propose_pendingScheduled → Pending再把运行交接给FlowRunExecutorFlowRunExecutor随后在生命周期第 1 步调用propose_submittingPending → Submitting 子状态默认propose_submittingTrue。这是两次独立转换不能合并——拆分是为了让监听 Pending 状态的自动化automation能在 executor 开始工作前正确触发。StateProposer_state_proposer.py是无状态服务封装所有 API 状态转换提议propose_pending、propose_submitting、propose_crashed等。其中propose_submitting对Abort会重新读取权威状态以区分已经是 Submittingworker/bundle 路径与真正的拒绝。三个调用方通过FlowRunExecutorContext.create_executor(propose_submittingFalse)关闭第二次提交——它们都已把运行推进过 Pending 状态再提议 Submitting 就是错的直接ProcessWorker.run()执行BaseWorker已提议 Submittingprefect flow-run executeCLI 路径由 worker 调用prefect.bundles.execute中的execute_bundle()由 bundle dispatch 调用。注意取消预检第 1a 步即使propose_submittingFalse也会无条件执行。Worker 执行路径与门面边界ProcessWorker 的执行路径区分如下直接 ProcessWorker 运行使用FlowRunExecutorContext且propose_submittingFalse生成的命令使用WorkspaceResolvingEngineCommandStarter命令选择跟随 workspace 准备显式配置的命令使用EngineCommandStarter保留自身的 pull-step 语义。executor 统一负责进程跟踪、取消、Attempt Control Session 证据、崩溃推断与规范化基础设施状态。ad hoc ProcessWorker 提交仍调用已废弃的Runner.execute_bundle()是待迁移目标。存储拉取与 Git 仓库的安全边界BlockStorageAdapter每次 pull 都先清空目标目录BlockStorageAdapter.pull_code()每次拉取都会清空目标目录而不仅是首次。若目标已存在先删除全部子项目录用shutil.rmtree文件与符号链接用unlink再通过block.get_directory()写入新内容。为什么必须这样get_directory通常调用shutil.copytree(..., dirs_exist_okTrue)它无法覆盖只读文件——git pack 文件权限为0o444不清空的话第二次 pull 会触发PermissionError。符号链接处理目标目录内的目录符号链接用unlink()删除而不是rmtree()——这删除的是链接本身保留链接目标。不要改成rmtree()那会跟随符号链接并删除目标目录。GitRepository 输入校验防参数注入GitRepository.__init__storage.py强制两条不显眼的约束commit_sha必须匹配^[0-9a-fA-F]{4,64}$——任何不匹配的值包括--upload-pack...这类 git 选项字符串都会抛出ValueError源码见 src/prefect/runner/storage.py。分支/标签名必须走branch参数directories中以--开头的条目触发UserWarning但不拒绝。这些值会传给git sparse-checkout set --带--分隔符防止标志注入。警告存在是因为此类路径不寻常但合法用途被允许。这些校验的目的就是防止 git 参数注入程序化构造GitRepository时不得绕过。GitRepository 并发拉取保护FileLockGitRepository.pull_code()通过FileLockprefect.locking._filelock序列化并发调用。锁文件紧挨目标目录放置destination.parent / (destination.name .lock)见 storage.py。崩溃进程留下的陈旧锁通过 PID 检查自动恢复。宽泛地 monkeypatchpathlib.Path.exists的测试必须考虑到目标目录旁会创建这个锁文件。Storage Base Path 作用域deployment.path中的$STORAGE_BASE_PATH来自RunnerDeployment.from_storage()。对于 work-pool deployment创建时path被设为None存储被序列化进pull_steps而非 path见 src/prefect/deployments/runner.py 中对应逻辑。load_flow_from_flow_run()仅在pull_steps缺席时才做$STORAGE_BASE_PATH替换见 src/prefect/flows.py。推论CLI 的prefect flow-run execute路径基于 worker总有pull_steps不需要tmp_dir/PREFECT__STORAGE_BASE_PATH只有 Runner 自服务的 deployment无 work pool才使用该替换。Workspace Resolver 子进程_workspace_resolver.py在受管 workspace supervisor 内准备 flow run 的 workspace存储拉取、pull steps、CWD/环境/sys.path捕获。Stdout 只保留 JSONPreparedWorkspaceResult载荷。pull step 的输出含继承自子进程的 stdout被重定向到 stderr。解析process.stdout取结果、process.stderr看诊断——违反这一点会静默破坏调用方。调用方 API当调用方自己负责 workspace 准备时应使用WorkspaceResolvingEngineCommandStarter_workspace_starter.py而不是直接调用 resolver。将starter.hook_runner作为hook_runner参数传给FlowRunExecutorContext.create_executor()——hook runner 读取准备好的运行时清单并用与引擎相同的命令前缀、工作目录和环境在子进程中加载钩子。显式配置的命令被视为不透明opaqueProcessWorker 用EngineCommandStarter启动它们避免重复其既有 pull 行为。环境隔离workspace 准备、引擎执行、钩子加载都停留在子进程中。不要把已准备 workspace 的环境、CWD 或sys.path安装进长驻的 Runner 进程。本地路径就地执行image-baked 模式当未配置 pull steps 且 entrypoint 文件存在于deployment.path但不在 workspace 根目录时prepare_workspace把working_directory设为本地路径而非workspace_root下。workspace 目录仍会创建但不是执行根目录——代码位于固定本地路径例如容器内而非对象存储。Pull-step 目录回退当 pull steps 执行了但没有任何一个产生directory输出或改变 CWD 时_ensure_entrypoint_in_workspace把存储复制到 workspace 根目录作为回退。这覆盖了纯 setup 型 pull steps例如用于环境准备的run_shell_script不控制工作目录的情况。自动依赖安装uv默认情况下WorkspaceResolvingEngineCommandStarter不会在启动 flow run 前安装依赖。当未传入显式命令且PREFECT_RUNNER_AUTO_INSTALL_DEPENDENCIES为 true 时它会自动选择uv run --no-default-groups --project project_root配合 worker 提供的 workspace bootstrap但必须同时满足三个条件pyproject.toml存在于project_rootproject.dependencies列表包含prefect通过 workspace 的PATH环境变量能找到uv不是系统 PATH。bootstrap 在prefect.flow_engine暴露可执行_main时使用准备好的 entrypoint更老的选定运行时回退到prefect.enginePREFECT__FLOW_ENTRYPOINT以保留已准备的 workspace。若设置为 false 或任一条件失败命令回退为None显式命令永远优先。小结给 Runner 演进者的实践清单职责归属新生命周期行为放进抽取类FlowRunExecutor/ProcessManager/StateProposer/CancellationManager/HookRunner门面只做委托服务顺序新服务进入Runner.__aenter__时必须遵守 AsyncExitStack 的 LIFO 序列client 最后退出取消依赖 Attempt Control Session 预植意图但 kill 信号仍是最终触发源新增意图要同步更新prefect._internal.attempt_control的共享字节映射与引擎分发事件降级任何 WebSocket 连接失败都要降级为NullEventsClient事件绝不能让 flow run 崩溃状态转换propose_pending与propose_submitting是两次独立转换已被 worker/bundle 推进过 Pending 的路径必须传propose_submittingFalse存储安全commit_sha校验、--目录警告、FileLock 并发保护、pull 前清空目标目录这些是防注入与防只读冲突的底线不得绕过。相关测试分布在 tests/runner/如test__cancellation_manager.py、test__control_channel.py、test__process_manager.py、test__flow_run_executor.py、test_runner.py、test_storage.py等重构的完整设计与取舍见 plans/completed/2026-02-18-runner-refactor.md。【免费下载链接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.项目地址: https://gitcode.com/GitHub_Trending/pr/prefect创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考