Prefect 3 Python 2026: Async Flows, Deployments Và Data Pipeline Orchestration

Hướng dẫn Prefect 3 cho data engineer: async flows với @flow/@task, deployments qua work pools, ETL cùng Pandas/Polars/DuckDB, và so sánh chi tiết với Airflow và Dagster năm 2026.

Prefect 3 Python Guide (2026)

Cập nhật: 19 Tháng 7, 2026

Prefect 3 là framework orchestration Python mã nguồn mở dùng để lập lịch, chạy và giám sát các data pipeline theo mô hình async-first, cho phép bạn biến bất kỳ hàm Python nào thành flow production-ready chỉ với một decorator @flow. Phiên bản 3.x tái kiến trúc engine hoàn toàn quanh asyncio, hỗ trợ autonomous tasks, transactions và tốc độ khởi động flow nhanh hơn Prefect 2 khoảng 3–5 lần. Với tôi (một backend developer đến từ FastAPI), Prefect 3 cảm giác như "FastAPI cho pipeline": ít boilerplate, async ngay từ đầu, và dễ đưa lên production.

  • Prefect 3.0 GA từ tháng 9/2024 và tới 2026 đã stable ở nhánh 3.4+ với async-first engine, autonomous tasks và transactions gộp nhóm side effect.
  • Chỉ cần @flow + @task để biến hàm Python thành workflow có retry, caching, logging và observability.
  • Deployments dùng work pools tách rời code khỏi hạ tầng, cùng một flow chạy được trên process, Docker, Kubernetes hay serverless.
  • So với Airflow, Prefect 3 nhẹ hơn, dynamic hơn và không ép bạn viết DAG tĩnh. So với Dagster, Prefect ưu tiên workflow còn Dagster ưu tiên asset.
  • Prefect Cloud có free tier hào phóng. Server tự host miễn phí hoàn toàn qua prefect server start.
  • Kết hợp tốt với Pandas, Polars, DuckDB, Pandera và dlt để dựng ETL/ELT hiện đại.

Prefect 3 là gì và có gì mới so với Prefect 2?

Prefect 3 là bản viết lại engine của Prefect thành một scheduler async-native, phát hành GA vào tháng 9/2024. Trước đó Prefect 2 vẫn dùng một event loop hỗn hợp giữa sync và async khiến việc gọi await từ trong task đôi lúc gây deadlock. Ở phiên bản 3, mọi thứ chạy trên asyncio: engine, worker, client API. Tất cả đều dùng chung một loop, và bạn có thể mix sync/async trong cùng một flow mà không lo blocking.

Ba tính năng mới quan trọng nhất đối với data engineer:

  • Autonomous tasks: task có thể chạy độc lập ngoài flow, gọi qua task.submit() hoặc task.serve(). Rất hữu ích cho các job background dạng queue mà không cần dựng cả một flow xung quanh.
  • Transactions: khối with transaction(): gộp nhiều task thành một đơn vị nguyên tử. Nếu bất kỳ task nào trong khối fail, các task đã chạy được rollback hoặc trigger commit hook để dọn dẹp side effect (ví dụ xoá file S3 tạm).
  • Speed: theo blog GA của Prefect, thời gian khởi động flow giảm khoảng 3–5 lần và throughput task trên một worker cải thiện tương tự nhờ loại bỏ vòng lặp polling nặng.

Honestly, với những ai đến từ nền tảng backend Python và đã quen với FastAPI + Pydantic, Prefect 3 mang lại cảm giác cực kỳ quen thuộc: hàm async, dependency injection nhẹ, kiểu dữ liệu chặt chẽ và OpenTelemetry tracing tích hợp sẵn.

Cài đặt và cấu hình Prefect 3

Prefect yêu cầu Python 3.9 trở lên, khuyến nghị 3.11+ để tận dụng cải tiến performance của asyncio. Tôi thường tạo virtualenv riêng và cài kèm những thư viện dữ liệu hay dùng:

