בניית צינורות ETL בפייתון עם Prefect 3: מדריך מעשי לשנת 2026

מדריך מעשי לבניית צינור ETL אמין בפייתון עם Prefect 3.4: פונקציות עם @task ו-@flow, retries, caching, deployments, בדיקות pytest, ו-footguns של פרודקשן לשנת 2026.

Prefect 3 ETL בפייתון: מדריך 2026

עודכן: 20 ביולי 2026

כדי לבנות צינור ETL אמין בפייתון עם Prefect 3 מגדירים פונקציות פייתון רגילות ומוסיפים את הדקורטורים @task ו-@flow, מפעילים את הזרימה מקומית, ואז יוצרים Deployment עם לוח זמנים ו-work pool להרצה בפרודקשן. Prefect 3, שהגיע ל-GA בספטמבר 2024 ומעודכן לגרסה 3.4 בשנת 2026, מספק retries אוטומטיים, caching, transactions ומעקב מלא, בלי צורך ב-DAGs סטטיים כמו ב-Airflow. במדריך הזה נבנה pipeline מלא של extract-transform-load עם pandas 3.0, נכתוב בדיקות שרצות בכל pull request, ונדבר על ה-footguns של הרצה בפרודקשן.

  • Prefect 3 מבוסס על פונקציות פייתון רגילות. אין DSL, אין קבצי YAML, ואין DAG סטטי. כל @flow הוא פונקציה שאפשר להריץ ולדבג כמו כל קוד אחר.
  • Retries, timeout, caching ו-concurrency limits מוגדרים כפרמטרים בדקורטור (לא כתשתית חיצונית). זה מוזיל את עלות התחזוקה מול Airflow.
  • לצורך ריצה בפרודקשן צריך שני רכיבים: Deployment (הגדרת ה-flow ולוח הזמנים) ו-Work Pool עם Worker שמושך משימות ומריץ אותן.
  • בדיקות pipeline אינן בונוס. הן חובה. ב-Prefect 3 קל להריץ flows בתוך pytest עם fixtures שמזרימות נתוני test דטרמיניסטיים.
  • אינטגרציה עם dbt-core 1.9, DuckDB ו-Great Expectations 1.x הופכת את Prefect 3 לתחליף רציני ל-Airflow עבור צוותים קטנים ובינוניים.

מה זה Prefect 3 ולמה להשתמש בו ל-ETL?

אז בואו נצלול לזה. Prefect 3 הוא מנוע אורקסטרציה של workflows בפייתון, בפועל שכבה דקה שהופכת פונקציות פייתון רגילות לצינורות ניתנים לתזמון, ניתנים למעקב וניתנים לניסיון-חוזר. במקום להגדיר DAG סטטי כמו ב-Airflow, אתם כותבים פונקציה, שמים מעליה @flow, ו-Prefect דואג לכל השאר: retries, לוגים, מטריקות, מצב ריצה, ומעקב וויזואלי דרך ה-UI.

בעולם ה-ETL הקלאסי (שאיפה מ-API, טרנספורמציה עם pandas, טעינה ל-warehouse), היתרונות המעשיים של Prefect 3 מורגשים מיד. אחרי שהעברתי בעצמי צוות משלושה pipelines של Airflow לזרימות Prefect, מה שבלט הכי הרבה זה שהזמן הממוצע לפתרון תקלה ירד בערך בחצי, פשוט מפני שהקוד רץ מקומית באותה צורה בדיוק שבה הוא רץ בפרודקשן. אין SubDagOperator, אין XCom, ואין הפתעות בגלל serialization.

עוד נקודה שחשוב לציין. Prefect 3 בנוי סביב transactions. אם step של טרנספורמציה נכשל אחרי שכבר הייתה כתיבה חלקית ל-staging, אפשר להגדיר on_rollback שינקה אחריו, משהו שב-Airflow דרש לרוב סקריפט bash חיצוני או task נפרד.

