Execution Order of Features in Prefect's Task Engine: Retries, Caching, and Persistence
Prefect's task engine processes every task run through an 18-step deterministic pipeline where parameter resolution and concurrency setup precede execution, followed by transaction-based caching and result persistence, with retry logic evaluated only after failed persistence attempts.
Prefect's task engine transforms a Python function decorated with @task into a managed TaskRun through a strictly ordered orchestration lifecycle. Understanding the execution order of features in Prefect's task engine—specifically how retries, caching, and persistence interact—is essential for debugging flow failures and optimizing data pipeline performance. This analysis examines the source code in PrefectHQ/prefect to map the exact sequence implemented in src/prefect/task_engine.py.
The 18-Step Task Engine Pipeline
The SyncTaskRunEngine and AsyncTaskRunEngine classes (both inheriting from BaseTaskRunEngine) execute every task run through identical phases. The following sections trace the pipeline using line references from src/prefect/task_engine.py.
Initialization and Setup (Steps 1-7)
Before user code executes, the engine prepares the runtime environment through seven setup phases:
-
Engine Construction:
BaseTaskRunEngine.__post_init__(lines 52-55) instantiates the engine with theTaskobject, parameters, and optionalwait_fordependencies. -
Parameter Resolution:
_resolve_parameters(lines 299-327) recursively resolves input values to their final results. This ensures downstream caching sees concrete inputs rather than futures. -
Custom Task-Run Naming:
_set_custom_task_run_name(lines 329-444) resolves dynamic names after parameters are ready, allowing runtime values to influence the task run name. -
Wait-For Dependencies:
_wait_for_dependencies(lines 446-558) guarantees that anywait_forobjects complete before execution begins. -
Local TaskRun Creation:
_create_task_run_locally(lines 33-126) builds a lightweightTaskRunobject in memory with an initial Pending state. -
Pending Event Emission:
emit_task_run_state_change_event(lines 528-558) notifies the Prefect server that the task is queued. -
Concurrency Lease Acquisition: The
_concurrencycontext manager (lines 489-555) acquires a lease for the task's tags, enforcing global resource limits before the task occupies a worker slot.
Execution and State Management (Steps 8-11)
The engine transitions to active execution and manages the user function lifecycle:
-
Transition to Running:
begin_run(lines 416-432) creates aRunningstate, recordsstart_time, and pushes the state to the server. -
Retry State Handling: If the task enters AwaitingRetry or Paused,
begin_runcontains a while-loop (lines 442-452) that polls the state usingclamped_poisson_intervalbackoff until it becomes Running. -
Context Installation and Timeout:
run_context(lines 990-1026) installs the task-run context and wraps execution in a timeout guard, raisingTaskRunTimeoutErrorif the task exceeds its time limit. -
User Function Invocation:
call_task_fn(lines 274-280) invokes the user function viacall_with_parameters, handing the raw result tohandle_successor routing exceptions tohandle_exception.
Caching and Persistence (Steps 12-14)
After the user function returns, the engine handles storage through a transactional pattern:
-
Transaction and Cache Key Computation:
transaction_context(lines 964-990) opens a transaction. The engine computes a cache key viacompute_transaction_key(lines 71-98) based on the task'sCachePolicy. If the key exists in the cache, the stored result is returned immediately; otherwise, the new result is staged. -
Result Persistence:
return_value_to_state_sync(sync engine) orreturn_value_to_state(async engine) writes the output to the configuredResultStore(lines 361-424). This step is controlled byshould_persist_result()and respects the task'spersist_resultflag. -
Transaction Commit: If
transaction.is_committed()succeeds (lines 447-456), the terminal state is renamed "Cached" to indicate the result was stored for future runs.
Finalization and Error Handling (Steps 15-18)
The engine closes the loop and handles failure scenarios:
-
Final State Update:
set_state(lines 556-607) updates the TaskRun's denormalized fields (state_id,run_count, etc.) and emits the final state-change event. -
Retry Evaluation: If the task raised an exception,
handle_retry(lines 601-695) evaluatestask.retries,retry_delay_seconds, andretry_condition_fn. If retries remain, it sets AwaitingRetry or Retrying state and loops back to step 8. -
Failure Recording: Unrecoverable exceptions trigger
handle_exception(lines 714-728) orhandle_crash(lines 730-751), transforming them intoFailedorCrashedstates while recording telemetry. -
Cleanup: The
finallyblock instart(lines 622-642) callslog_finished_messageand resets internal flags, releasing the concurrency lease.
Why Execution Order Matters
The sequencing in src/prefect/task_engine.py protects data integrity through three critical design choices:
-
Retries occur after persistence attempts: The engine only evaluates
handle_retry(step 16) after attempting to persist results (step 13). This prevents stale cache entries from failed runs from polluting the result store. -
Caching is transaction-protected: The cache key computation (step 12) happens within
transaction_context, ensuring atomic read/write operations. A failed transaction never writes partial data to the cache. -
State transitions drive orchestration: Every phase change emits events (steps 6, 8, 15), keeping the Prefect server synchronized with the local engine state via the WebSocket connection.
Code Examples
The following examples demonstrate how the execution order affects real-world usage.
Configuring Retries
from prefect import task, flow
@task(retries=2, retry_delay_seconds=5)
def fragile(x: int) -> int:
if x < 3:
raise ValueError("x is too small")
return x * 2
@flow
def demo():
# Step 16 (handle_retry) triggers twice with 5-second delays
result = fragile.submit(1)
result.wait()
print(result.result()) # Prints 8 after retries complete
demo()
Implementing Caching with CachePolicy
from prefect import task, flow
from prefect.cache_policies import CachePolicy
@task(cache_policy=CachePolicy())
def compute_heavy(a: int, b: int) -> int:
return a ** b
@flow
def cache_demo():
# Step 12 computes cache key; step 14 commits to cache
first = compute_heavy.submit(2, 10)
first.wait()
# Second call hits cache at step 12, skips user function
second = compute_heavy.submit(2, 10)
second.wait()
assert first.result() == second.result()
Disabling Result Persistence
from prefect import task, flow
@task(persist_result=False) # Skips step 13
def transient(x: int) -> int:
return x + 1
@flow
def persistence_demo():
r = transient.submit(41)
r.wait()
# No result stored in ResultStore
Core Source Files
| File | Description |
|---|---|
src/prefect/task_engine.py |
Core SyncTaskRunEngine and AsyncTaskRunEngine implementation containing the full 18-step pipeline, retry handling, and state management. |
src/prefect/cache_policies.py |
CachePolicy definition and compute_key method used in step 12 for cache key generation. |
src/prefect/results.py |
ResultStore and ResultRecord utilities for step 13 persistence operations. |
src/prefect/utilities/engine.py |
Helpers for parameter resolution (collect_task_run_inputs_sync) and state-change event emission. |
src/prefect/states.py |
State definitions (Pending, Running, AwaitingRetry, Retrying, Completed, etc.) that the engine transitions between. |
src/prefect/concurrency/_asyncio.py |
Async concurrency manager used in step 7 for lease acquisition. |
src/prefect/concurrency/_sync.py |
Sync concurrency manager used in step 7. |
Summary
- The execution order of features in Prefect's task engine follows 18 deterministic phases defined in
src/prefect/task_engine.py. - Parameter resolution (step 2) and concurrency leasing (step 7) occur before any user code executes to ensure inputs are concrete and resources are available.
- Caching (step 12) and persistence (step 13) happen immediately after the user function returns, protected by
transaction_contextto prevent race conditions. - Retry logic (step 16) evaluates only after persistence attempts, ensuring failed runs do not corrupt the result store or cache.
- All state transitions emit events, maintaining synchronization between the local engine and the Prefect server.
Frequently Asked Questions
Does Prefect check the cache before or after running the task function?
The engine computes the cache key in step 12 (compute_transaction_key, lines 71-98) immediately after the user function executes but before persisting results. If the key matches an existing entry, the transaction returns the cached value immediately, skipping function invocation on subsequent runs. On the first run, the function executes, then the result is stored if the transaction commits successfully.
What happens if a task fails after result persistence begins?
The transaction context (step 12) protects against partial writes. If the user function raises an exception, control passes to handle_exception (lines 714-728) or handle_retry (lines 601-695) before the transaction commits. Since return_value_to_state_sync (step 13) only executes within handle_success, failed runs never write partial results to the ResultStore.
How does Prefect handle retries with cached results?
Retries are evaluated in step 16 (handle_retry) only after a failed execution attempt. If the retry policy permits another attempt, the engine transitions to AwaitingRetry or Retrying state and loops back to step 8 (begin_run). The cache key remains constant across retry attempts, ensuring that once a successful run completes and commits (step 14), future executions with identical inputs hit the cache instead of retrying the logic.
Can I disable persistence but keep caching enabled?
Yes. Setting persist_result=False on a task skips step 13 (return_value_to_state_sync) entirely. However, the transaction and cache key computation in step 12 still occur. This configuration allows in-memory caching within a single flow run while preventing writes to persistent storage, though it limits cache availability to the current process.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →