How Do Prefect Transactions Ensure Atomicity in Data Pipelines

Prefect transactions ensure atomicity by wrapping task groups in a context manager that stages all changes and commits them only when the entire block succeeds, automatically rolling back partial work if any task fails.

Prefect transactions provide ACID-like guarantees for data workflows in PrefectHQ/prefect, ensuring that grouped operations either complete entirely or leave no trace. By implementing a sophisticated state machine and tight integration with ResultStore, these transactions prevent partial failures from corrupting pipeline state. Understanding how Prefect transactions ensure atomicity in data pipelines is essential for building reliable, production-grade data systems that handle failures gracefully.

The Transaction State Machine

Atomicity relies on explicit lifecycle tracking. In src/prefect/transactions.py, the TransactionState enum defines a strict progression that every transaction follows: PENDING → ACTIVE → STAGED → COMMITTED or ROLLED_BACK (lines 59-64).

This state machine ensures that work cannot be partially persisted. When you enter a transaction context, the state transitions from PENDING to ACTIVE as the system calls prepare_transaction() and begin(). Only when all tasks complete successfully does the state move to STAGED and finally COMMITTED. If any exception occurs, the state immediately shifts to ROLLED_BACK, triggering cleanup logic.

Context Manager Lifecycle and Commit Modes

The Transaction class implements __enter__ and __exit__ methods (lines 64-81 and 70-88) to handle the atomic boundary. When entering the context, prepare_transaction() validates isolation level compatibility and begin() initializes the transaction state. On exit, the context manager evaluates whether to call commit() or rollback() based on exception status and the configured CommitMode.

Prefect supports three commit modes defined in the CommitMode enum (lines 53-56):

  • EAGER: Commits immediately when the transaction context exits successfully
  • LAZY: Defers commitment until explicitly requested or until the parent transaction commits
  • OFF: Disables automatic commitment, requiring manual control

These modes allow downstream tasks to decide when the atomic block should persist, giving fine-grained control over transaction boundaries in complex flows.

Isolation Levels for Concurrent Safety

Atomicity extends to concurrent execution through isolation levels. The IsolationLevel enum (lines 48-51) supports READ_COMMITTED (default) and SERIALIZABLE.

With READ_COMMITTED, transactions see previously committed results but do not block concurrent writes. When using SERIALIZABLE, the Transaction.begin method (lines 99-107) acquires a distributed lock on the transaction key using a LockManager (such as FileSystemLockManager or MemoryLockManager). This prevents race conditions by ensuring only one transaction with a given key can execute at a time, effectively serializing access to shared resources.

ResultStore Integration for Durability

Transactions integrate with ResultStore (from src/prefect/results.py) to ensure durability without compromising atomicity. During execution, results are staged in memory; the Transaction.commit logic (lines 57-66) writes to the ResultStore only when write_on_commit=True at commit time.

If the transaction rolls back, the store remains untouched, guaranteeing no partial data persists. This "all-or-nothing" persistence model ensures that downstream systems never see incomplete or inconsistent states, even if the pipeline fails midway through a multi-step transformation.

Rollback Hooks and Side Effects

Real-world pipelines often interact with external systems that lack native transaction support. Prefect addresses this through rollback hooks (on_rollback) and commit hooks (on_commit) executed via Transaction.run_hook (lines 93-106).

These hooks let you attach compensating actions that execute atomically with the transaction outcome. For example, if a task writes to a file or sends a message to a queue, you can register a cleanup function that removes the file or sends a compensating message if the transaction rolls back.

Nested Transaction Support

Data pipelines often require composition where sub-flows need atomic guarantees within larger atomic blocks. Prefect supports nested transactions where child transactions inherit the parent's context.

If a child transaction fails, calling Transaction.reset (lines 121-130) propagates the rollback signal to all parent transactions in the stack. This ensures that a failure anywhere in the nested hierarchy rolls back the entire atomic set, maintaining atomicity across complex, multi-layered workflows.

Implementing Atomic Blocks in Practice

Basic Atomic Execution with Rollback

The following example demonstrates a transaction that writes a file only if quality checks pass, with automatic cleanup on failure:

from prefect import task, flow
from prefect.transactions import transaction
import os

@task
def write_file(contents: str):
    with open("side-effect.txt", "w") as f:
        f.write(contents)

@write_file.on_rollback
def delete_file(txn):
    os.unlink("side-effect.txt")

@task
def quality_test():
    with open("side-effect.txt") as f:
        if len(f.readlines()) < 2:
            raise ValueError("Not enough data!")