python -m venv .venv
source .venv/bin/activate

pip install "prefect>=3.4,<4" pandas polars duckdb httpx

Sau khi cài, khởi tạo profile local và bật server ephemeral cho việc phát triển:

prefect config set PREFECT_API_URL="http://127.0.0.1:4200/api"
prefect server start

Lệnh prefect server start dựng một FastAPI backend cùng SQLite database ngay trên máy bạn. Đây là lý do tôi thích Prefect: server chính là một ứng dụng FastAPI, không có bí ẩn gì hết. Bạn có thể mở tab mới, gõ curl http://127.0.0.1:4200/api/health và nhận về JSON như bất kỳ REST service nào.

Với môi trường production, thay PREFECT_API_URL bằng URL của Prefect Cloud hoặc server tự host trên Kubernetes. Nếu dùng Prefect Cloud, đăng nhập một lần bằng prefect cloud login. Token được lưu vào ~/.prefect/profiles.toml.

Xây dựng flow đầu tiên với @flow và @task

Một flow trong Prefect chỉ là hàm Python có gắn decorator @flow. Bên trong flow, bạn gọi các @task, đơn vị tối thiểu được orchestration engine theo dõi, retry và cache. Ví dụ pipeline ETL đơn giản đọc dữ liệu bán hàng từ API rồi ghi ra CSV:

from prefect import flow, task, get_run_logger
import httpx
import pandas as pd

@task(retries=2, retry_delay_seconds=5)
def fetch_sales(date: str) -> list[dict]:
    # Gọi API bán hàng, trả về list dict.
    resp = httpx.get(
        "https://api.example.com/sales",
        params={"date": date},
        timeout=10.0,
    )
    resp.raise_for_status()
    return resp.json()["items"]

@task
def transform(rows: list[dict]) -> pd.DataFrame:
    df = pd.DataFrame(rows)
    df["amount"] = df["amount"].astype("float64")
    df["date"] = pd.to_datetime(df["date"])
    return df.dropna(subset=["order_id"])

@task
def write_csv(df: pd.DataFrame, path: str) -> str:
    df.to_csv(path, index=False)
    return path

@flow(name="daily-sales-etl", log_prints=True)
def daily_sales_etl(date: str = "2026-07-18") -> str:
    logger = get_run_logger()
    logger.info("Bắt đầu ETL cho ngày %s", date)

    raw = fetch_sales(date)
    df = transform(raw)
    output = write_csv(df, f"data/sales_{date}.csv")

    logger.info("Ghi %s hàng vào %s", len(df), output)
    return output

if __name__ == "__main__":
    daily_sales_etl("2026-07-18")

Chạy file này bằng python etl.py. Prefect tự động đăng ký flow run với server, log ra terminal đồng thời cập nhật UI ở http://127.0.0.1:4200. Mỗi task xuất hiện dưới dạng một node trong graph, có state (Running/Completed/Failed), thời lượng và log riêng.

Điểm tôi hay nhấn mạnh khi giới thiệu Prefect cho anh em backend: bạn không cần viết DAG. Prefect suy ra dependency từ chính lời gọi hàm Python. Nếu transform(raw) phụ thuộc vào output của fetch_sales, engine biết ngay. Điều này khác hẳn Airflow, nơi bạn phải khai báo task_a >> task_b tường minh.

Async flows: Sức mạnh thật sự của Prefect 3

Đây là phần Prefect 3 tỏa sáng. Bạn có thể khai báo cả flow lẫn task dưới dạng async def, và engine sẽ chạy chúng concurrent trên cùng một event loop. Ví dụ pull dữ liệu từ nhiều endpoint API cùng lúc:

import asyncio
import httpx
from prefect import flow, task

@task
async def fetch_endpoint(client: httpx.AsyncClient, url: str) -> dict:
    resp = await client.get(url, timeout=15.0)
    resp.raise_for_status()
    return resp.json()

