Prefect2中任务失败时如何挂起后续任务、流程及调度?
Great question! Let's break down exactly how to implement your requirements in Prefect 2—stopping downstream tasks when my_favorite_function() fails, pausing the current flow run, and halting future scheduled executions until the issue is fixed. We'll also cover specific exception handling patterns tailored to Prefect's model.
1. Pausing Downstream Tasks & Current Flow Run
By default, Prefect marks downstream tasks as "Skipped" if an upstream task fails, but you want to pause the entire flow instead (so you can manually intervene before resuming). Here's how to do that by catching the specific ValueError and returning a Paused state:
First, update your code to include flow-level exception handling:
from prefect import flow, task, get_run_logger from prefect.states import Paused @task def my_favorite_function(): raise ValueError("This flow immediately fails") @task def one_return(): return 1 @task def tow_return(): return 2 @flow def run_flow(): logger = get_run_logger() try: # Run your critical task first my_favorite_function() # These tasks will only run if the critical task succeeds one_return() tow_return() except ValueError as e: logger.error(f"Critical error hit: {str(e)}") # Return a Paused state to halt the flow and prevent downstream tasks return Paused(message=f"Flow paused - fix the error in my_favorite_function() to resume")
What this does:
- If
my_favorite_function()throws theValueError, thetry-exceptblock catches it immediately. - Returning
Pausedtells Prefect to stop executing downstream tasks (one_return()andtow_return()) and leave the flow run in a paused state in the UI. - You can then manually resume the flow run once the error is fixed (or restart it from scratch).
2. Halting Future Scheduled Runs
To stop the 5-minute recurring schedule from triggering new runs while the error persists, you need to interact with Prefect's API to disable the flow's schedule. Here's how to extend the code to pause scheduling automatically:
Add these imports first:
from prefect.client import get_client from prefect.orion.schemas.actions import FlowUpdate
Then update the flow function to pause scheduling on failure:
@flow async def run_flow(): # Note: We make the flow async to use the client logger = get_run_logger() try: my_favorite_function() one_return() tow_return() except ValueError as e: logger.error(f"Critical error hit: {str(e)}") # Pause the current flow run pause_state = Paused(message=f"Flow paused - fix the error in my_favorite_function() to resume") # Disable future scheduled runs async with get_client() as client: # Fetch the flow by name flow_obj = await client.get_flow_by_name("run_flow") if flow_obj: # Update the flow to disable scheduling await client.update_flow( flow_obj.id, FlowUpdate(is_scheduled_active=False) ) logger.info("Flow schedule has been paused until the error is resolved.") return pause_state
What this does:
- The flow is now async to allow using Prefect's async client.
- When the
ValueErroris caught, we use the client to fetch the flow and setis_scheduled_active=False, which turns off all future scheduled runs. - Once you fix the error in
my_favorite_function(), you'll need to re-enable the schedule manually via the Prefect UI or API.
3. Task-Level Exception Handling (Alternative Approach)
If you prefer handling the exception directly at the task level (instead of the flow), you can use a state handler to catch the error and pause the task (which will block downstream tasks):
from prefect import task, get_run_logger from prefect.states import Paused def critical_error_handler(task, old_state, new_state): logger = get_run_logger() # Check if the task failed with our specific ValueError if isinstance(new_state.result, ValueError): logger.error(f"Task {task.name} failed critically - pausing task.") # Return a Paused state for the task return Paused(message=f"Task failed: {str(new_state.result)}") return new_state @task(state_handlers=[critical_error_handler]) def my_favorite_function(): raise ValueError("This flow immediately fails")
With this setup, if my_favorite_function() throws the ValueError, the task will enter a paused state. Downstream tasks will wait for this task to complete (since they depend on it by default), effectively halting the flow until you manually resume the task.
Key Notes
- Resuming: After fixing the error, you'll need to:
- Resume the paused flow/task run via the Prefect UI.
- Re-enable the flow's schedule (if you disabled it) to restore the 5-minute cadence.
- Specific Exceptions: The examples above target
ValueError, but you can adjust theexceptblock or state handler to catch any specific exception type your task might throw. - Logging: Always use
get_run_logger()to track what's happening—this will help you debug and confirm that the pause logic is working as expected.
内容的提问来源于stack exchange,提问作者Alexander