@flow
def pipeline(contents: str):
    with transaction() as txn:
        write_file(contents)
        quality_test()  # If this fails, delete_file runs automatically

Idempotent Execution Using Transaction Keys

Transactions support idempotency through unique keys, preventing duplicate execution of expensive operations:

from prefect import flow, task
from prefect.transactions import transaction

@task
def download_data():
    return "expensive data"

@task
def write_data(data: str):
    with open("data.txt", "w") as f:
        f.write(data)

@flow
def idempotent_pipeline():
    with transaction(key="download-and-write-v1") as txn:
        if txn.is_committed():
            print("Already completed – skipping.")
            return
        data = download_data()
        write_data(data)

Preventing Race Conditions with Serializable Isolation

For concurrent flows accessing shared resources, use SERIALIZABLE isolation with a lock manager:

import threading
from prefect import flow, task
from prefect.transactions import transaction, IsolationLevel
from prefect.locking.filesystem import FileSystemLockManager
from prefect.results import ResultStore
from prefect.settings import PREFECT_HOME

@task
def download_data():
    return f"{threading.current_thread().name} data"

@task
def write_file(contents: str):
    with open("shared-resource.txt", "w") as f:
        f.write(contents)

@flow
def concurrent_pipeline(key: str):
    with transaction(
        key=key,
        isolation_level=IsolationLevel.SERIALIZABLE,
        store=ResultStore(
            lock_manager=FileSystemLockManager(
                lock_files_directory=PREFECT_HOME.value() / "locks"
            )
        ),
    ) as txn:
        if txn.is_committed():
            return
        data = download_data()
        write_file(data)

# Only one thread will execute the critical section

key = "shared-work-unit"
threading.Thread(target=concurrent_pipeline, args=(key,)).start()
threading.Thread(target=concurrent_pipeline, args=(key,)).start()

Summary

Prefect transactions ensure atomicity in data pipelines through several coordinated mechanisms:

  • Explicit state management via TransactionState enum prevents partial commits by enforcing strict lifecycle transitions in src/prefect/transactions.py
  • Configurable commit modes (EAGER, LAZY, OFF) control when staged changes persist to external systems
  • Isolation levels (READ_COMMITTED, SERIALIZABLE) prevent race conditions through distributed locking when using FileSystemLockManager or similar implementations
  • ResultStore integration guarantees that data is written only upon successful commit, with write_on_commit ensuring no partial persistence
  • Rollback hooks provide compensating transactions for external side effects, executing cleanup logic automatically on failure
  • Nested transaction support ensures that failures in child transactions propagate rollback to parent contexts, maintaining atomicity across complex workflow hierarchies

Frequently Asked Questions

What is the difference between EAGER and LAZY commit modes in Prefect transactions?

EAGER mode commits the transaction immediately when the context manager exits successfully, making results available to downstream tasks right away. LAZY mode defers the commit until the parent transaction commits (or until explicitly triggered), allowing you to batch multiple child transactions into a single atomic unit. According to the source code in src/prefect/transactions.py, this distinction is handled in the __exit__ method based on the CommitMode enum value.

How does SERIALIZABLE isolation prevent race conditions in concurrent flows?

When you set isolation_level=IsolationLevel.SERIALIZABLE, the Transaction.begin method acquires a lock on the transaction key using the configured LockManager (such as FileSystemLockManager). This ensures that only one transaction with that specific key can execute at a time across all concurrent workers. As implemented in src/prefect/transactions.py (lines 99-107), this serializes access to shared resources, preventing the "check-then-act" race conditions that occur when multiple flows attempt to write to the same destination simultaneously.

Can I use Prefect transactions with asynchronous flows?

Yes. Prefect provides AsyncTransaction with async/await support through the atransaction context manager. The implementation mirrors the synchronous Transaction class, offering identical atomicity guarantees for asynchronous workflows. The __aenter__ and __aexit__ methods handle the async lifecycle, while still maintaining the same state machine and ResultStore integration found in the synchronous version.

What happens to data in ResultStore if a transaction rolls back?

If a transaction rolls back, the ResultStore remains completely untouched. The transaction only writes staged data to the store when commit() is called with write_on_commit=True. If an exception occurs or rollback() is explicitly invoked, the transaction exits without calling the store's write methods, ensuring that no partial or corrupted data persists. This behavior provides the durability guarantee that once a transaction commits, the data is safe, but until then, the system remains in its original state.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →