Prefect
Prefect is a modern workflow orchestration framework that makes it easy to build, run, and monitor data pipelines. Unlike Airflow's DAG-based approach, Prefect emphasizes Python-native flows with automatic retries, caching, and a cloud observability platform.
Core Concepts
Flows and Tasks
In Prefect, a flow is a Python function that defines a workflow. A task is a discrete unit of work within a flow. Tasks are decorated with @task and flows with@flow.
from prefect import flow, task
from prefect.task_runners import ConcurrentTaskRunner
import httpx
@task(retries=3, retry_delay_seconds=60)
def fetch_weather(city: str) -> dict:
response = httpx.get(
f"https://api.weather.gov/points/{city}",
headers={"User-Agent": "my-pipeline/1.0"},
)
response.raise_for_status()
return response.json()
@task
def transform(data: dict) -> dict:
properties = data.get("properties", {})
return {
"city": properties.get("relativeLocation", {}).get("properties", {}).get("city"),
"forecast": properties.get("forecast"),
}
@task(cache_key_fn=lambda ctx: "weather-aggregation")
def aggregate(results: list[dict]) -> dict:
return {"cities": results, "count": len(results)}
@flow(log_prints=True)
def weather_pipeline(cities: list[str]):
results = []
for city in cities:
raw = fetch_weather(city)
transformed = transform(raw)
results.append(transformed)
return aggregate(results)
if __name__ == "__main__":
weather_pipeline(["37.7749,-122.4194", "40.7128,-74.0060"])Automatic Retries
Prefect's retry mechanism is built-in and configurable at the task level. You specify how many times a task should retry and how long to wait between attempts. Prefect tracks retry history in its metadata store.
Caching
Tasks can be cached using a cache key function. If the cache key matches a previous run, the cached result is returned without executing the task. This is useful for expensive or idempotent operations.
Prefect Server vs Prefect Cloud
| Feature | Prefect Server (OSS) | Prefect Cloud |
|---|---|---|
| Hosting | Self-hosted (Docker/K8s) | Managed by Prefect |
| UI | Included (Docker Compose) | Cloud-hosted |
| Automations | Webhook-based | Advanced (Slack, PagerDuty, etc.) |
| Work pools | Yes | Yes + managed workers |
| SLA dashboards | Basic | Advanced analytics |
Work Pools and Workers
Prefect 2 introduced work pools and workers as the execution model. A work pool defines how flow runs are submitted (e.g., to a Kubernetes cluster, a local process, or a serverless function). A worker polls the work pool and executes runs.
# Start a local worker for the "default" work pool prefect worker start --pool "default" # Or use Docker prefect worker start --pool "docker" --type "docker"Deployments and Schedules
A deployment packages a flow with its configuration and makes it available for scheduling or API-triggered runs.
from prefect.deployments import Deployment
from prefect.filesystems import GitHub
weather_pipeline_deploy = Deployment.build_from_flow(
flow=weather_pipeline,
name="hourly-weather",
schedule={"cron": "0 * * * *", "timezone": "America/New_York"},
storage=GitHub(repository="kiersten/data-pipelines"),
work_queue_name="default",
)
weather_pipeline_deploy.apply()Best Practices
- Use flow-level logging with
log_prints=True - Set retries for tasks with network calls or external dependencies
- Use caching for expensive aggregations that don't change often
- Keep tasks small and focused - one task, one responsibility
- Define typed parameters for flow inputs