如何让Prefect 2中任务等待前置任务成功执行?
问题
在Windows系统上使用Prefect 2配置多任务执行时,部分任务偶尔出错,修复后需手动重新运行。希望其他任务仅等待前置任务成功完成后再启动,无需手动触发失败的调度。当前代码中task_that_should_wait会等待task_that_fails结束,但无论其成功或失败都会继续执行。由于两个任务调度时间不同,属于独立Flow,无法合并。请问如何配置task_that_should_wait使其仅等待task_that_fails的成功完成?
当前代码:
from prefect import flow, task from prefect.deployments import Deployment from prefect.server.schemas.schedules import CronSchedule @task def task_that_fails(): raise(Exception("Failed here")) # 假设这是会手动修复的bug return True @flow def flow_that_fails(): task_that_fails() deploy_that_fails = Deployment.build_from_flow( name="task_that_fails", flow=flow_that_fails, schedule=CronSchedule(cron="36 21 * * MON,TUE,WED,THU,FRI", timezone='US/Central'), ) deploy_that_fails.apply() @task def task_that_should_wait(): print("hey!") return True @flow def flow_that_waits(): task_that_should_wait(wait_for=['task_that_fails']) deploy_that_waits = Deployment.build_from_flow( name="deploy_that_waits", flow=flow_that_waits, schedule=CronSchedule(cron="35 21 * * MON,TUE,WED,THU,FRI", timezone='US/Central'), ) deploy_that_waits.apply()
解决方案
要实现仅等待前置Flow成功完成的逻辑,需要利用Prefect 2的状态感知任务依赖和Flow运行查询功能,具体实现如下:
- 将等待Flow改为异步模式,使用Prefect客户端API查询目标Flow的最新成功运行记录
- 过滤状态为
COMPLETED的运行实例,仅当存在该实例时才执行后续任务;否则终止当前Flow,等待下一次调度
修改后的代码:
from prefect import flow, task, get_client from prefect.deployments import Deployment from prefect.server.schemas.schedules import CronSchedule from prefect.server.schemas.filters import FlowRunFilter, FlowRunFilterState from prefect.server.schemas.states import StateType @task def task_that_fails(): raise(Exception("Failed here")) # 假设这是会手动修复的bug return True @flow def flow_that_fails(): task_that_fails() deploy_that_fails = Deployment.build_from_flow( name="task_that_fails", flow=flow_that_fails, schedule=CronSchedule(cron="36 21 * * MON,TUE,WED,THU,FRI", timezone='US/Central'), ) deploy_that_fails.apply() @task def task_that_should_wait(): print("hey!") return True @flow async def flow_that_waits(): # 获取Prefect客户端连接 async with get_client() as client: # 查询flow_that_fails的最新成功运行记录 flow_runs = await client.read_flow_runs( flow_filter=FlowRunFilter(name={"equals": "flow_that_fails"}), state_filter=FlowRunFilterState(type={"equals": StateType.COMPLETED}), sort="-created", limit=1 ) if not flow_runs: # 无成功记录时终止当前Flow,等待下一次调度或手动触发 raise RuntimeError("flow_that_fails尚未成功执行,当前Flow终止") # 获取最新成功的Flow运行实例,让任务等待该实例(确保依赖的是成功状态) successful_run = flow_runs[0] task_that_should_wait(wait_for=[successful_run]) deploy_that_waits = Deployment.build_from_flow( name="deploy_that_waits", flow=flow_that_waits, schedule=CronSchedule(cron="35 21 * * MON,TUE,WED,THU,FRI", timezone='US/Central'), ) deploy_that_waits.apply()
关键逻辑说明
- 异步Flow:因为Prefect客户端API为异步实现,所以将
flow_that_waits改为异步Flow以兼容查询操作 - 状态过滤:通过
FlowRunFilterState筛选出状态为COMPLETED的Flow运行实例,确保依赖的是成功执行的记录 - 失败处理:如果未找到目标Flow的成功记录,直接抛出异常终止当前Flow,避免无意义的执行
内容的提问来源于stack exchange,提问作者Mefitico
相关产品推荐
相关产品推荐

