Prefect 3 is a Python workflow orchestration framework that lets you turn plain Python functions into resilient, observable data pipelines using decorators for flows and tasks, deployments that decouple code from infrastructure, and work pools that route runs to Docker, Kubernetes, or serverless targets. As of Prefect 3.8.5 (released September 3, 2026), the framework ships with first-class event-driven automations, transactional task semantics with on_commit and on_rollback hooks, and an in-process PrefectDbtRunner that treats every dbt model as a lineage node.
Prefect 3.8.5 (September 3, 2026) is the current stable release; it requires Python 3.9+ and is a full rewrite of the 2.x engine focused on runtime task graphs, native async, and event-driven automations.
Flows are ordinary Python functions decorated with @flow; tasks (@task) get retries, caching, mapping, and transactional commit/rollback semantics without any DAG boilerplate.
Deployments decouple code from infrastructure: flow.serve(), prefect deploy, and prefect.yaml all attach schedules and infra config, and work pools route runs to Docker, Kubernetes, ECS, ACI, or Google Cloud Run.
Automations react to events (flow state changes, custom emits, webhooks) and trigger actions like starting a downstream flow, and the same event bus is now open source, not just Prefect Cloud.
Prefect acquired Dagster Labs on July 13, 2026; both products keep their names, licenses, and roadmaps, and Prefect 3 remains the recommended path for Python-native pipelines.
The new PrefectDbtRunner executes dbt Core in-process with per-node lineage emission, replacing the older DbtCoreOperation and ShellOperation patterns.
What's new in Prefect 3 in 2026
Prefect 3.0 went generally available on September 3, 2024, and the 3.x line has spent the two years since then hardening the runtime and pulling former Cloud-only features into the OSS engine. As of the current 3.8.5 release, the pieces of the platform I actually touch every week are the runtime task graph (task dependencies are resolved at execution time, not parse time), the native async engine, transactional semantics on every task, and the events and automations system.
The single biggest shift compared to Prefect 2 is that events and automations are now part of the open-source package. Previously you needed Prefect Cloud to react to arbitrary events; today the same trigger-to-action model runs against a self-hosted server. I have a pipeline at work that used to be a cron-scheduled Airflow DAG polling S3 every five minutes; the Prefect 3 version listens for a custom bronze.file_landed event and starts within a couple of seconds of the object being written.
The other announcement that matters for planning: Prefect acquired Dagster Labs on July 13, 2026. Nothing changed for users of either tool (Dagster kept its name and Apache 2.0 license), but if you were putting off a Prefect adoption because you were worried about consolidation, that worry has settled. Both products continue under their own governance, and Prefect 3 remains the recommended path for teams that want a Python-native, code-first orchestrator.
Install Prefect 3 and write your first flow
Prefect 3 requires Python 3.9+ and installs cleanly with any modern package manager. I use uv at work because we run a lot of ephemeral containers and I need lockfile-based reproducibility. Plain pip works fine too.
Once installed, a flow is just a Python function with a decorator. There are no DAG classes, no operator inheritance, and no separate "declare then run" step. Here's the pipeline I used the first time I onboarded a new engineer on my team, trivial enough to explain over Slack, real enough to catch the common gotchas:
fetch_prices.map(symbols) fans the task out with a Prefect future per element. In Prefect 2 you had to reason about the DaskTaskRunner or ConcurrentTaskRunner to get real parallelism. In 3.x, the default engine runs task submissions concurrently on a thread pool, and I only reach for a task runner when I have GPU work or need Ray/Dask specifically.
Tasks, retries, and caching that survives restarts
The task decorator is where most of the operational value lives. The parameters I set on almost every task in production are retries, retry_delay_seconds (which accepts a list for explicit backoff, or a callable for jitter), timeout_seconds, and cache_key_fn. Prefect 3 supports exponential backoff, jitter, and per-exception retry predicates via retry_condition_fn, which is useful for making a task retry only on 429/5xx and fail fast on 4xx.
The caching model is the piece I explain most often to teams migrating from Airflow. Prefect writes cache entries at transaction commit, not at task completion. In practice that means a cached upstream result is only reused if the downstream task also succeeded. The first time we ran a bronze-to-silver refresh in production and it half-failed, I expected the cache to be poisoned, and it wasn't. The successful upstream had staged its result but not committed it. Rerunning the flow re-executed both tasks together, which is exactly what I wanted.
So, if you're coming from a scripted-orchestrator background and have never thought hard about caching versus idempotency, our async ETL in Python with httpx and asyncio guide walks through the same concurrency patterns without an orchestrator on top. Worth a read before you decide where a task boundary belongs.
Deployments and work pools explained
A flow is the code. A deployment is a durable server-side record of how, where, and when that code should run. You can trigger the same flow from a Python REPL forever, but the moment you want a schedule, an API to fire runs from outside Python, or infrastructure that spins up on demand, you need a deployment.
There are three ways to create one, and they trip people up on day one, so worth naming them explicitly.
flow.serve() keeps a long-running process alive that polls for scheduled runs of one specific flow. Great for a laptop demo or a small always-on VM; not what you want in Kubernetes.
flow.deploy() (called inside a Python script) or prefect deploy (CLI, reads prefect.yaml) register a deployment against the server, then a worker polls a work pool and executes runs on infrastructure the pool describes. This is what I run in production.
prefect.yaml is the declarative config file that prefect deploy consumes. Pin schedules, parameters, work pool, and infra overrides here, commit it to git, and CI runs prefect deploy --all on merge to main.
A minimal prefect.yaml for the market data flow above:
Work pools are the abstraction that lets a deployment describe "run this somewhere" without the orchestrator ever touching your data. A pool has a type (process, docker, kubernetes, ecs, azure-container-instance, cloud-run-v2, vertex-ai) and a default job template. Workers running in your VPC pull scheduled runs from that pool and execute them in the described infrastructure. The Prefect server or Cloud never opens an inbound connection into your network. This is the same hybrid pattern Airflow 3's Edge Executor adopted, and it's why Prefect works cleanly for teams with strict network egress rules.
Automations and event-driven pipelines
The pattern that changed how I design pipelines in 2026 is event-driven orchestration. Prefect 3 emits events for every flow lifecycle change (started, running, completed, failed, cancelled), for changes to variables/blocks/artifacts, and for anything you emit yourself with emit_event(). An automation binds a trigger (an event pattern) to an action (start a flow run, send a Slack message, pause a deployment, cancel matching runs). No cron, no polling.
from prefect.events import emit_event
@task
def promote_partition(partition: str) -> None:
# ... move object from staging to gold ...
emit_event(
event="gold.partition.landed",
resource={"prefect.resource.id": f"gold.partition.{partition}"},
payload={"partition": partition, "rows": 1_204_991},
)
On the Prefect UI or via the CLI, I bind a trigger ("when gold.partition.landed fires for any resource") to an action: "start deployment materialize_dashboards/prod with the payload as parameters." The downstream flow starts within a few seconds. The upstream flow doesn't need to know it exists.
Events also arrive from outside Prefect. A webhook endpoint (prefect webhook create) turns any HTTP POST from Fivetran, dbt Cloud, an SNS-to-SQS bridge, or a GitHub Action into a Prefect event. Combine that with an S3 EventBridge rule and you have the "file lands in bucket, run pipeline" pattern without ever writing a poller. If the downstream flow needs to validate that the payload matches the contract the producer promised, our data contracts with Pydantic and Pandera writeup covers the schema-enforcement layer I put on top of every event trigger.
Transactions, commit hooks, and rollback
Every task in Prefect 3 runs inside a transaction. The lifecycle has four phases: BEGIN (compute the cache key and look up any existing record; if present, skip execution and return the cached result), STAGE (execute the function and stage a result at the configured storage location), COMMIT (persist the record), and ROLLBACK (on any error inside the transaction, discard staged results and run on_rollback hooks).
Tasks can be grouped into a larger transaction with the transaction context manager. If any task inside the block fails, the entire block rolls back, including successful earlier tasks. It's the closest thing Python has to a database transaction across side-effectful I/O.
from prefect import flow, task
from prefect.transactions import transaction
from pathlib import Path
@task
def write_partition(partition: str, rows: list[dict]) -> Path:
p = Path(f"/mnt/silver/{partition}.parquet")
p.write_bytes(_to_parquet(rows))
return p
@write_partition.on_rollback
def delete_partition(txn):
p: Path = txn.get("staged_path")
if p and p.exists():
p.unlink()
@task
def dq_check(path: Path) -> None:
# any assertion failure will roll back the enclosing transaction
assert _row_count(path) > 0
@flow
def silver_refresh(partition: str, rows: list[dict]) -> None:
with transaction() as txn:
p = write_partition(partition, rows)
txn.set("staged_path", p)
dq_check(p) # if this fails, delete_partition fires
The on_rollback hook is the piece I care about most. It fires when the enclosing transaction fails, not just when the task itself throws. That means a downstream data-quality assertion can undo a file that a successful upstream task wrote. In one of my own pipelines, we use this to reverse a Snowflake MERGE if the post-merge invariant check fails: the merge succeeded, the invariant didn't, and the rollback hook issues a compensating UPDATE that restores the previous state. The alternative is a stale-partition alert at 3am, which I no longer accept.
Running dbt Core from Prefect with PrefectDbtRunner
If you already have a dbt Core project (I have several, and contributed patches to a couple), the shortest wiring in 2026 is PrefectDbtRunner, the new in-process interface in prefect-dbt. It replaces the older DbtCoreOperation and the generic ShellOperation pattern, and it emits one lineage event per dbt node so the flow-run page shows every model, seed, and snapshot as an individual resource.
from prefect import flow
from prefect_dbt import PrefectDbtRunner
@flow(name="analytics-build")
def analytics_build(select: str = "tag:daily", threads: int = 8):
runner = PrefectDbtRunner() # picks up DBT_PROFILES_DIR etc from env
runner.invoke(["build", "--select", select, "--threads", str(threads)])
if __name__ == "__main__":
analytics_build()
Two behaviours are worth calling out. First, PrefectDbtRunner raises on failure by default; you can set raise_on_failure=False if you treat failed tests differently from failed models. Second, environment variables prefixed with DBT_ are auto-detected via PrefectDbtSettings, a Pydantic BaseSettings subclass, so a CI job that already exports DBT_PROFILES_DIR and DBT_TARGET needs no further configuration.
Honestly, for fine-grained orchestration where I want per-node retries and selection filtering inside the flow, I reach for PrefectDbtOrchestrator with ExecutionMode.PER_NODE, which materializes each dbt model as its own Prefect task. That gives me task-level retries on transient adapter errors, and I've seen enough BigQuery quota flaps to know I want them.
If you write dbt tests as part of your build (and you should; see our dbt unit testing guide), Prefect's per-node lineage makes it obvious which failing test blocked which downstream marts. That's a step-change from tailing run_results.json after the fact.
Prefect vs Airflow vs Dagster in 2026
The three tools have converged more than the marketing pages let on. All three now support hybrid execution, event-driven triggers, and Python-first APIs. The choice is less about capability and more about model. Prefect optimizes for turning existing Python into orchestrated Python with the smallest possible ceremony; Dagster optimizes for a typed asset graph with software-defined data assets at the center; Airflow 3 keeps its DAG mental model but has adopted the TaskFlow API, assets, and dynamic task mapping.
Dimension
Prefect 3.8
Airflow 3.3
Dagster 1.9
Latest release
3.8.5 (Sep 2026)
3.3.1 (Nov 2025)
1.9.x (2026)
Primary abstraction
Flow (function) + Task
DAG + Task
Software-defined Asset
Runtime task graph
Yes (built at run)
Yes (TaskFlow + expand())
Yes
Event-driven triggers
Native (open source)
Asset watchers
Sensors + auto-materialize
Transactions / rollback
Native (on_rollback)
No first-class primitive
Asset checks + partitions
Hybrid execution
Workers + work pools
Edge Executor (AIP-69)
Code locations + agents
dbt integration
PrefectDbtRunner
DbtTaskGroup (Cosmos)
dbt assets (native)
Best for
Python-heavy pipelines, ML/data eng convergence
Established DAG shops, hybrid stacks
Type-checked lakehouse teams
If you want the deep dive on how the three compare on developer experience and total cost, our Airflow vs Prefect vs Dagster comparison benchmarks each on the same workload. My short answer today: pick Prefect if your team's day-to-day language is Python and you don't want to teach anyone a DAG DSL; pick Dagster if you're building a lakehouse and want asset lineage first; pick Airflow 3 if you already run Airflow and the migration cost outweighs the marginal wins.
Production pitfalls I keep hitting
Four things bite me repeatedly, and I have not yet found a way to make Prefect prevent them at the API level.
Entrypoints are relative to CWD. Running prefect deploy from a subdirectory silently registers a deployment whose entrypoint resolves against the wrong root, and the first run fails with ModuleNotFoundError. Always run prefect deploy from the project root, and enforce it in CI with if [[ ! -f prefect.yaml ]]; then exit 1; fi.
Cron schedules default to UTC. A 0 9 * * * schedule attached to a deployment fires at 09:00 UTC, not local time, unless you set timezone: in prefect.yaml. I've had a payment-recon pipeline miss its SLA for a week because someone assumed "9am" meant local time.
Docker deployments and push=True. If you set push: true on a Docker work pool job template, Prefect calls docker push at deploy time. If your CI runner isn't logged in to the registry, docker push fails silently, the deployment succeeds, and workers later fail with "image not found." Add a docker login step in CI and check the push exit code explicitly.
ShellOperation vs PrefectDbtRunner cancellation. If you still use ShellOperation or the older DbtCoreOperation, Prefect's cancellation sends SIGTERM but doesn't escalate to SIGKILL. dbt subprocesses can hang and the flow run sits in Running with no new logs. The in-process PrefectDbtRunner avoids this entirely, and it's the reason I've migrated every dbt-in-Prefect pipeline I own.
Yes. The prefect Python package and the self-hosted Prefect server are Apache 2.0 licensed. Prefect Cloud is the paid hosted product with managed multi-tenancy and long-term event retention, but every core feature (flows, tasks, deployments, work pools, automations, transactions) is available in the OSS package as of 3.0.
Does Prefect 3 replace Airflow?
Not automatically. Prefect and Airflow solve overlapping problems with different mental models: Prefect starts from a Python function and adds orchestration around it, while Airflow starts from a DAG and treats Python callables as one operator among many. For greenfield Python-native pipelines Prefect 3 is usually a faster path; for shops with hundreds of existing DAGs, staying on Airflow 3 and adopting TaskFlow is often the pragmatic choice.
What is the difference between a flow and a deployment in Prefect?
A flow is the Python function decorated with @flow, which is code you can run from a script or a notebook. A deployment is a server-side record that pins a specific flow entrypoint, a schedule, parameters, and infrastructure config to a work pool, so it can be triggered by the API, the UI, an automation, or a cron. You can have many deployments per flow.
How do you handle secrets in a Prefect deployment?
Use Prefect Blocks for structured secrets (a Secret block, or provider blocks like SnowflakeConnector) and reference them by name from your flow code. In production I prefer to keep the block record in Prefect but source the actual value from AWS Secrets Manager or Vault via a small init container, so the secret never lives in the Prefect database.
Can Prefect 3 run dbt Core natively?
Yes. Install prefect-dbt and use PrefectDbtRunner to execute dbt in-process from a flow, with per-model lineage emitted as Prefect events. For per-node retries and selection filtering inside a flow, PrefectDbtOrchestrator with ExecutionMode.PER_NODE materializes each dbt model as its own Prefect task.
Daniel is a staff data engineer with 13 years across fintech and logistics. He spent four years at Plaid building the transaction-enrichment pipeline (Python + Kafka + Snowflake), three years before that at Flexport on the freight-visibility data platform, and started his career at IBM doing DB2 performance work he still grudgingly draws on.
He writes about the gluework of modern Python data stacks: Prefect 2 flow design, dbt run orchestration from Python, Pydantic-based contract validation between Bronze and Silver layers, and the operational realities of running polars in containers with strict memory limits. He has contributed patches to dbt-core and to the prefect-snowflake integration.
Daniel is based in Lagos and Lisbon depending on the quarter, holds AWS Solutions Architect Professional, and writes a small newsletter about data-platform postmortems.
A hands-on tour of Apache Airflow 3.3 with the airflow.sdk Task SDK: TaskFlow decorators, task groups, dynamic task mapping with .expand(), assets, native DAG versioning, and a 2.x migration playbook drawn from moving 900 DAGs.
A practical 2026 guide to dbt unit tests: given/expect syntax, dict/csv/sql fixture formats, testing incremental models and macros, and running the suite in CI on DuckDB.
Compare GPTQ, AWQ, bitsandbytes, and GGUF for LLM quantization in Python. Real H100 benchmarks, kernel choices, and a production-ready decision tree for 2026.