Prefect 3 מול Airflow ו-Dagster: השוואה

לפני שקופצים לקוד, שווה להבין למה לבחור ב-Prefect 3 ולא ב-Airflow או Dagster. שלושתם כלים לגיטימיים, אבל הם מכוונים לצוותים שונים ולתרחישים שונים. הטבלה למטה מסכמת את ההשוואה שאני בדרך כלל עורכת מול צוותי דאטה שוקלים מעבר בשנת 2026.

פיצ׳רPrefect 3Airflow 3Dagster 1.9
מודל תכנותפונקציות פייתון עם דקורטוריםDAG אימפרטיבי + OperatorsAssets מבוססי-נתונים (data-aware)
עקומת למידהנמוכה. אם אתם כותבים פייתון, אתם מוכניםבינונית-גבוהה. צריך ללמוד Airflow כמושגגבוהה. פרדיגמת asset שונה
הרצה מקומיתמיידית: python flow.pyדורש scheduler + webserverדורש dagster dev
Retries ו-cachingברמת הדקורטור (מובנה)ברמת ה-Operatorברמת ה-op ו-asset
Data lineage מובנהחלקי (Artifacts)לאכן, הפיצ׳ר הדגל
עלות הרצה עצמיתנמוכה (server + workers)גבוהה (scheduler, webserver, executor, DB)בינונית
מתאים לצוות של…3–20 מהנדסים, ETL קלאסיארגונים גדולים, DAGs מורכביםצוותים שמתמקדים ב-data assets

הכלל שאני עובדת לפיו: אם המשימה היא בעיקר לגרד API, לעבד עם pandas או Polars, ולטעון ל-warehouse, אז Prefect 3 יחסוך לכם שעות רבות של setup. אם יש לכם 400 DAGs היסטוריים ב-Airflow ופיצ׳רים כמו KubernetesExecutor בייצור, מעבר לא ישתלם. Dagster מבריק כשלב הבא של תפיסת data-mesh, אבל דורש שינוי תפיסתי גדול.

התקנה והרצת ה-flow הראשון

ההתקנה של Prefect 3 היא pip install אחד. אני ממליצה על סביבת virtualenv נפרדת לכל pipeline (לא בגלל Prefect, אלא כי צינור ETL בייצור נוגע בהמון ספריות כמו pandas, requests, sqlalchemy, boto3 שנוטות להתעדכן בקצב שונה).

python -m venv .venv
source .venv/bin/activate
pip install "prefect==3.4.*" "pandas==3.0.*" "requests" "sqlalchemy" "duckdb"
prefect --version
# 3.4.7

אחרי ההתקנה, יוצרים flow מינימלי. שימו לב איך הכל הוא פונקציית פייתון רגילה. אין מחלקות, אין metaclass magic, ובעיקר אין קובץ YAML.

from prefect import flow, task
from prefect.logging import get_run_logger

@task(retries=3, retry_delay_seconds=10)
def greet(name: str) -> str:
    logger = get_run_logger()
    logger.info(f"Greeting {name}")
    return f"Hello, {name}"

@flow(name="hello-etl")
def hello_flow(name: str = "world"):
    return greet(name)

if __name__ == "__main__":
    print(hello_flow("data-eng"))

מריצים עם python hello_flow.py. המסך יראה לוג של Prefect, וב-~/.prefect יתעדכן מצב הריצה. כדי לפתוח את ה-UI המקומי:

prefect server start
# פתחו את http://127.0.0.1:4200 בדפדפן

ב-UI תראו את ה-flow שהרצתם, את משך הזמן של כל task, ואת ה-retries אם היו. זה כבר יותר תצפית ממה שהיה לי בפרויקטים שלמים של Airflow לפני שהוספתי StatsD ידנית.

בניית צינור ETL מלא עם pandas 3.0

עכשיו לחלק המעניין: pipeline אמיתי. נבנה זרימה שמוציאה נתוני מכירות מ-API ציבורי, מנקה אותם עם pandas 3.0, ומטעינה ל-DuckDB, פלטפורמה שאני מעדיפה עבור אנליטיקה מקומית עד גדלים של עשרות ג׳יגה. אם עדיין לא הכרתם את היתרונות של DuckDB לפייפליינים, כדאי לקרוא את המדריך שלנו על Polars מול pandas בפייתון. DuckDB משתלב עם שניהם באופן חלק.

from __future__ import annotations
import pandas as pd
import requests
import duckdb
from pathlib import Path
from prefect import flow, task
from prefect.logging import get_run_logger

API_URL = "https://api.example.com/sales/2026"
DB_PATH = Path("./warehouse.duckdb")

@task(retries=3, retry_delay_seconds=[5, 30, 120], timeout_seconds=60)
def extract_sales(since: str) -> pd.DataFrame:
    logger = get_run_logger()
    logger.info(f"Extracting sales since {since}")
    resp = requests.get(API_URL, params={"since": since}, timeout=30)
    resp.raise_for_status()
    df = pd.DataFrame(resp.json()["records"])
    logger.info(f"Extracted {len(df)} rows")
    return df

@task
def transform_sales(df: pd.DataFrame) -> pd.DataFrame:
    logger = get_run_logger()
    if df.empty:
        logger.warning("No rows to transform")
        return df
    df = df.copy()
    df["order_date"] = pd.to_datetime(df["order_date"], utc=True)
    df["revenue_ils"] = df["revenue_usd"].astype("float64") * 3.65
    df = df.dropna(subset=["customer_id", "order_id"])
    df = df.drop_duplicates(subset=["order_id"])
    logger.info(f"Transformed to {len(df)} clean rows")
    return df

@task
def load_to_duckdb(df: pd.DataFrame, table: str = "sales_2026") -> int:
    logger = get_run_logger()
    with duckdb.connect(str(DB_PATH)) as con:
        con.execute(f"CREATE TABLE IF NOT EXISTS {table} AS SELECT * FROM df LIMIT 0")
        con.register("staging", df)
        con.execute(f"INSERT INTO {table} SELECT * FROM staging")
        count = con.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0]
    logger.info(f"Loaded {len(df)} rows; table total = {count}")
    return count

@flow(name="sales-etl", log_prints=True)
def sales_etl(since: str = "2026-01-01") -> int:
    raw = extract_sales(since)
    clean = transform_sales(raw)
    total = load_to_duckdb(clean)
    return total

if __name__ == "__main__":
    sales_etl("2026-07-01")

מספר דברים חשובים לשים לב אליהם. retry_delay_seconds מקבל רשימה עם exponential backoff (5, 30, 120 שניות), משהו שבנייה ידנית ב-Airflow דרשה custom operator שלם. ה-timeout_seconds על ה-extract מגן מפני API שנתקע (וקורה לעיתים יותר קרובות ממה שהיינו רוצים). ואחרי שהעברתי pipeline דומה ל-צינור ניקוי נתונים עם pandas 3.0, אני שמה drop_duplicates על ה-order_id תמיד. API שמחזיר את אותה שורה בשני עמודים שונים זה scenario אמיתי, לא תיאורטי.

איך בודקים flows של Prefect?

לבדוק pipeline ETL זו לא בונוס, זו חובה. הכלל אצלי בצוות: אם אין pytest שרץ בכל pull request, ה-pipeline לא נכנס לפרודקשן. הסיבה פשוטה. Pipelines נכשלים בשלוש בלילה, וכשמישהו מתעורר לפתור, הוא לא רוצה לגלות שה-schema של ה-API השתנה שבועיים לפני זה בלי שאף אחד שם לב.

Prefect 3 מאפשר להריץ flow בתוך pytest כמו כל פונקציה. ה-fixture הסטנדרטית שלי מזרימה DataFrame דטרמיניסטי במקום לקרוא ל-API אמיתי, בעזרת monkeypatch:

import pandas as pd
import pytest
from sales_pipeline import sales_etl, extract_sales, transform_sales

@pytest.fixture
def sample_sales() -> pd.DataFrame:
    return pd.DataFrame({
        "order_id": ["o1", "o2", "o2", "o3"],  # duplicate on purpose
        "customer_id": ["c1", "c2", "c2", None],  # null on purpose
        "order_date": ["2026-07-01", "2026-07-02", "2026-07-02", "2026-07-03"],
        "revenue_usd": [100.0, 250.0, 250.0, 40.0],
    })

def test_transform_removes_dupes_and_nulls(sample_sales):
    out = transform_sales.fn(sample_sales)  # .fn runs the function without Prefect
    assert len(out) == 2
    assert "revenue_ils" in out.columns
    assert out["revenue_ils"].iloc[0] == pytest.approx(365.0)

def test_flow_end_to_end(monkeypatch, sample_sales, tmp_path):
    monkeypatch.setattr(
        "sales_pipeline.extract_sales",
        lambda since: sample_sales,
    )
    monkeypatch.setattr("sales_pipeline.DB_PATH", tmp_path / "test.duckdb")
    total = sales_etl(since="2026-07-01")
    assert total == 2

שני הטריקים הכי שימושיים כאן: task.fn קורא לפונקציה המקורית בלי ה-wrapper של Prefect, כך שהבדיקה לא צריכה server פעיל. וה-monkeypatch על DB_PATH מבודד את בדיקת ה-load ל-DuckDB זמני. אחרי שנים של pipelines שנשברו בפרודקשן, הפכתי את הכלל של לפחות בדיקה אחת per task לחובה בסקירת הקוד.

Deployments, schedules ו-work pools

הרצת flow מקומית זה נחמד, אבל לא ETL אמיתי עד שיש schedule ו-worker שרצים ללא התערבות ידנית. ב-Prefect 3, שלושת הרכיבים ש-חייבים להבין הם: Deployment (איך ה-flow ירוץ), Work Pool (סוג התשתית: process, docker, k8s), ו-Worker (התהליך שמושך משימות מה-pool ומריץ אותן).

הדרך המומלצת לגרסה 3.4 היא flow.serve() לפיתוח, או flow.deploy() לפרודקשן. הנה איך הופכים את ה-sales_etl שלמעלה ל-deployment שרץ כל שעה:

from sales_pipeline import sales_etl

if __name__ == "__main__":
    sales_etl.serve(
        name="sales-etl-hourly",
        cron="0 * * * *",             # every full hour
        tags=["etl", "sales", "prod"],
        description="Hourly sales ETL from external API to DuckDB",
        parameters={"since": "2026-01-01"},
    )

הפקודה serve() יוצרת בו-זמנית deployment, לוח זמנים, ו-worker פנימי. לפרודקשן אמיתי (עם k8s, ECS, או שרת EC2 נפרד) משתמשים ב-deploy() עם work pool נפרד:

prefect work-pool create --type process my-process-pool
prefect worker start --pool my-process-pool
sales_etl.deploy(
    name="sales-etl-prod",
    work_pool_name="my-process-pool",
    cron="0 * * * *",
    image="ghcr.io/mycompany/sales-etl:2026.07",
    push=True,
)

ההפרדה בין ה-deployment ל-worker מאפשרת לשמור על principle חשוב: ה-Prefect server לא צריך לגשת ל-warehouse. ה-worker רץ בתוך ה-VPC שלכם עם הרשאות ל-DuckDB / Snowflake / BigQuery, וה-server (בין אם Prefect Cloud או self-hosted) רק שולח לו הודעות מה להריץ. זה גם הופך את הבידוד של סודות (secrets) לפשוט משמעותית.

Retries, caching ותצפית (observability)

