ETL Pipelines στην Python: Πλήρης Οδηγός με Airflow, Prefect και Dagster (2026)

Πλήρης πρακτικός οδηγός για ETL pipelines στην Python το 2026: αναλυτική σύγκριση Apache Airflow 3.0, Prefect 3 και Dagster με παραδείγματα κώδικα, data quality checks και best practices production.

Οδηγός ETL Python: Airflow, Prefect (2026)

Ενημερώθηκε: 13 Ιουλίου 2026

Ένα ETL pipeline στην Python είναι μια αυτοματοποιημένη ροή εργασιών που εξάγει (Extract) δεδομένα από πηγές (APIs, βάσεις δεδομένων, αρχεία), τα μετασχηματίζει (Transform) σύμφωνα με τους επιχειρηματικούς κανόνες και τα φορτώνει (Load) σε μια αποθήκη δεδομένων ή άλλο προορισμό. Το 2026, οι τρεις κυρίαρχες βιβλιοθήκες που καθορίζουν πώς φτιάχνουμε ETL pipelines είναι το Apache Airflow 3.0, το Prefect 3 και το Dagster. Έχω στήσει pipelines και με τα τρία σε production, οπότε θα δούμε πώς λειτουργεί το καθένα, ποιο ταιριάζει σε ποια περίπτωση χρήσης, και θα γράψουμε παραγωγικό κώδικα για κάθε ένα.

  • Το Apache Airflow 3.0 (κυκλοφορία Απρίλιος 2025) είναι ο βετεράνος του orchestration με τεράστια κοινότητα και ώριμο οικοσύστημα providers, ιδανικό για εταιρικά, batch-oriented workloads.
  • Το Prefect 3 προσφέρει καθαρή Pythonic σύνταξη με decorators @flow και @task, είναι φιλικό για data scientists και εξαιρετικό για ταχεία ανάπτυξη.
  • Το Dagster εισάγει τη φιλοσοφία των software-defined assets: το pipeline περιγράφει τα δεδομένα ως προϊόντα, όχι απλώς τις εργασίες που τα παράγουν.
  • Για μικρές ροές ETL, η pandas ή η Polars με ένα απλό script και cron μπορεί να αρκεί. Δεν χρειάζεσαι orchestrator για κάθε δουλειά.
  • Για production, οι data quality checks με Great Expectations ή Pandera είναι εξίσου κρίσιμοι με την ίδια τη ροή ETL.
  • Οι σύγχρονες ροές χρησιμοποιούν PyArrow και Parquet ως ενδιάμεση αναπαράσταση για ταχύτητα και συμβατότητα ανάμεσα σε εργαλεία.

Τι είναι ETL Pipeline και γιατί το χρειάζεσαι

Το ETL (ακρωνύμιο του Extract, Transform, Load) είναι το θεμελιώδες μοτίβο για τη μετακίνηση δεδομένων ανάμεσα σε συστήματα. Στο στάδιο Extract ανακτούμε ακατέργαστα δεδομένα από πηγές: REST APIs, PostgreSQL, MongoDB, S3 buckets, CSV αρχεία, event streams όπως Kafka. Στο στάδιο Transform καθαρίζουμε, φιλτράρουμε, εμπλουτίζουμε ή συγκεντρώνουμε τα δεδομένα. Εδώ ζει η επιχειρηματική λογική. Στο Load γράφουμε το τελικό dataset σε ένα data warehouse (Snowflake, BigQuery, Redshift), σε ένα data lake ή σε μια λειτουργική βάση δεδομένων.

Η Python κυριαρχεί στον χώρο των ETL pipelines γιατί έχει το πιο ώριμο οικοσύστημα εργαλείων: βιβλιοθήκες για κάθε πηγή δεδομένων, ισχυρές μηχανές μετασχηματισμού όπως pandas και Polars, και τρεις εξαιρετικούς orchestrators. Στην πράξη, ένα production ETL pipeline χρειάζεται πολύ περισσότερα από τρία βήματα: retry logic όταν αποτυγχάνει μια κλήση API, χρονοπρογραμματισμό, παρακολούθηση, alerting, έλεγχους ποιότητας, εξαρτήσεις μεταξύ εργασιών και ιστορικό εκτελέσεων. Εδώ ακριβώς μπαίνουν οι orchestrators όπως το Airflow, το Prefect και το Dagster.

Πριν προχωρήσουμε στον κώδικα, αξίζει να αναφέρουμε το σύγχρονο δίδυμο ELT (Extract, Load, Transform), όπου τα δεδομένα φορτώνονται πρώτα στην αποθήκη και ο μετασχηματισμός γίνεται εκεί (συνήθως με SQL μέσω dbt). Οι Python orchestrators χρησιμοποιούνται εξίσου και για τις δύο προσεγγίσεις.

Πώς φτιάχνεις ένα βασικό ETL pipeline με Python και Pandas

Πριν επενδύσουμε σε έναν orchestrator, ας δούμε την πιο απλή μορφή ενός ETL pipeline: ένα plain Python script που χρησιμοποιεί την pandas. Αυτή η προσέγγιση είναι απόλυτα βιώσιμη για μικρά, καθημερινά jobs που τρέχουν σε cron. Το παρακάτω παράδειγμα εξάγει δεδομένα πωλήσεων από ένα CSV, τα καθαρίζει και τα φορτώνει σε μια SQLite βάση.

import pandas as pd
from sqlalchemy import create_engine
from pathlib import Path
import logging

logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
log = logging.getLogger("etl")

def extract(source: Path) -> pd.DataFrame:
    log.info("Extract: reading %s", source)
    df = pd.read_csv(source, parse_dates=["order_date"])
    log.info("Extracted %d rows", len(df))
    return df

def transform(df: pd.DataFrame) -> pd.DataFrame:
    log.info("Transform: cleaning and enriching")
    # Καθαρισμός: αφαίρεση διπλότυπων και NaN σε κρίσιμες στήλες
    df = df.drop_duplicates(subset=["order_id"])
    df = df.dropna(subset=["order_id", "customer_id", "total"])

    # Μετασχηματισμοί: κατηγοριοποίηση αξίας παραγγελίας
    df["value_bucket"] = pd.cut(
        df["total"],
        bins=[0, 50, 200, 1000, float("inf")],
        labels=["low", "mid", "high", "premium"],
    )
    df["order_month"] = df["order_date"].dt.to_period("M").astype(str)
    return df

def load(df: pd.DataFrame, db_url: str, table: str) -> None:
    log.info("Load: writing %d rows to %s", len(df), table)
    engine = create_engine(db_url)
    df.to_sql(table, engine, if_exists="replace", index=False)

if __name__ == "__main__":
    raw = extract(Path("data/orders.csv"))
    clean = transform(raw)
    load(clean, "sqlite:///data/warehouse.db", "orders_daily")

Το script δουλεύει, αλλά αν το βάλεις σε production θα σου λείψουν πολύ σύντομα κρίσιμα χαρακτηριστικά. Τι γίνεται αν αποτύχει το extract βήμα λόγω δικτύου; Ποιος θα σε ειδοποιήσει; Πώς θα κάνεις retry μόνο το βήμα που απέτυχε, χωρίς να ξαναδιαβάσεις όλο το CSV; (Έχω φάει το κεφάλι μου με αυτό ακριβώς το σενάριο σε μια 3π.μ. εκτέλεση.) Αν σε ενδιαφέρει η μηχανή μετασχηματισμού καθαυτή, δες τον αναλυτικό μας οδηγό καθαρισμού δεδομένων στην Python για κάθε σκέψιμη τεχνική με pandas και Scikit-Learn. Οι orchestrators που ακολουθούν προσφέρουν αυτά τα production concerns out of the box.

Apache Airflow 3.0: Ο βετεράνος του orchestration

Το Apache Airflow 3.0, που κυκλοφόρησε τον Απρίλιο του 2025, είναι η μεγαλύτερη αναβάθμιση του project από την πρώτη κυκλοφορία του. Οι σημαντικότερες αλλαγές: πλήρης διαχωρισμός των tasks από τον scheduler μέσω του νέου Task Execution API, edge executor για hybrid cloud deployments, DAG versioning, και ένα ανασχεδιασμένο React-based UI. Το Airflow είναι η κυρίαρχη επιλογή σε μεγάλες επιχειρήσεις γιατί έχει το πιο εκτενές οικοσύστημα providers: πάνω από 90 επίσημα integrations με AWS, GCP, Azure, Snowflake, dbt, Kubernetes και δεκάδες άλλα.

Το βασικό δομικό στοιχείο είναι το DAG (Directed Acyclic Graph). Στο Airflow 3 προτιμούμε το TaskFlow API με decorators αντί για κλασικά operators, γιατί ο κώδικας γίνεται πιο Pythonic και πιο κοντά στη σύνταξη του Prefect.