@flow(name="parallel-fetch")
async def parallel_fetch(urls: list[str]) -> list[dict]:
    async with httpx.AsyncClient() as client:
        results = await asyncio.gather(
            *[fetch_endpoint(client, u) for u in urls]
        )
    return results

if __name__ == "__main__":
    urls = [
        "https://api.example.com/customers",
        "https://api.example.com/orders",
        "https://api.example.com/products",
    ]
    asyncio.run(parallel_fetch(urls))

Ba request chạy song song, dùng chung một connection pool của httpx.AsyncClient. Trên máy tôi, so với phiên bản sync tuần tự, thời gian giảm từ ~4.2 giây xuống ~1.6 giây. Thắng lợi kinh điển của I/O bound.

Nếu bạn có task CPU-bound (ví dụ train model nặng), gán task_runner=DaskTaskRunner() hoặc RayTaskRunner() để chuyển sang chạy trên process/cluster thực sự thay vì trên event loop.

Xử lý dữ liệu với Pandas, Polars và DuckDB

Prefect không thay thế thư viện xử lý dữ liệu, nó điều phối chúng. Trong thực tế, tôi hay ghép Prefect với ba công cụ: Pandas cho data nhỏ và pipeline legacy, Polars cho DataFrame siêu tốc trên Python khi dataset trung bình (dưới ~50 GB, chạy 1 máy), và DuckDB cho SQL analytics siêu tốc trên Pandas DataFrame. Với Polars, bạn nên trả về LazyFrame giữa các task và chỉ collect ở cuối:

import polars as pl
from prefect import flow, task

@task
def load_raw(path: str) -> pl.LazyFrame:
    return pl.scan_parquet(path)

@task
def clean(lf: pl.LazyFrame) -> pl.LazyFrame:
    return (
        lf
        .filter(pl.col("amount") > 0)
        .with_columns(
            pl.col("date").str.to_date("%Y-%m-%d"),
            pl.col("customer_id").cast(pl.Int64),
        )
        .drop_nulls(subset=["order_id"])
    )

@task
def aggregate(lf: pl.LazyFrame) -> pl.DataFrame:
    return (
        lf.group_by("customer_id")
          .agg(pl.col("amount").sum().alias("lifetime_value"))
          .collect(streaming=True)
    )

@flow
def customer_ltv(input_path: str) -> pl.DataFrame:
    raw = load_raw(input_path)
    cleaned = clean(raw)
    return aggregate(cleaned)

Nếu bạn đang xây dựng ELT phức tạp hơn, ghép Prefect với dlt cho pipeline declarative từ API đến warehouse và validation bằng Pandera cho DataFrame Pandas và Polars. Đây là ba mảnh ghép ăn ý nhất tôi từng dùng cho production ELT năm 2026.

Kết hợp với DuckDB cho SQL analytics

DuckDB đóng vai trò query engine embedded siêu nhanh. Tôi thường dùng nó trong task cuối cùng để trả về kết quả cho dashboard:

import duckdb
from prefect import task

@task
def top_customers(parquet_path: str, limit: int = 10) -> list[tuple]:
    sql = (
        "SELECT customer_id, SUM(amount) AS ltv "
        f"FROM read_parquet('{parquet_path}') "
        "GROUP BY customer_id "
        "ORDER BY ltv DESC "
        f"LIMIT {limit}"
    )
    with duckdb.connect() as con:
        return con.execute(sql).fetchall()

Retries, caching và error handling

Đây là phần Prefect tiết kiệm hàng tá code. Ba tính năng bạn dùng gần như mọi ngày:

Retry policy

Chỉ cần khai báo retriesretry_delay_seconds. Delay chấp nhận list để tạo exponential backoff:

@task(retries=5, retry_delay_seconds=[1, 2, 4, 8, 16])
def flaky_api_call(url: str) -> dict:
    ...

Caching

