Prefect¶
A Python-first workflow orchestrator: add
@flowand@taskto ordinary functions and get retries, scheduling, observability and event-driven runs.
Last reviewed · Download PDF
Prerequisites: Python for DE · Apache Airflow
Related: Dagster · Docker · Data Ingestion & CDC · Glossary
Overview¶
Challenge: Most pipelines start as a Python script. As they grow you need retries when an API blips, a schedule, visibility into failures, parameters, and a way to run the same code on different machines, without rewriting the script into a framework's DAG syntax.
Solution: Prefect keeps your code as normal Python. A function decorated with @flow becomes a tracked, observable run, and @task marks the units inside it that need retries, caching or concurrency. Control flow is plain Python (if, for, try), not a graph definition, so dynamic pipelines are natural.
flowchart LR
C["Your Python code<br/>@flow / @task"] --> D["Deployment<br/>(schedule, parameters,<br/>work pool)"]
D --> S["Prefect server<br/>or Prefect Cloud<br/>(state, UI, automations)"]
S --> W["Worker in a work pool<br/>(process, Docker, Kubernetes)"]
W --> R["Flow run:<br/>states, logs, retries"]
R --> S
Relevance to data engineering: Prefect fits ingestion jobs, API pulls, ML and reverse-ETL workflows, and glue between systems, where the logic is Python and events matter more than a fixed nightly DAG.
On this page
Basic - Flows and Tasks - Retries, Timeouts and Logging - Running Locally
Intermediate - Parameters and Dynamic Workflows - Concurrency - Caching and Results - Deployments and Schedules
Advanced - Work Pools and Workers - Blocks, Variables and Secrets - Automations and Events - Testing - Prefect vs Airflow vs Dagster
Reference - Common Pitfalls - Cheat Sheet - Interview Questions - Further Reading
Flows and Tasks¶
from prefect import flow, task, get_run_logger
@task
def extract(day: str) -> list[dict]:
return [{"order_id": 1, "amount": 10.0}, {"order_id": 1, "amount": 10.0}, {"order_id": 2, "amount": 5.0}]
@task
def transform(rows: list[dict]) -> list[dict]:
return list({r["order_id"]: r for r in rows}.values()) # deduplicate by key
@task
def load(rows: list[dict]) -> int:
get_run_logger().info(f"loaded {len(rows)} rows")
return len(rows)
@flow(name="orders-daily", log_prints=True)
def orders_daily(day: str = "2024-03-15") -> int:
rows = extract(day)
n = load(transform(rows))
print(f"loaded {n} rows for {day}")
return n
if __name__ == "__main__":
orders_daily()
Run it with python orders.py. Prefect starts a temporary local server, records the flow and task runs and their states (Pending, Running, Completed, Failed, Retrying, ...), and prints them.
| Piece | Role |
|---|---|
| Flow | The unit you run, schedule and observe. Entry point of a workflow |
| Task | A step inside a flow. Adds retries, caching, timeouts and its own state |
| Flow run / task run | One execution, with a state history and logs |
| State | Where a run is: Completed, Failed, Cancelled, Crashed... Your code can react to it |
A flow can call other flows (subflows), which is how you compose large pipelines from testable pieces.
Retries, Timeouts and Logging¶
@task(retries=3, retry_delay_seconds=[1, 2, 4], timeout_seconds=30)
def extract(day: str) -> list[dict]:
... # a transient failure is retried after 1 s, 2 s, then 4 s
- Set retries on the tasks that touch flaky things: APIs, networks, shared databases. Do not retry non-idempotent work such as an unguarded
INSERT. retry_delay_secondsaccepts a list for backoff.retry_jitter_factoradds randomness so many retries do not hit a server at the same moment.- Flows have
retriestoo, for re-running the whole workflow. get_run_logger()sends logs to the run in the UI, andlog_prints=Truecapturesprintoutput.- Timeouts and blocking calls: a
timeout_secondson a synchronous task cannot interrupt a blocking call (a network request,time.sleep) in a worker thread. Use an async task, or set timeouts on the underlying client (for examplehttpx.Client(timeout=10)).
Running Locally¶
prefect server start # UI at http://127.0.0.1:4200
export PREFECT_API_URL=http://127.0.0.1:4200/api # point runs at that server
python orders.py
Without a server running, Prefect uses a temporary one per run, which is fine for trying things. Use the local server when you want to browse flow runs in the UI. For a hosted control plane, log in to Prefect Cloud (prefect cloud login).
Parameters and Dynamic Workflows¶
Flow parameters are type-validated, so a schedule or a UI run can pass day="2024-03-16" and Prefect checks it against your annotations.
Because the workflow is plain Python, loops and conditions need no special syntax:
@flow
def load_all(days: list[str]):
for day in days:
if day_has_data(day): # any Python condition
orders_daily(day) # call a flow like a function (a subflow)
Use .map() to fan a task out over a list, running the copies concurrently:
@flow
def extract_many(days: list[str]):
futures = extract.map(days) # one task run per day
return [f.result() for f in futures]
Concurrency¶
Tasks called with .submit() or .map() return futures and run on a task runner (a thread pool by default), so independent tasks overlap.
@flow
def parallel(days: list[str]):
futures = [extract.submit(d) for d in days] # start all
results = [f.result() for f in futures] # wait for all
Limit pressure on a shared system with a global concurrency limit: create it once, then hold a slot while touching the resource.
from prefect import task
from prefect.concurrency.sync import concurrency
@task
def load_partition(day: str) -> None:
with concurrency("warehouse", occupy=1): # waits here if five loads are already running
... # write to the warehouse
If the limit named in concurrency(...) does not exist, Prefect logs a warning and skips the acquisition, so create the limit before you rely on it.
For CPU-bound work, use a process-based task runner (prefect-dask, prefect-ray) or push the work to Spark or a warehouse rather than a thread pool.
Caching and Results¶
Caching skips a task when its inputs have not changed since a previous run.
from datetime import timedelta
from prefect import task
from prefect.cache_policies import INPUTS
@task(cache_policy=INPUTS, cache_expiration=timedelta(hours=6), persist_result=True)
def fetch_reference_data(country: str) -> dict:
... # an expensive call that is fine to reuse for six hours
persist_result=True stores the task result (locally or in remote storage) so a retry or a later run can reuse it. Persisted results can contain sensitive data, so set the storage location deliberately.
Deployments and Schedules¶
A deployment turns a flow into something Prefect can run remotely on a schedule, or on demand from the UI or API. The simplest way, useful for a single machine or container, is serve:
if __name__ == "__main__":
orders_daily.serve(name="orders-daily-6am", cron="0 6 * * *", parameters={"day": "2024-03-15"})
This starts a long-running process that creates the deployment and runs scheduled flow runs itself. For production you usually want the code pulled from Git and run on infrastructure Prefect provisions:
from prefect import flow
flow.from_source(
source="https://github.com/your-org/your-repo.git",
entrypoint="flows/orders.py:orders_daily",
).deploy(
name="orders-daily",
work_pool_name="docker-pool",
cron="0 6 * * *",
)
Or declare it in a prefect.yaml file and run prefect deploy. Schedules can be cron, interval or RRule.
Work Pools and Workers¶
A work pool is a queue of flow runs with a defined infrastructure type. A worker is a small process you run in your own environment that polls the pool and starts each run there.
| Pool type | Runs flows as |
|---|---|
process |
Subprocesses on the worker's machine |
docker |
Containers |
kubernetes |
Kubernetes jobs |
| Serverless (ECS, Cloud Run, Azure Container Instances) | Serverless containers |
This split matters for security: your code and data stay in your infrastructure. The control plane only stores metadata about runs.
Blocks, Variables and Secrets¶
| Object | Use |
|---|---|
| Variable | Non-secret configuration, such as a bucket name or a feature flag: Variable.get("bucket") |
| Secret block | An encrypted credential: Secret.load("warehouse-password").get() |
| Blocks | Saved, reusable configuration objects (a database connection, a cloud storage location) in the UI or code |
| Environment variables | Configuration read by the worker environment. Preferred for deployment-specific settings |
Never hard-code credentials in the flow or prefect.yaml. Load them from a secret block or the environment at run time.
Automations and Events¶
Prefect emits an event for every state change and can receive custom ones. Automations trigger actions when an event pattern occurs, for example:
- If a flow run of
orders-dailyends inFailed, send a Slack message. - If an expected flow has not run within 25 hours, alert (a "missing run" trigger).
- When a webhook arrives, start a deployment.
This is what makes Prefect event-driven: pipelines can start on "a file landed" or "an upstream system finished" rather than only on a clock. Notifications on failure are usually the first automation to set up.
Testing¶
Flows and tasks are ordinary functions, so test the logic without an orchestrator:
def test_transform_deduplicates():
rows = [{"order_id": 1}, {"order_id": 1}, {"order_id": 2}]
assert len(transform.fn(rows)) == 2 # .fn runs the underlying function without Prefect
def test_flow_runs():
assert orders_daily("2024-03-15") == 2 # calling the flow runs it for real
Use prefect.testing.utilities.prefect_test_harness as a fixture to run flows against a throwaway local database, so tests do not touch your real server.
Prefect vs Airflow vs Dagster¶
| Prefect | Airflow | Dagster | |
|---|---|---|---|
| Style | Python functions with decorators, dynamic control flow | Static DAG definitions (dynamic with mapping) | Declarative assets |
| Best at | Event-driven and dynamic Python workflows, quick start | Broad integrations, mature scheduling at large scale | Table and dbt centred platforms with lineage |
| Scheduling | Deployments with cron, interval, RRule, and events | Cron, datasets/assets, timetables | Schedules, sensors, automation conditions |
| Runs on | Your infrastructure via workers, control plane hosted or self-hosted | You operate scheduler, workers and metadata DB (or a managed service) | Dagster+ or self-hosted |
| Learning curve | Lowest | Highest | Medium |
See Apache Airflow and Dagster.
Common Pitfalls¶
| Pitfall | Symptom | Fix |
|---|---|---|
| Retrying a non-idempotent task | Duplicate rows after a retry | Make the load idempotent (MERGE, delete-then-insert) before adding retries |
timeout_seconds on a blocking sync task |
Task runs far past its timeout | Set timeouts on the client, or use an async task |
| Passing large data between tasks | Slow runs, high memory | Write to storage and pass a path or table name |
Forgetting .result() on futures |
Task runs still in progress when the flow ends, or wrong values | Call .result() (or wait) on every future you depend on |
| Flow works locally, fails in the work pool | Missing packages or credentials in the run environment | Bake dependencies into the image, and read secrets from the environment or blocks |
| Not running a worker | Scheduled runs sit in Late |
Start and monitor a worker for each pool |
| No failure notification | Failures found days later | Add an automation for failed and missing runs |
| Everything in one giant flow | Hard to retry or test | Split into subflows and tasks with a clear boundary |
Cheat Sheet¶
| Task | Command / Syntax |
|---|---|
| Flow / task | @flow, @task |
| Retries | @task(retries=3, retry_delay_seconds=[1, 2, 4]) |
| Run in parallel | task.submit(x), task.map(xs), future.result() |
| Local UI | prefect server start |
| Serve on a schedule | flow.serve(name="x", cron="0 6 * * *") |
| Deploy from Git | flow.from_source(url, entrypoint="f.py:flow").deploy(name=..., work_pool_name=...) |
| Create a pool | prefect work-pool create my-pool --type docker |
| Start a worker | prefect worker start --pool my-pool |
| Run a deployment | prefect deployment run "orders-daily/orders-daily" |
| Limit concurrency | prefect gcl create warehouse --limit 5, then with concurrency("warehouse"): |
| Unit test a task | my_task.fn(args) |
Interview Questions¶
Q: How is Prefect different from Airflow? A: Airflow expects a static DAG of operators that a scheduler parses. Prefect runs your Python function and observes it, so loops, conditions and dynamic fan-out are plain code. Airflow has a far bigger operator ecosystem and is entrenched at many companies. Prefect is faster to start with and better for event-driven and dynamic workflows.
Q: What are work pools and workers? A: A work pool is a queue of runs bound to an infrastructure type (process, Docker, Kubernetes, serverless). A worker is a process in your environment that polls a pool and launches the runs. Code and data stay on your infrastructure; the control plane only holds metadata and schedules.
Q: How do you make a Prefect pipeline reliable?
A: Make each task idempotent so retries are safe, add retries with backoff on the flaky calls, keep large data out of task results, alert on failed and missing runs with automations, and cover the logic with unit tests via .fn. Then use deployments with a Git source so what runs in production is the reviewed code.
Q: How would you run tasks concurrently?
A: Use .submit() or .map() to get futures on the task runner, then collect with .result(). For a shared resource such as a database, add a global concurrency limit, and for CPU-heavy work use a Dask or Ray runner or push the work to an engine built for it.
Further Reading¶
Previous: Dagster · Next: DuckDB & Polars · Back to: Index