DAG লেখা ও schedule
এই পাঠে যা শিখবেন
- 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-এ।
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()
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।
execution_date-এর নাম change হয়ে logical_date এবং data_interval_start / data_interval_end introduce। নতুন DAG-এ এই terminology preferred।
Common cron:
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
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।
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
৬ · 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 অপরিবর্তিত।
৭ · 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-এ অসাধারণ।
৯ · 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-errorsCI-তে।
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_endintroduced।{{ ds }}এখনোdata_interval_startpoint করে।- "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_failedstate-এ 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=TrueDAG-এ — একটি 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_valueoverride। - বড় 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এtenacitylibrary দিয়ে 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) ভাবার সময়।
অনুশীলন
-
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")। -
লিখুন: বাংলাদেশের ৩টি 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() -
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-এ আপনার পরবর্তী পদক্ষেপ
- পাঠ ১৫ · dbt — analytics engineering পরবর্তী পাঠ Airflow + dbt একসাথে — SQL-based transformation pipeline।
- পাঠ ১৩ · Apache Airflow পরিচিতি আগের পাঠ Concept বুঝে আসা ভাল — তারপর code।
- পাঠ ২৫ · Data quality testing এই পাঠের সাথে সম্পর্কিত প্রতিটি DAG-এ data quality check — Great Expectations, dbt tests।
- সব AI Courses দেখুন ABCL TECH Python, ML, DL, NLP, CV, GenAI, RL, MLOps — সব AI কোর্স একসাথে।
pip install apache-airflow করে DAG file syntax check করতে পারেন।