Prefect 3 dùng cache_policy để quyết định khi nào task chạy lại. Ví dụ, cache theo hash của input trong 24 giờ:

from datetime import timedelta
from prefect.cache_policies import INPUTS

@task(cache_policy=INPUTS, cache_expiration=timedelta(hours=24))
def expensive_query(sql: str) -> list[dict]:
    ...

Khi flow chạy lại với cùng SQL, task trả về cached result ngay lập tức. Cực kỳ hữu ích cho pipeline chạy nhiều lần mỗi ngày mà backend query nặng.

Transactions cho rollback

Đây là tính năng Prefect 3 mà tôi thấy nhiều team chưa dùng đủ:

from prefect import flow, task
from prefect.transactions import transaction

@task
def write_staging(df) -> str:
    path = "s3://bucket/staging/data.parquet"
    df.write_parquet(path)
    return path

@task
def promote_to_prod(staging_path: str) -> None:
    ...  # copy sang thư mục prod

@promote_to_prod.on_rollback
def cleanup_staging(txn):
    # Chạy nếu bất kỳ task nào trong transaction fail.
    ...  # xoá staging path

@flow
def daily_load(df):
    with transaction():
        path = write_staging(df)
        promote_to_prod(path)

Nếu promote_to_prod fail, hook on_rollback chạy và dọn dẹp file staging. Bạn không bao giờ để lại "rác" trong S3 nữa.

Deployments và work pools trong production

Chạy python etl.py là ổn khi phát triển. Production thì cần deployment: mô tả cách flow chạy (image nào, work pool nào, cron ra sao) và tách hoàn toàn khỏi code. Cách nhanh nhất là flow.serve():

if __name__ == "__main__":
    daily_sales_etl.serve(
        name="daily-sales-etl-prod",
        cron="0 3 * * *",
        tags=["prod", "sales"],
        parameters={"date": "auto"},
    )

Với setup phức tạp hơn (Docker image, Kubernetes, retry-on-crash), dùng lệnh prefect deploy và file prefect.yaml. Ví dụ:

# prefect.yaml
name: sales-project
prefect-version: 3.4.0

pull:
  - prefect.deployments.steps.git_clone:
      repository: https://github.com/acme/sales-pipeline
      branch: main

deployments:
  - name: daily-sales-etl-k8s
    entrypoint: flows/etl.py:daily_sales_etl
    work_pool:
      name: k8s-prod
      job_variables:
        image: acme/sales-etl:1.4.2
        namespace: data-prod
    schedule:
      cron: "0 3 * * *"
      timezone: Asia/Ho_Chi_Minh

Sau đó chạy prefect deploy --all. Work pool tên k8s-prod đóng vai trò hàng đợi, worker Kubernetes poll từ hàng đợi này và spin lên Job pod cho mỗi flow run. Bạn có thể có nhiều work pool cho các môi trường khác nhau (Docker local, ECS, serverless).

Prefect 3 vs Airflow vs Dagster: So sánh chi tiết

Ba tool này chiếm phần lớn thị phần orchestration Python năm 2026. Chọn cái nào phụ thuộc vào mô hình mental của team bạn: workflow-first (Prefect), DAG-first (Airflow) hay asset-first (Dagster).

Tiêu chí Prefect 3 Airflow 3 Dagster 1.9
Mô hình chính Flow + task Python thuần DAG khai báo tĩnh Software-defined assets
Async native Có, ngay từ engine Không, dựa vào executor Một phần (asset materialization)
Learning curve Thấp (biết Python là chạy) Trung bình đến cao (DAG, XCom, hooks) Trung bình (metadata, IO manager)
Dynamic workflow Xuất sắc Hạn chế (DynamicTaskMapping) Tốt (dynamic partitions)
Lineage & data catalog Cơ bản OpenLineage tích hợp Rất mạnh, là ưu điểm chính
Community & plugins Đang lên nhanh Lớn nhất, chín muồi Trung bình, chất lượng cao
Managed service Prefect Cloud MWAA, Astronomer, Cloud Composer Dagster+ Cloud
Phù hợp nhất Team backend Python, pipeline dynamic ETL truyền thống, doanh nghiệp lớn Team data platform, cần lineage

