পাঠ ১৪ · ২৯-এর মধ্যে · মডিউল ২

DAG লেখা ও schedule

Writing DAGs — TaskFlow, schedule, XCom, idempotency
৮ মিনিট পড়া মাঝারি · Intermediate Python সহ

এই পাঠে যা শিখবেন

  • Classic operator API ও modern TaskFlow API — দু'টি style
  • Schedule interval, start_date, catchup — প্রায়ই বুঝতে কঠিন trio
  • XCom — task-এর মধ্যে data পাঠানো
  • retry, params, templating — production DAG-এর গুরুত্বপূর্ণ অংশ

১ · Classic operator vs TaskFlow — দু'টি style

Airflow-এ DAG দু'ভাবে লেখা যায়। পুরাতন (Airflow ১.x থেকে): operator class instantiate করে dependency define। নতুন (২.০+): Python decorator।

Classic style:

extract = PythonOperator(task_id="extract", python_callable=extract_fn)
transform = PythonOperator(task_id="transform", python_callable=transform_fn)
extract >> transform

TaskFlow style (recommended):

@task
def extract():
    return data

@task
def transform(data):
    return result

transform(extract())   # dependency auto-inferred
কোনটি বাছবেন

১) TaskFlow: Python function-heavy DAG-এ। কম boilerplate, type hint, native return value।
২) Classic: Bash/SQL/Spark operator-এ — যেখানে কোনো Python function নেই।
৩) Mix-and-match: বেশিরভাগ production DAG-এ দু'টোই থাকে।

২ · একটি সম্পূর্ণ DAG — Daraz daily summary

নিচের DAG-টি Daraz-এর order data থেকে daily revenue, top-10 product summary তৈরি করে — TaskFlow style-এ।

Python · Airflow TaskFlow
from airflow.decorators import dag, task
from datetime import datetime, timedelta
import pandas as pd

@dag(
    dag_id="daraz_daily_summary",
    description="Daily revenue + top products",
    schedule="0 2 * * *",          # রাত ২টা প্রতিদিন
    start_date=datetime(2026, 5, 1),
    catchup=False,
    default_args={
        "owner": "data-team",
        "retries": 3,
        "retry_delay": timedelta(minutes=5),
    },
    tags=["daraz", "daily"],
)
def daraz_daily_summary():

    @task
    def extract(ds: str) -> str:
        """MySQL → /tmp/orders_{date}.parquet"""
        path = f"/tmp/orders_{ds}.parquet"
        # ... query MySQL where date = ds, write parquet ...
        return path

    @task
    def compute_revenue(orders_path: str, ds: str) -> dict:
        df = pd.read_parquet(orders_path)
        return {
            "date": ds,
            "total_revenue_bdt": float(df["amount"].sum()),
            "order_count": int(len(df)),
        }

    @task
    def compute_top_products(orders_path: str) -> list:
        df = pd.read_parquet(orders_path)
        top = df.groupby("product_id")["amount"].sum() \
                .nlargest(10).reset_index()
        return top.to_dict("records")

    @task
    def load_to_warehouse(revenue: dict, top_products: list):
        # UPSERT into Snowflake / BigQuery
        # idempotent: same date overwrites
        print(f"Revenue: {revenue['total_revenue_bdt']:,.0f} BDT")
        print(f"Top: {top_products[0]}")

    # Dependency graph (auto-built from function call)
    orders = extract("{{ ds }}")
    rev = compute_revenue(orders, "{{ ds }}")
    top = compute_top_products(orders)
    load_to_warehouse(rev, top)

dag = daraz_daily_summary()

    
DAG graph: extract → [compute_revenue, compute_top_products] → load_to_warehouse। Airflow auto-detect করে কারণ orders দু'টি function-এ pass হয়েছে। Parallel-এ revenue ও top চলবে — tighter execution।

৩ · Schedule interval — সবচেয়ে confusing trio

তিন parameter এক সাথে কাজ করে — অনেক ভুলের উৎস:

  • start_date: DAG-এর প্রথম "logical" date।
  • schedule: কত ঘন ঘন চালাবে।
  • catchup: পুরাতন missed run চালাবে কি?

Airflow-এর মাইন্ড-বেন্ডিং rule: একটি DAG run-এর execution_date = period-এর শুরু, run হয় period শেষে।