אחד הפיצ׳רים החזקים ביותר של Prefect 3 הוא ה-caching. אם ה-extract משך 100k רשומות ואז ה-transform נכשל בגלל bug, אתם לא רוצים למשוך שוב מה-API. עם cache_key_fn ו-cache_expiration, Prefect שומר את הפלט של task ומדלג עליו ב-run הבא אם הפרמטרים זהים.

from datetime import timedelta
from prefect.tasks import task_input_hash

@task(
    retries=3,
    retry_delay_seconds=[5, 30, 120],
    cache_key_fn=task_input_hash,
    cache_expiration=timedelta(hours=1),
)
def extract_sales(since: str) -> pd.DataFrame:
    ...

אם קראתם ל-extract_sales("2026-07-01") פעמיים באותה שעה, הקריאה השנייה תחזיר את הפלט המקורי בלי לגשת ל-API. זה מציל pipelines בפרודקשן, במיוחד כשמפעילים אותם ידנית לצורך debug אחרי כישלון.

לצד caching, יש שלושה מנגנוני תצפית שאני מגדירה בכל pipeline לפני שהוא יוצא לפרודקשן:

  • Artifacts: כתיבת סיכום markdown לסיום כל flow עם ספירות שורות, זמנים ולינק ל-dashboard. פשוט create_markdown_artifact(...).
  • State hooks: on_failure, on_completion, on_cancellation שרצים כשה-flow מסיים במצב מסוים. אני מחברת on_failure ל-Slack webhook כדי לקבל התראה אמיתית ולא רק אימייל שאיש לא קורא.
  • Automations: כללים ב-UI שמפעילים אירוע (למשל: אם flow נכשל 3 פעמים ברצף, השהה את ה-schedule).

עבור צוותים שמריצים גם מודלי ML, כדאי לחבר את ה-Prefect flows ל-pipeline המודל. לדוגמה, אם ה-flow של ETL מסיים בהצלחה, אפשר להפעיל טריגר על flow אחר שמעדכן את המודל שנמצא ב-production. ראו את המדריך שלנו על פריסת מודלי scikit-learn עם FastAPI ו-ONNX להבנת הצד השני של ה-pipeline הזה.

Footguns של הרצה בפרודקשן

כמה מלכודות שנתקלתי בהן בפרויקטים אמיתיים, וכדאי לדעת עליהן לפני שאתם מכניסים Prefect 3 ל-production. הכי כואב לי היה על pipeline של לקוח שקרס בערב שישי בגלל time zone. מאז אני שומרת רשימה קצרה שאני מריצה בכל code review לפני שקוד עולה.

1. Concurrency ברמת ה-worker לא מגן על ה-warehouse

אם ה-load שלכם כותב ל-DuckDB או SQLite מקומי, שני workers שרצים במקביל יגרמו ל-lock. השתמשו ב-concurrency limits של Prefect (הגדרה גלובלית לפי tag): prefect concurrency-limit create sales-etl 1, ואז הוסיפו tags=["sales-etl"] ל-flow.

2. Time zones הם עדיין בעיה

ה-cron="0 * * * *" ב-Prefect עובד לפי UTC כברירת מחדל. אם ה-schedule שלכם הוא "כל יום ב-3 בבוקר שעון ישראל", השתמשו במפורש ב-timezone="Asia/Jerusalem". שווה גם לקרוא את המסמכים הרשמיים של Prefect על schedules.

3. dbt-core משתלב, אבל דורש worker נכון

אם אתם משלבים dbt-core 1.9 בתוך flow של Prefect (וזה שילוב אלגנטי מאוד), ודאו שה-worker רץ בסביבה עם ה-profiles.yml הנכון וה-service account המתאים. שגיאה נפוצה: dbt רץ טוב מקומית וקורס בפרודקשן כי משתני הסביבה של ה-warehouse חסרים ב-worker container.

4. Backfills דורשים תכנון מראש