from datetime import datetime, timedelta
from airflow.decorators import dag, task
import pandas as pd
import requests

DEFAULT_ARGS = {
    "owner": "data-team",
    "retries": 3,
    "retry_delay": timedelta(minutes=5),
}

@dag(
    dag_id="daily_sales_etl",
    default_args=DEFAULT_ARGS,
    start_date=datetime(2026, 1, 1),
    schedule="0 6 * * *",  # καθημερινά στις 06:00 UTC
    catchup=False,
    tags=["etl", "sales"],
)
def daily_sales_etl():

    @task
    def extract_from_api(day: str) -> list[dict]:
        r = requests.get(
            "https://api.example.com/orders",
            params={"date": day},
            timeout=30,
        )
        r.raise_for_status()
        return r.json()["orders"]

    @task
    def transform(orders: list[dict]) -> list[dict]:
        df = pd.DataFrame(orders)
        df = df.drop_duplicates(subset=["order_id"])
        df["total_eur"] = df["total_usd"] * 0.92
        df["order_month"] = pd.to_datetime(df["order_date"]).dt.to_period("M").astype(str)
        return df.to_dict(orient="records")

    @task
    def load_to_warehouse(records: list[dict]) -> int:
        from sqlalchemy import create_engine
        engine = create_engine("postgresql+psycopg2://user:pass@warehouse:5432/analytics")
        pd.DataFrame(records).to_sql("orders", engine, if_exists="append", index=False)
        return len(records)

    orders = extract_from_api("{{ ds }}")
    transformed = transform(orders)
    load_to_warehouse(transformed)

daily_sales_etl()

Prefect 3: Το σύγχρονο Pythonic ETL framework

Το Prefect 3 κυκλοφόρησε τον Σεπτέμβριο του 2024 και προσφέρει μια εντελώς διαφορετική φιλοσοφία από το Airflow. Αντί για DAGs που ορίζονται μία φορά και εκτελούνται στατικά, τα Prefect flows είναι κανονικές Python συναρτήσεις που δηλώνουν το DAG τους δυναμικά κατά τον χρόνο εκτέλεσης. Αυτό σημαίνει ότι μπορείς να έχεις conditional logic, loops και dynamic task generation με απλό Python, χωρίς hacks. Το Prefect 3 προσθέτει transactions με semantics για ατομικές λειτουργίες, ενισχυμένο caching, και ενσωματωμένη υποστήριξη για event-driven flows.

Το παρακάτω παράδειγμα δείχνει το ίδιο σενάριο πωλήσεων ως Prefect flow. Πρόσεξε πόσο κοντά στο "κανονικό Python" είναι.

from prefect import flow, task, get_run_logger
from prefect.tasks import task_input_hash
from datetime import timedelta, date
import pandas as pd
import requests

@task(retries=3, retry_delay_seconds=60, cache_key_fn=task_input_hash,
      cache_expiration=timedelta(hours=1))
def extract_from_api(day: str) -> pd.DataFrame:
    log = get_run_logger()
    log.info("Fetching orders for %s", day)
    r = requests.get(
        "https://api.example.com/orders",
        params={"date": day},
        timeout=30,
    )
    r.raise_for_status()
    return pd.DataFrame(r.json()["orders"])

@task
def transform(df: pd.DataFrame) -> pd.DataFrame:
    df = df.drop_duplicates(subset=["order_id"])
    df["total_eur"] = df["total_usd"] * 0.92
    df["order_month"] = pd.to_datetime(df["order_date"]).dt.to_period("M").astype(str)
    return df

@task
def load_to_warehouse(df: pd.DataFrame) -> int:
    from sqlalchemy import create_engine
    engine = create_engine("postgresql+psycopg2://user:pass@warehouse:5432/analytics")
    df.to_sql("orders", engine, if_exists="append", index=False)
    return len(df)

@flow(name="daily-sales-etl", log_prints=True)
def daily_sales_etl(target_day: str | None = None):
    day = target_day or date.today().isoformat()
    orders = extract_from_api(day)
    if orders.empty:
        print(f"Δεν βρέθηκαν παραγγελίες για {day}")
        return 0
    clean = transform(orders)
    return load_to_warehouse(clean)

if __name__ == "__main__":
    daily_sales_etl.serve(name="daily-sales", cron="0 6 * * *")

Πρόσεξε το .serve() στο τέλος: μια εντολή που καταχωρεί το flow, ρυθμίζει τον χρονοπρογραμματισμό και εκκινεί έναν worker. Δεν χρειάζεσαι ξεχωριστό scheduler ή webserver όπως στο Airflow. Ειλικρινά, όταν το πρωτοδοκίμασα ήταν από τα πράγματα που με έκαναν να αναρωτηθώ γιατί το Airflow απαιτεί τόσο setup. Για μεγαλύτερα production setups, χρησιμοποιείς το Prefect Cloud (SaaS) ή self-hosted server, αλλά το .serve() είναι ιδανικό για μικρές ομάδες. Αν έρχεσαι από πιο κλασικά data engineering εργαλεία και θέλεις να συγκρίνεις μηχανές μετασχηματισμού, ρίξε μια ματιά στον οδηγό μας Polars vs Pandas. Το Polars συνεργάζεται εξαιρετικά με το Prefect για γρήγορους μετασχηματισμούς.

Dagster: Asset-based data orchestration

Το Dagster εισάγει μια θεμελιωδώς διαφορετική νοοτροπία. Αντί να ορίζεις εργασίες (tasks) που τρέχουν με μια σειρά, ορίζεις software-defined assets, δηλαδή τα ίδια τα δεδομένα ως πρωταρχικά αντικείμενα. Ένα asset είναι ένας πίνακας, ένα ML μοντέλο, ή ένα report που θέλεις να παράγεις. Το Dagster υπολογίζει αυτόματα τι πρέπει να τρέξει για να πάρεις ένα asset up-to-date, βασισμένο στις εξαρτήσεις.

Αυτή η προσέγγιση ταιριάζει εξαιρετικά με σύγχρονες πρακτικές όπως τα data contracts, γιατί κάθε asset έχει σαφή schema, metadata και freshness policies. Το Dagster έχει επίσης εξαιρετικό developer experience με το τοπικό Dagster UI (πρώην Dagit), όπου μπορείς να δεις τη ροή, να τρέξεις partial refreshes, και να δεις lineage. Το παρακάτω παράδειγμα δείχνει τα ίδια τρία βήματα ως assets.

from dagster import asset, AssetExecutionContext, Definitions, ScheduleDefinition, define_asset_job
import pandas as pd
import requests

@asset(group_name="sales")
def raw_orders(context: AssetExecutionContext) -> pd.DataFrame:
    day = context.partition_key if context.has_partition_key else "2026-07-13"
    r = requests.get("https://api.example.com/orders", params={"date": day}, timeout=30)
    r.raise_for_status()
    df = pd.DataFrame(r.json()["orders"])
    context.add_output_metadata({"row_count": len(df), "day": day})
    return df

@asset(group_name="sales")
def cleaned_orders(raw_orders: pd.DataFrame) -> pd.DataFrame:
    df = raw_orders.drop_duplicates(subset=["order_id"])
    df["total_eur"] = df["total_usd"] * 0.92
    df["order_month"] = pd.to_datetime(df["order_date"]).dt.to_period("M").astype(str)
    return df

@asset(group_name="sales")
def warehouse_orders(context: AssetExecutionContext, cleaned_orders: pd.DataFrame) -> int:
    from sqlalchemy import create_engine
    engine = create_engine("postgresql+psycopg2://user:pass@warehouse:5432/analytics")
    cleaned_orders.to_sql("orders", engine, if_exists="append", index=False)
    context.add_output_metadata({"rows_loaded": len(cleaned_orders)})
    return len(cleaned_orders)

daily_job = define_asset_job("daily_sales_job", selection=["raw_orders", "cleaned_orders", "warehouse_orders"])
daily_schedule = ScheduleDefinition(job=daily_job, cron_schedule="0 6 * * *")

defs = Definitions(
    assets=[raw_orders, cleaned_orders, warehouse_orders],
    schedules=[daily_schedule],
)

Airflow vs Prefect vs Dagster: Ολοκληρωμένη σύγκριση

Και τα τρία εργαλεία μπορούν να λύσουν το ίδιο πρόβλημα, αλλά ξεκινούν από διαφορετικές αρχές. Ο παρακάτω πίνακας συνοψίζει τις κρίσιμες διαφορές που πρέπει να σκεφτείς πριν επιλέξεις.