যেমন: schedule="@daily", start_date=2026-05-01। প্রথম run হবে ২ মে ০০:০০-এ — কিন্তু তার execution_date = 2026-05-01। কারণ ১ তারিখের data ১ তারিখ শেষে available।

Airflow ২.৪+ থেকে execution_date-এর নাম change হয়ে logical_date এবং data_interval_start / data_interval_end introduce। নতুন DAG-এ এই terminology preferred।

Common cron:

Cron expression — Airflow
schedule="@daily"        # 0 0 * * *  (midnight UTC)
schedule="@hourly"       # 0 * * * *
schedule="0 2 * * *"     # প্রতিদিন রাত ২টা UTC
schedule="0 */6 * * *"   # প্রতি ৬ ঘণ্টা
schedule="0 18 * * 5"    # প্রতি শুক্রবার রাত ৬টা
schedule="@weekly"       # 0 0 * * 0  (রবিবার midnight)
schedule="@monthly"      # 0 0 1 * *
schedule=None            # manual trigger only
schedule="@once"         # একবার চলে শেষ

# বাংলাদেশ time zone (BST = UTC+6)
# UTC 02:00 = BST 08:00
# UTC 18:00 = BST 00:00 (পরদিন)
# DAG-এ timezone explicitly set করুন:
import pendulum
start_date=pendulum.datetime(2026, 5, 1, tz="Asia/Dhaka")
schedule="0 8 * * *"   # BST 8 AM

    
Airflow default UTC-তে চলে। Bangladesh team-এ timezone confusion এড়াতে pendulum.timezone("Asia/Dhaka") সবসময় ব্যবহার করুন। তবে UTC-তে রাখা অনেক সহজ — DST issue নেই।

৪ · catchup — পুরাতন run চালানোর রহস্য

catchup=True: DAG-এর start_date থেকে আজ পর্যন্ত প্রতিটি missed run চলবে। যদি start_date=2025-01-01 এবং আজ ২০২৬-০৫-০৯ — ৪৯৪টি daily run!

catchup=False (recommended default): শুধু latest run। পুরাতন data দরকার হলে manual backfill।

Pathao-এর rider weekly settlement DAG কল্পনা করুন। ছুটি কাটিয়ে ২ মাস পরে এসে DAG enable করলেন — catchup=True হলে — Airflow ৮ সপ্তাহ একসাথে run করবে! Cluster overload, double-payment risk। সাধারণত catchup=False।

৫ · XCom — task-এর মধ্যে data pass

XComXCom"Cross-Communication" — task-এর মধ্যে ছোট data exchange। Default: metadata DB-তে JSON serialized। Custom backend-এ S3 বা GCS রাখা যায়। = Cross-Communication। Task A-এর return value Task B-তে use করতে — XCom।

TaskFlow-এ এটি transparent — function return value auto-push হয়, parameter auto-pull:

@task
def extract():
    return {"count": 1000, "date": "2026-05-09"}

@task
def transform(data):    # data = upper return value
    return data["count"] * 2
XCom-এর সীমা: default-এ metadata DB (PostgreSQL)-এ JSON store। প্রতি XCom ~১ MB-র নিচে রাখুন। বড় DataFrame XCom-এ পাঠাবেন না — file path বা S3 URI পাঠান।

৬ · Idempotency — {{ ds }} macro

প্রতিটি Airflow task context পায় — যেখানে {{ ds }} = execution date (YYYY-MM-DD)। SQL/Bash command-এ embed করলে — প্রতি run-এ আলাদা।

BashOperator(
    task_id="extract",
    bash_command="python extract.py --date={{ ds }}",
)

# SQL templating:
PostgresOperator(
    task_id="upsert",
    sql="""
        DELETE FROM revenue WHERE date = '{{ ds }}';
        INSERT INTO revenue
        SELECT date, sum(amount) FROM orders
        WHERE date = '{{ ds }}' GROUP BY date;
    """,
)

এতে retry/backfill safe — একই date বারবার চললে result অপরিবর্তিত।

Daraz daily summary DAG — execution extract → [revenue, top_products] → load ⏰ schedule fires 02:00 UTC daily 📥 extract MySQL → parquet 💰 compute_revenue sum(amount) by date 🏆 top_products groupBy product, top 10 📤 load UPSERT trigger XCom XCom ❌ যদি কোনো task fail হয় retry × 3 (5 min gap each) retry শেষ → email + Slack alert downstream task SKIPPED Idempotent design = retry safe + backfill possible
একটি TaskFlow DAG-এর execution flow। Parallel branch + automatic retry + XCom data passing।

