You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让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()

关键逻辑说明

  1. 异步Flow:因为Prefect客户端API为异步实现,所以将flow_that_waits改为异步Flow以兼容查询操作
  2. 状态过滤:通过FlowRunFilterState筛选出状态为COMPLETED的Flow运行实例,确保依赖的是成功执行的记录
  3. 失败处理:如果未找到目标Flow的成功记录,直接抛出异常终止当前Flow,避免无意义的执行

内容的提问来源于stack exchange,提问作者Mefitico

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 03:17:12