ΧαρακτηριστικόApache Airflow 3.0Prefect 3Dagster
Τρέχουσα έκδοση (2026)3.0.x3.x1.9+
Μοντέλο εκτέλεσηςΣτατικά DAGsΔυναμικά flowsSoftware-defined assets
Καμπύλη μάθησηςΑπότομηΉπιαΜέτρια
Ecosystem / Providers90+ providersGrowing collectionΚαλή, εστίαση σε data stack
UIΝέο React UI (v3)Καθαρό, cloud-firstExcellent asset lineage
Local developmentΑπαιτεί docker-composeΈνας φάκελος, .serve()Το καλύτερο τοπικό DX
Πληρωμένη υπηρεσίαAstronomer, MWAA, Cloud ComposerPrefect CloudDagster+
Ιδανικό γιαEnterprise batch pipelinesData science teams, ML workflowsModern data platforms με dbt

Πρακτικός κανόνας επιλογής: αν η ομάδα σου ήδη τρέχει Airflow, μείνε στο Airflow — η αναβάθμιση σε 3.0 φέρνει αρκετά νέα χαρακτηριστικά ώστε να μην χρειάζεται migration. Αν ξεκινάς από την αρχή με μικρή ομάδα, το Prefect έχει το λιγότερο overhead. Αν χτίζεις σοβαρή data platform με πολλές πηγές και θες lineage/data catalog από την πρώτη μέρα, το Dagster είναι η πιο μελλοντοστραφής επιλογή.

Data Quality Checks και Monitoring σε production

Ένα ETL pipeline χωρίς data quality checks είναι κρυμμένη τεχνική βλάβη περιμένοντας να σκάσει. Οι δύο κυρίαρχες βιβλιοθήκες για ελέγχους ποιότητας το 2026 είναι το Great Expectations 1.x (πλέον με πιο απλοποιημένο API) και η Pandera, που είναι πιο ελαφριά και κατάλληλη για inline checks μέσα σε pandas/Polars pipelines. Και τα δύο ενσωματώνονται εύκολα και στους τρεις orchestrators.

import pandera as pa
from pandera.typing import Series, DataFrame

class OrderSchema(pa.DataFrameModel):
    order_id: Series[str] = pa.Field(unique=True, nullable=False)
    customer_id: Series[str] = pa.Field(nullable=False)
    total_eur: Series[float] = pa.Field(ge=0, le=100_000)
    order_month: Series[str] = pa.Field(str_matches=r"^\d{4}-\d{2}$")

    class Config:
        strict = True  # απαγορεύει επιπλέον στήλες

@pa.check_types
def transform_validated(df: DataFrame) -> DataFrame[OrderSchema]:
    df = df.drop_duplicates(subset=["order_id"])
    df["total_eur"] = df["total_usd"] * 0.92
    df["order_month"] = df["order_date"].dt.to_period("M").astype(str)
    return df

Όταν το transform_validated επιστρέψει DataFrame που παραβιάζει το schema, η Pandera πετάει εξαίρεση με λεπτομερές report. Σε συνδυασμό με τον orchestrator, ο έλεγχος αποτυγχάνει καθαρά και ενεργοποιεί alerts. Για την πλευρά του monitoring, χρησιμοποίησε OpenTelemetry traces (υποστηρίζεται εγγενώς και από τα τρία εργαλεία), εξαγωγή μετρικών σε Prometheus/Grafana, και integration με ένα incident management σύστημα όπως το PagerDuty. Δες την επίσημη τεκμηρίωση OpenTelemetry για Python για τη ρύθμιση των exporters.

Βέλτιστες πρακτικές για παραγωγικά ETL pipelines

Ας το πούμε ξεκάθαρα. Ανεξάρτητα από τον orchestrator που θα επιλέξεις, υπάρχουν αρχές που εφαρμόζονται παντού και ξεχωρίζουν τα ερασιτεχνικά pipelines από τα production-grade. Ορίστε αυτές που εγώ βρίσκω πιο κρίσιμες.

1. Idempotency: κάθε task πρέπει να μπορεί να ξανατρέξει με ασφάλεια

Χρησιμοποίησε MERGE ή UPSERT αντί για INSERT στο βήμα load, ώστε ένα retry να μην δημιουργήσει διπλότυπα. Για partitioned tables, γράψε ολόκληρη partition σε staging table και μετά κάνε atomic swap.

2. Ενδιάμεσα αρχεία σε Parquet, όχι JSON/CSV

