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 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:
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:
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():
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.
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.
dlt là thư viện Python open-source để viết ELT pipeline declarative chỉ trong vài chục dòng code. Bài viết hướng dẫn từ setup, incremental merge, schema evolution, tích hợp dbt đến deploy production trên Airflow và Dagster, với ví dụ thực tế.
So sánh LitServe, BentoML và FastAPI cho ML serving năm 2026: dynamic batching, multi-GPU, LLM streaming và cost-per-prediction từ kinh nghiệm ship 3 dự án production.
Ibis Framework Python 2026 giúp bạn viết một cú pháp DataFrame duy nhất và chạy trên 20+ backend như DuckDB, BigQuery, Snowflake, Polars. Hướng dẫn cài đặt, lazy execution, so sánh với Pandas/Polars, pipeline ETL thực tế và cách xử lý lỗi thường gặp.