৭ · Retry, alert, ও SLA

default_args-এ DAG-wide setting:

default_args = {
    "retries": 3,                          # ব্যর্থ হলে ৩ বার retry
    "retry_delay": timedelta(minutes=5),    # প্রতি retry-এর মধ্যে ৫ মিনিট
    "retry_exponential_backoff": True,      # 5, 10, 20 min
    "max_retry_delay": timedelta(hours=1),
    "email_on_failure": True,
    "email": ["data-oncall@daraz.com.bd"],
    "sla": timedelta(hours=2),              # ২ ঘণ্টায় শেষ না হলে alert
}

SLA miss = task expected time-এ শেষ হয়নি কিন্তু এখনো চলছে (বা এখনো শুরু হয়নি)। UI-তে red mark, email পাঠায়।

৮ · Params — DAG runtime input

Manual trigger-এ user থেকে input নিতে — Params API:

from airflow.models.param import Param

@dag(
    params={
        "target_date": Param(type="string", format="date"),
        "limit": Param(default=100, type="integer", minimum=1),
    },
)
def my_dag():
    @task
    def process(**ctx):
        date = ctx["params"]["target_date"]
        limit = ctx["params"]["limit"]
        # ...

UI-তে "Trigger DAG w/ config" button-এ form দেখাবে। Manual backfill, ad-hoc run-এ অসাধারণ।

Production tip: Sensitive value (DB password, API key) Params-এ রাখবেন না। Airflow Connections বা Variables ব্যবহার করুন — encrypted store।

৯ · Best practice — production checklist

  • Top-level কোড সরান: heavy import, DB query — task-এর ভেতরে।
  • catchup=False: default থেকে নিয়ে — surprise backfill এড়ান।
  • Idempotent task: {{ ds }} ব্যবহার করুন, now() না।
  • Owner ও tags: বড় team-এ filtering সহজ।
  • doc_md: DAG-এ markdown documentation embed — UI-তে দেখা যাবে।
  • Test: airflow tasks test <dag_id> <task_id> <date> — local-এ একটি task run।
  • CI/CD: DAG file Git-এ, পরিবর্তনে linter (ruff) ও airflow dags list-import-errors CI-তে।
Airflow ২-এ default schedule semantics বদলেছে। নতুন DAG-এ schedule= argument ব্যবহার করুন (পুরাতন schedule_interval= deprecated)। উভয়ই কাজ করে কিন্তু documentation দেখলে confusion হয়।

ভাবনার প্রশ্ন

প্রতিটি প্রশ্ন নিজে কিছুক্ষণ ভাবুন — তারপর "→ উত্তর" চাপুন।

প্র ০১ "Execution date" নিয়ে নতুনদের সবচেয়ে বেশি confusion। ৫ মে DAG run হলে কেন {{ ds }} = 2026-05-04? এই design কেন বেছে নিল Airflow?

এটি Airflow-এর "interval-end" semantics — পরিকল্পিত design choice, ভুল নয়। বুঝতে গেলে batch processing-এর historical context দরকার।

মূল প্রশ্ন: "৫ মে-র data" কখন available?

  • ৫ মে শেষ হওয়ার আগে — পুরো দিনের data নেই।
  • ৫ মে শেষ হলে (অর্থাৎ ৬ মে ০০:০০-এ) — পুরো ৫ মে-র data যত হওয়ার তত।
  • তাই "৫ মে process" মানে — ৬ মে চালান।

Airflow-এর rule:

  • প্রতিটি DAG run একটি "interval" cover করে: [start, end)।
  • Run trigger হয় interval-এর end-এ।
  • execution_date (পুরাতন নাম) / data_interval_start = interval-এর start।
  • {{ ds }} = data_interval_start-এর YYYY-MM-DD।

ফলাফল:

  • "@daily" + start_date 2026-05-01 →
  • প্রথম run: trigger ২ মে ০০:০০-এ, ds = 2026-05-01, ১ মে-র data process।
  • দ্বিতীয় run: trigger ৩ মে ০০:০০, ds = 2026-05-02।

