DE Wikiconcepts / prefect

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

FeaturePrefect Server (OSS)Prefect Cloud
HostingSelf-hosted (Docker/K8s)Managed by Prefect
UIIncluded (Docker Compose)Cloud-hosted
AutomationsWebhook-basedAdvanced (Slack, PagerDuty, etc.)
Work poolsYesYes + managed workers
SLA dashboardsBasicAdvanced 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

Resources

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

FeaturePrefect Server (OSS)Prefect Cloud
HostingSelf-hosted (Docker/K8s)Managed by Prefect
UIIncluded (Docker Compose)Cloud-hosted
AutomationsWebhook-basedAdvanced (Slack, PagerDuty, etc.)
Work poolsYesYes + managed workers
SLA dashboardsBasicAdvanced 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

Resources