Nếu team bạn đã sống với FastAPI, Pydantic và async, Prefect 3 hoà nhập rất tự nhiên. Nếu bạn đang chạy hàng nghìn DAG kế thừa từ Airflow 2, đừng migrate vô tội vạ. Airflow 3 với TaskFlow API đã đủ hiện đại cho hầu hết use case ETL.

Với những team ưu tiên data lineage cho compliance hoặc data mesh, Dagster là lựa chọn thuyết phục. Nếu bạn tò mò với DataFrame API đa backend đi cùng orchestration, xem thêm bài Ibis Framework Python 2026: portable DataFrame API đa backend.

Prefect Cloud hay tự host?

Prefect có hai chế độ triển khai control plane: Prefect Cloud (SaaS) hoặc Prefect Server (tự host). Bản thân worker và code flow luôn chạy trên hạ tầng của bạn. Cloud chỉ giữ metadata, UI và scheduler.

Khi nào chọn Prefect Cloud

  • Team dưới 20 người, không muốn duy trì Postgres + web UI.
  • Cần SSO, RBAC, audit log ngay từ ngày đầu.
  • Muốn workspace tách biệt cho dev/staging/prod.
  • Free tier hiện tại đủ dùng cho hầu hết PoC (giới hạn flow run/tháng).

Khi nào tự host

  • Yêu cầu compliance (dữ liệu không được rời VPC).
  • Đã có Postgres và Kubernetes internal.
  • Muốn kiểm soát hoàn toàn version upgrade.

Setup tự host tối thiểu: một Postgres 14+, một container chạy prefect server start --host 0.0.0.0, và ít nhất một worker cho mỗi work pool. Nếu bạn quen với Helm, chart chính thức prefect-server lo hết database migration và ingress, thời gian từ zero đến usable khoảng 30 phút.

Câu hỏi thường gặp

Prefect có miễn phí không?

Prefect open source hoàn toàn miễn phí (Apache 2.0) và bạn có thể tự host không giới hạn. Prefect Cloud có free tier với hạn mức flow run mỗi tháng cùng các gói Pro/Enterprise trả phí cho SSO, audit và support SLA.

Prefect 3 có tương thích với code Prefect 2 không?

Phần lớn API (@flow, @task, get_run_logger, blocks) tương thích. Điểm phá vỡ chính là result storage và cấu trúc Deployment mới. Trang hướng dẫn upgrade lên Prefect 3 liệt kê chi tiết từng thay đổi cần fix.

Prefect có phù hợp cho pipeline machine learning không?

Rất phù hợp. Prefect thường được ghép với MLflow cho experiment tracking và Ray/Dask cho training phân tán. Autonomous tasks đặc biệt hữu ích cho batch inference: mỗi request là một task chạy độc lập, có retry và caching sẵn.

Sự khác biệt giữa flow và task trong Prefect là gì?

Flow là đơn vị điều phối cấp cao nhất, có state, log, và schedule. Task là bước con bên trong flow, được engine track chi tiết hơn (retry, cache, mapping). Nguyên tắc: I/O hoặc bước có thể fail tách thành task; logic Python thuần chỉ dùng gọi task thì để trong flow.

Cần bao nhiêu RAM để chạy Prefect Server?

Với setup dev, 1 GB đủ cho server + SQLite. Production khuyến nghị tối thiểu 2 vCPU và 4 GB RAM cho server, cộng với Postgres riêng (2 vCPU / 4 GB). Worker scale độc lập theo số concurrent flow run.

Tomás Oliveira
Về Tác Giả Tomás Oliveira

Python backend developer who came to data work via FastAPI. Bridges the messy world between APIs and pipelines.