কেন এই design?

  • Idempotency: WHERE date = '{{ ds }}' — সবসময় completed day-এর data।
  • Backfill: পুরাতন interval পুনরায় run করা সহজ — ds পরিচিত।
  • Hourly/sub-hourly: একই pattern — interval শেষে process।

Practical implication:

  • "আজকের data" পেতে — {{ ds }} চলবে না (গতকাল process)।
  • "Real-time" দরকার হলে — Airflow নয়, streaming।
  • BI report-এ users কে বুঝান: "৫ মে-র data" = ৬ মে সকালে available।

Airflow ২.৪+ পরিবর্তন:

  • data_interval_start ও data_interval_end introduced।
  • {{ ds }} এখনো data_interval_start point করে।
  • "Run At" UI-তে actual trigger time, "Logical Date" = ds।

Bangladesh fintech example: bKash daily settlement — ৫ মে-র সব transaction settle হবে ৬ মে ভোরে। DAG schedule "0 5 * * *" (UTC), ds = 2026-05-05, query "WHERE txn_date = '{{ ds }}'"। ৬ মে ১১:০০ AM BST-তে stakeholder report পায় — যা ৫ মে-র complete data।

মূল কথা: "data interval" mental model adopt করুন। তারিখ নিয়ে confusion আর হবে না।

প্র ০২ আপনার DAG-এ ১০টি task। Task ৪ fail হলো। Airflow auto-retry শুরু হবে। কিন্তু task ৫-৭ ইতিমধ্যে চলছে — এদের কী হবে? আপনি কীভাবে handle করবেন?

এটি real production question। আগে DAG topology বুঝে নিই: ১০ task-এ অনেক branch থাকতে পারে।

Scenario analyze:

  • Task 4 fail এবং Task 5-7 চলছে — মানে এরা Task 4-এর downstream নয় (parallel branch)।
  • Task 4-এর downstream (যেমন Task 8) Task 4 success-এর জন্য wait করছিল।

Default behavior:

  • Task 5-7 চলতে থাকবে — এরা Task 4-এর dependent না।
  • Task 4 retry করবে (retries=3)।
  • Task 4-এর downstream task (যেমন Task 8) upstream_failed state-এ skip হবে যদি retry-ও fail।

Trigger rules — দশটি option:

  • all_success (default): সব upstream success হলে চলবে।
  • all_failed: সব fail হলে — recovery task।
  • one_success: যেকোনো ১টি success — branch merge।
  • one_failed: যেকোনো ১টি fail — alert task।
  • none_failed: কেউ fail না হলে (success বা skipped)।
  • all_done: সব finish হলে (যেকোনো state) — cleanup task।

Production pattern: cleanup + alert

extract >> transform >> load
load >> cleanup       # trigger_rule="all_done"
load >> alert_failure # trigger_rule="one_failed"

"Stop everything on failure" pattern:

  • fail_fast=True DAG-এ — একটি fail হলেই DAG-এর সব running task kill।
  • Use case: critical sequential pipeline যেখানে partial completion বিপজ্জনক।
  • Default off — overusage থেকে সাবধান।

"Continue despite failure" pattern:

  • Downstream task-এ trigger_rule="all_done"।
  • Failed branch-এর result null/skipped, কিন্তু others process।
  • Use case: non-critical enrichment fail হলেও main pipeline চলবে।

Manual intervention:

  • UI-তে failed task → "Clear" → re-queue।
  • Downstream-ও clear করতে পারেন: cascade।
  • "Mark Success" — manually fix data, তারপর mark।
  • CLI: airflow tasks clear <dag_id> --task-ids extract।

Bangladesh-এর reality: Pathao analytics DAG-এ — ride extract fail হলেও driver settlement চালাতে চান। দু'টি independent branch বানান, alert task সব fail capture। main pipeline-এ partial result acceptable করুন।

মূল কথা: Task failure design-time-এ ভাবা — operational nightmare এড়ায়। trigger_rule + retry + alert-এর সঠিক combination = robust DAG।

প্র ০৩ আপনার pandas DataFrame ৫০০ MB। Task A থেকে Task B-তে XCom-এ পাঠাতে চান। কী হবে? কী alternative?

XCom-এ ৫০০ MB — disaster। কেন এবং কী করণীয়, ব্যাখ্যা করি।