Το Parquet με PyArrow είναι 5 έως 10 φορές ταχύτερο για ανάγνωση/εγγραφή, μικρότερο σε μέγεθος, και διατηρεί τύπους δεδομένων. Είναι το de facto format για data lakes και υποστηρίζεται από κάθε σοβαρή αποθήκη δεδομένων. Δες την επίσημη τεκμηρίωση PyArrow Parquet για βέλτιστες παραμέτρους συμπίεσης και partitioning.

3. Χρήση secrets managers, ποτέ hardcoded credentials

Ολα τα credentials πρέπει να έρχονται από AWS Secrets Manager, GCP Secret Manager, HashiCorp Vault, ή τουλάχιστον από περιβαλλοντικές μεταβλητές μέσω .env που δεν κάνει commit. Το Airflow έχει το Connections abstraction, το Prefect τα Blocks, και το Dagster τα Resources.

4. Παρακολούθηση freshness και data volume

Ρύθμισε alerts όταν ο όγκος δεδομένων μειώνεται απότομα (π.χ. -50% σε σχέση με το κινούμενο μέσο όρο). Αυτό είναι συχνά σημάδι σπασμένης πηγής, όχι πραγματικής πτώσης της κίνησης. Για μαθήματα και συνδυασμούς αυτών των εργαλείων με ML training pipelines, ρίξε μια ματιά στον οδηγό Μηχανική Μάθηση στην Python με Scikit-Learn Pipelines.

Συχνές ερωτήσεις

Ποια είναι η καλύτερη βιβλιοθήκη ETL για Python το 2026;

Δεν υπάρχει καθολικά "καλύτερη". Εξαρτάται από την περίπτωσή σου. Το Apache Airflow 3.0 είναι η ασφαλής enterprise επιλογή, το Prefect 3 ταιριάζει σε μικρές, ευέλικτες ομάδες, και το Dagster είναι ιδανικό για σύγχρονες data platforms με έμφαση στο asset lineage και το dbt.

Χρειάζομαι πάντα orchestrator ή μπορώ να χρησιμοποιήσω απλό Python script;

Για μικρές, ανεξάρτητες εργασίες που τρέχουν σε cron και δεν έχουν εξαρτήσεις, ένα απλό Python script με pandas ή Polars είναι μια χαρά. Οι orchestrators δικαιολογούνται όταν έχεις πολλά αλληλοεξαρτώμενα jobs, χρειάζεσαι retries/alerts, ή θέλεις κεντρικό UI για monitoring.

Είναι δωρεάν το Prefect;

Ναι, το Prefect Core είναι πλήρως ανοιχτού κώδικα (Apache 2.0). Η εταιρεία προσφέρει και το Prefect Cloud, μια πληρωμένη managed υπηρεσία με δωρεάν tier για μικρές ομάδες, αλλά μπορείς να τρέξεις το self-hosted Prefect server χωρίς κόστος.

Ποια είναι η διαφορά ανάμεσα σε ETL και ELT;

Στο ETL, ο μετασχηματισμός γίνεται πριν από τη φόρτωση στην αποθήκη, οπότε τα δεδομένα φτάνουν έτοιμα προς κατανάλωση. Στο ELT, φορτώνεις πρώτα raw δεδομένα και μετασχηματίζεις εντός της αποθήκης (συνήθως με SQL/dbt). Το ELT κυριαρχεί στις σύγχρονες cloud data warehouses επειδή αξιοποιεί την υπολογιστική τους ισχύ.

Μπορώ να τρέξω Airflow τοπικά για να το δοκιμάσω;

Ναι. Ο πιο εύκολος τρόπος είναι το docker-compose setup από το επίσημο repo, ή το astro dev init από το Astronomer CLI που στήνει τοπικό περιβάλλον με ένα command. Για μεμονωμένη δοκιμή, το airflow standalone εκκινεί όλα τα components (webserver, scheduler, DB) σε μία εντολή.

Πώς αντιμετωπίζω αστοχίες σε ETL pipelines;

Ρύθμισε τρεις γραμμές άμυνας: (1) Retries με exponential backoff στο επίπεδο του task για παροδικά σφάλματα δικτύου. (2) Data quality checks με Great Expectations ή Pandera μεταξύ των βημάτων. (3) Alerts μέσω Slack/PagerDuty όταν αποτυγχάνει ή τρέχει πολύ αργά ένα job.

Editorial Team
Σχετικά με τον Συγγραφέα Editorial Team

Our team of expert writers and editors.