Prefect 3 לא מספק backfill מובנה כמו Airflow. הדפוס המומלץ הוא לבנות את ה-flow כך שיקבל פרמטר run_date ואז לכתוב סקריפט ניפרד שקורא לו בלולאה: for d in date_range(...): sales_etl.deploy_run(parameters={"run_date": d}). עצה מהחוויה: הוסיפו idempotency ברמת ה-warehouse (DELETE WHERE run_date = ? לפני INSERT). אחרת backfill יכול לייצר כפילויות.

5. Prefect Cloud vs self-hosted

ההחלטה בין Prefect Cloud (SaaS) ל-self-hosted היא בעיקר operational. ה-server היא שכבה דקה, אבל אני ראיתי צוותים משקיעים שבועות בתחזוקת PostgreSQL, Redis ו-workers רק בשביל להימנע מ-$400 בחודש ל-Cloud. אם אתם צוות של פחות מ-10 מהנדסים, ה-ROI של Cloud כמעט תמיד חיובי. את ההשוואה המלאה של המחיר תמצאו ב-דף המחירים הרשמי של Prefect.

שאלות נפוצות

האם Prefect 3 טוב יותר מ-Airflow?

לצוותים קטנים ובינוניים שכותבים ETL קלאסי בפייתון, כן, לרוב. ה-setup פשוט יותר, ההרצה המקומית זהה לפרודקשן, ופיצ׳רים כמו retries ו-caching מובנים ברמת הדקורטור. לצוותים גדולים עם מאות DAGs קיימים ב-Airflow, מעבר לא ישתלם בטווח קצר. Prefect 3 גם צריך פחות תשתית תפעולית, אין scheduler, webserver ו-executor נפרדים.

מה זה Prefect Deployment ולמה צריך אותו?

Deployment ב-Prefect 3 הוא רישום של flow ב-server יחד עם קונפיגורציה של איך ומתי הוא ירוץ: שם, לוח זמנים, פרמטרים ברירת-מחדל, work pool ותגיות. בלי Deployment, ה-flow יכול לרוץ רק כשמפעילים אותו ידנית מקומית. עם Deployment, ה-scheduler של Prefect שולח משימות ל-worker אוטומטית לפי ה-cron/interval שהוגדר.

איך Prefect מטפל ב-retries אוטומטית?

מגדירים retries=N ו-retry_delay_seconds=X (או רשימה עבור exponential backoff) בדקורטור @task. אם ה-task מעלה exception, Prefect ינסה שוב עד N פעמים עם השהיה בין ניסיונות. את סוג ה-exception אפשר לסנן עם retry_condition_fn כך שרק שגיאות רשת יגרמו ל-retry, לא שגיאות סכימה.

האם Prefect 3 מתאים ל-production אמיתי?

כן. Prefect 3 הגיע ל-GA בספטמבר 2024 והוא בשימוש בייצור בחברות כמו Cash App, Uber ו-Slack. הכלים ליציבות פרודקשן (retries, caching, timeouts, state hooks, automations, ו-work pools מבודדים) כולם חלק מהליבה. הנקודה שכדאי לוודא לפני production היא backup ל-Prefect DB (אם self-hosted) ו-monitoring על ה-workers עצמם.

איך משתמשים ב-dbt יחד עם Prefect 3?

מתקינים את prefect-dbt ומשתמשים ב-DbtCoreOperation או ב-tasks ייעודיים כמו trigger_dbt_cli_command בתוך flow. הדפוס הנפוץ הוא flow אחד שמורכב מ-extract וטעינה ל-staging עם pandas, ואחריו task שמריץ dbt run ואז dbt test. אם dbt tests נכשלים, ה-flow עוצר וההתראה נשלחת אוטומטית.

Hannah Walsh
אודות הכותב Hannah Walsh

Data engineer making sure the pipelines feeding the models don't silently break at 3am. Big fan of dbt and bigger fan of testing.