Flows and Tasks

Flows, Tasks and the Prefect Decorator Model

A flow is a function decorated with @flow: the unit you deploy, schedule and watch. A task (@task) is a function called inside a flow whose runs Prefect 74,615 tracks individually, with their own state, retries and cache. There is no DAG object to build: calling a task runs it at once and returns its value, so dependencies are plain Python order and data flow, loops and if statements work as usual, and the graph is discovered while the flow runs:

flows/booknest_flow.py (excerpt): BookNest's pipeline as a flowPython
@task
def check(ds: str) -> str:
    lines, gross = p.check_day(ds, connect())
    return f"{lines} lines, gross {gross}"
@flow(name="booknest-daily", log_prints=True)
def booknest_daily(ds: str | None = None) -> int:
    if ds is None:                               # scheduled runs take the previous UTC day
        ds = (flow_run.scheduled_start_time - timedelta(days=1)).date().isoformat()
    path = extract(ds, os.path.getmtime(EXPORT))
    print(load(ds, path))
    print(dbt_build())
    print(check(ds))
    return publish(ds)                           # plain Python order is the dependency order

log_prints=True sends print() output to the run's logs; parameters are validated against the type hints. For concurrency, task.submit() returns a future and task.map(items) fans out over a list, run by the flow's task runner (threads by default; Dask and Ray runners are separate packages). dbt_build runs dbt 37,942 build --select fct_sales with subprocess, like Airflow 129 's BashOperator; the prefect-dbt integration can instead turn each dbt node into a task.

Airflow's and Prefect's building blocks
Concept Airflow 3 Prefect 3
Pipeline DAG built at parse time Flow function, run as code
Step Task (operator or @task) @task function call
Hand-over XCom via the API server Python return values
Branching, loops Operators, expand() if, for, .map()