কী ঘটবে:

  • Default XCom backend = metadata DB (PostgreSQL)।
  • JSON serialize → ৫০০ MB blob → DB-তে INSERT।
  • Failure cases:
  • (ক) Postgres default max_allowed_packet — limit hit।
  • (খ) Serialize-এ executor memory burst।
  • (গ) Task B start-এর সময় ৫০০ MB pull → memory pressure।
  • (ঘ) Metadata DB bloat — সব XCom retention period পর্যন্ত থাকে।
  • (ঙ) Webserver UI slow — XCom view 500MB load করতে চাইবে।

Airflow guidance: XCom < ১ MB। এর বেশি হলেই warning sign।

Alternative ১: External storage + path passing

@task
def transform(date: str) -> str:
    df = ...  # 500 MB DataFrame
    path = f"s3://daraz-data/staging/{date}.parquet"
    df.to_parquet(path)
    return path   # ছোট string XCom

@task
def load(path: str):
    df = pd.read_parquet(path)
    # ...

সবচেয়ে common pattern। S3, GCS, local path — সব কাজ করে। URL/path = ছোট string।

Alternative ২: Custom XCom backend

  • Airflow ২.০+ — BaseXCom.serialize_value override।
  • বড় object S3-তে save, metadata DB-তে শুধু URI।
  • Configuration: xcom_backend = my_module.S3XComBackend।
  • Code transparent — task থেকে দেখতে normal XCom।

Alternative ৩: Single task

  • Transform + load একই task-এ মার্জ — XCom লাগে না।
  • Trade-off: visibility কম (২ task-এর বদলে ১)।
  • Retry granularity কম।

Alternative ৪: Spark/cluster compute

  • ৫০০ MB pandas-এ — সম্ভবত Spark/Dask-এ যাওয়ার সময়।
  • Airflow শুধু Spark job submit করে — data Spark cluster-এ থাকে।
  • XCom-এ শুধু job_id বা output path।

সাধারণ ভুল:

  • return df.to_dict() — XCom-এ python dict, JSON serializable, কিন্তু still huge।
  • List of dict — Postgres TEXT field overflow।
  • Image bytes XCom-এ — totally wrong।

Best practice rule:

  • XCom = metadata, control flow, references (URI, ID, count)।
  • Real data = external storage।
  • "data lake first" mindset।

মূল কথা: Airflow orchestrator — data mover না। Heavy data work systems (Spark, dbt, Snowflake)-এ; Airflow শুধু coordinate।

প্র ০৪ Pathao-তে আপনি একটি DAG লিখছেন — প্রতি ১৫ মিনিটে driver location aggregate। কিন্তু কখনো API ৩-৪ মিনিট দেরি করে। কীভাবে handle করবেন? Sensor? retry? ভিন্ন pattern?

এটি real-world API integration challenge — multiple valid approach আছে, trade-off-ভিত্তিক।

Option ১: Simple retry

@task(retries=5, retry_delay=timedelta(minutes=1))
def fetch_locations(ds_nodash):
    response = requests.get(API_URL, timeout=30)
    response.raise_for_status()
    return response.json()
  • সরল, কাজ করে।
  • ৫ × ১ min = ৫ min retry window।
  • API down হলে — alert ৫ min পর।
  • Drawback: প্রতি retry-এ task slot occupy।

Option ২: Exponential backoff

retries=5,
retry_exponential_backoff=True,
retry_delay=timedelta(seconds=30),
max_retry_delay=timedelta(minutes=10),
  • Wait: 30s, 60s, 120s, 240s, 480s।
  • Transient failure-এ ভাল — API recovering গেলে দ্রুত পেয়ে যাবে।
  • Long outage-এ — ১৫ মিনিট পর্যন্ত wait।

Option ৩: HttpSensor + reschedule mode

wait_api = HttpSensor(
    task_id="wait_api",
    http_conn_id="pathao_api",
    endpoint="/health",
    poke_interval=30,
    timeout=600,  # 10 min max wait
    mode="reschedule",  # slot ছেড়ে দেয়
)
wait_api >> fetch_locations
  • Sensor ৩০ সেকেন্ডে একবার health check।
  • mode="reschedule" — wait time-এ worker slot free।
  • API ready হলেই fetch।

Option ৪: Increase task timeout

@task(execution_timeout=timedelta(minutes=10))
def fetch_locations(...):
    # API ধৈর্যের সাথে wait
    response = requests.get(API_URL, timeout=300)
  • সবচেয়ে সরল।
  • API connection long-poll করে।
  • Drawback: worker slot block।

Option ৫: Circuit breaker pattern

  • @task এ tenacity library দিয়ে complex retry logic।
  • API consecutive ৫ fail হলে — temporarily skip।
  • Production-grade resilience।

Option ৬: Async / decoupling

  • API webhook-এ event push, queue (SQS/Kafka)।
  • Airflow polling বদলে event-driven।
  • Architecture change — large scale-এ valuable।

কোনটা বাছবেন:

  • Pathao 15-min cycle: Option ২ (exponential) + Option ৪ (timeout) combine।
  • Reasoning: 15 min = tight schedule। দীর্ঘ wait lattice cycle miss করবে।
  • SLA decision: ১২ মিনিটে শেষ না হলে — current cycle drop, পরের cycle-এ catch up।

Monitoring critical:

  • Task duration metric — কখন বেশি wait?
  • Retry count — API quality indicator।
  • Failure rate — alert threshold।

মূল কথা: External API integration-এ retry/timeout/sensor combine করুন। কিন্তু ১৫ min batch-এ ১০ min wait — design pattern পরিবর্তন (event-driven) ভাবার সময়।

অনুশীলন

  1. Schedule লিখুন:
    • (ক) প্রতি কর্মদিবস (সোম-শুক্র) সকাল ৯টা BST।
    • (খ) প্রতি ১৫ মিনিটে।
    • (গ) প্রতি মাসের ১ তারিখে রাত ৩টা UTC।
    • (ক) schedule="0 3 * * 1-5" (UTC 3 AM = BST 9 AM, Mon-Fri)
    • (খ) schedule="*/15 * * * *"
    • (গ) schedule="0 3 1 * *" বা schedule="@monthly"

    BST-তে রাখতে চাইলে — start_date=pendulum.datetime(..., tz="Asia/Dhaka")।

  2. লিখুন: বাংলাদেশের ৩টি bank (DBBL, BRAC, City) থেকে statement download → merge → upload-এর TaskFlow DAG।
    from airflow.decorators import dag, task
    from datetime import datetime
    
    @dag(schedule="0 6 * * *", start_date=datetime(2026,5,1), catchup=False)
    def bank_reconciliation():
    
        @task
        def download(bank: str, ds: str) -> str:
            path = f"s3://recon/{bank}/{ds}.csv"
            # ... download from each bank API ...
            return path
    
        @task
        def merge(paths: list, ds: str) -> str:
            # combine all bank files into one parquet
            out = f"s3://recon/merged/{ds}.parquet"
            # ...
            return out
    
        @task
        def upload(path: str):
            # to Snowflake
            pass
    
        paths = [download(b, "{{ ds }}") for b in ["DBBL", "BRAC", "City"]]
        merged = merge(paths, "{{ ds }}")
        upload(merged)
    
    dag = bank_reconciliation()
  3. Debug: এই DAG কেন idempotent না? Fix করুন।
    @task
    def daily_revenue():
        today = datetime.now().date()
        df = query(f"SELECT * FROM orders WHERE date = '{today}'")
        df.to_csv(f"/data/revenue_{today}.csv")

    সমস্যা:

    • datetime.now() = run time, not logical date। Backfill ভাঙবে।
    • Manual retry-এ আজকের date — কিন্তু আসল target কাল।
    • Output path-ও datetime.now() উপর depend।

    Fix:

    @task
    def daily_revenue(ds: str):
        df = query(f"SELECT * FROM orders WHERE date = '{ds}'")
        df.to_csv(f"/data/revenue_{ds}.csv")
    
    # Call:
    daily_revenue("{{ ds }}")

    এখন {{ ds }} macro = logical date। Backfill, retry, manual run — সব consistent।

আরও পড়ুন · ABCL TECH-এ আপনার পরবর্তী পদক্ষেপ

হাতে-কলমে অনুশীলন: Local-এ Airflow ছাড়াও Google Colab-এ pip install apache-airflow করে DAG file syntax check করতে পারেন।
পূর্ববর্তী পাঠ
পাঠ ১৩ · Apache Airflow পরিচিতি