Apache Airflow পরিচিতি
এই পাঠে যা শিখবেন
- Workflow orchestration কী — কেন cron যথেষ্ট নয়
- DAG, task, operator, sensor — Airflow-এর ভিত্তি শব্দ
- Airflow architecture — চারটি component কীভাবে একসাথে কাজ করে
- কখন Airflow ব্যবহার করবেন, কখন অন্য tool
১ · Workflow orchestration কী — সমস্যাটা কী
কল্পনা করুন আপনি Daraz-এর data team-এ। প্রতি রাতে এই কাজগুলো হতে হবে:
- MySQL থেকে গতকালের order data Snowflake-এ copy
- সেই data clean — duplicate বাদ, type fix
- Customer dimension table-এর সাথে JOIN
- Daily revenue, top products report তৈরি
- Email-এ executive-দের পাঠাও, Slack-এ team-কে
প্রতিটি কাজের dependency আছে। ৩ চলবে শুধু ১ ও ২ শেষ হলে। ৪ শেষ হওয়ার আগে ৫ চললে — empty email।
cron শুধু সময়ে কাজ চালায়। কিন্তু:
১) Task-এর মধ্যে dependency নেই — শুধু সময়সূচি।
২) Task fail হলে retry নেই, alert নেই।
৩) কোন task কখন চলেছে, কত সময় নিয়েছে — কোনো log/UI নেই।
৪) Backfill (পুরাতন তারিখ পুনরায় run) করতে manual script।
AirflowApache AirflowAirbnb-এ ২০১৪-তে তৈরি, ২০১৬-তে Apache-এ donate। আজ পৃথিবীর সবচেয়ে জনপ্রিয় workflow orchestrator। Astronomer, AWS MWAA, GCP Composer — সবাই managed Airflow সরবরাহ করে। এই সব সমস্যার সমাধান। আপনি Python-এ একটি graph লেখেন — Airflow চালায়, retry করে, log রাখে, alert পাঠায়।
২ · DAG — Airflow-এর কেন্দ্রীয় ধারণা
DAGDirected Acyclic GraphDirected = তীরের দিক আছে; Acyclic = চক্র নেই (A → B → A হয় না)। Task-এর dependency representation-এর গাণিতিক structure। = Directed Acyclic Graph। Task-গুলো node, dependency-গুলো edge। "Acyclic" — চক্র নেই। ১ → ২ → ৩ → ১ কখনো নয়।
একটি Airflow DAG হলো একটি Python file যেখানে আপনি:
- DAG-এর schedule define করেন (যেমন প্রতিদিন রাত ২টা)
- Task-গুলো define করেন (Python function, SQL query, bash command)
- Task-এর মধ্যে dependency define করেন (
task1 >> task2)
৩ · Airflow architecture — চারটি component
Airflow চলতে চারটি জিনিস লাগে — প্রতিটি একটি স্বতন্ত্র process:
- Webserver: Flask-ভিত্তিক UI। DAG দেখা, manually trigger, log read, retry — সব এখানে।
- Scheduler: Airflow-এর "মস্তিষ্ক"। প্রতি কয়েক সেকেন্ডে DAG file scan করে — কোনটি এখন run করার সময়, কোন task ready।
- Executor: Scheduler-এর decision অনুযায়ী আসল task চালায়। বিভিন্ন রকম: SequentialExecutor (dev), LocalExecutor, CeleryExecutor (distributed), KubernetesExecutor (cloud-native)।
- Metadata Database: PostgreSQL/MySQL — সব state এখানে। DAG run history, task instance, connection, variable। UI ও Scheduler দু'টোই এখান থেকে read/write।
৪ · Operator ও Sensor — task-এর রকমভেদ
Task-এর ভেতরে আসল কাজটি কে করে? — OperatorOperatorAirflow-এর প্রতিটি task একটি Operator class-এর instance। ১,০০০+ built-in operator আছে — provider package হিসেবে।। Airflow-এ ১,০০০+ built-in operator আছে।
- BashOperator: shell command চালায়।
- PythonOperator: Python function call করে।
- PostgresOperator, MySqlOperator: SQL query চালায়।
- SparkSubmitOperator: Spark cluster-এ job submit।
- EmailOperator, SlackWebhookOperator: notification।
SensorSensorএকটি বিশেষ ধরনের Operator যা একটি condition true না হওয়া পর্যন্ত wait করে। File arrival, S3 object, external DAG completion — সব Sensor দিয়ে। = special operator যা একটি event-এর জন্য wait করে। যেমন:
- S3KeySensor: S3-তে একটি file আসা পর্যন্ত wait।
- HttpSensor: একটি API endpoint ready না হওয়া পর্যন্ত wait।
- ExternalTaskSensor: অন্য DAG-এর task শেষ না হওয়া পর্যন্ত wait।
mode="reschedule" ব্যবহার করুন — sensor task slot ছেড়ে wait time-এ ঘুমিয়ে যায়।
৫ · cron বনাম Airflow — তুলনা
সাধারণ ৩-task pipeline cron-এ:
# /etc/crontab — Daraz daily report (cron version)
# 02:00 — extract from MySQL
0 2 * * * /opt/scripts/extract.sh >> /var/log/extract.log 2>&1
# 02:30 — transform (আশা করছি extract শেষ!)
30 2 * * * /opt/scripts/transform.sh >> /var/log/transform.log 2>&1
# 03:00 — load to warehouse
0 3 * * * /opt/scripts/load.sh >> /var/log/load.log 2>&1
# সমস্যা:
# - extract যদি ৩৫ মিনিট লাগে? transform partial data পাবে
# - load fail হলে কিভাবে জানবেন? কেউ log না দেখলে — silent failure
# - গত মঙ্গলবারের data missing — কীভাবে backfill?
# - ৫০টি pipeline হলে — crontab অপাঠ্য
৬ · একই pipeline Airflow-এ
from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta
default_args = {
"owner": "data-team",
"retries": 3,
"retry_delay": timedelta(minutes=5),
"email_on_failure": True,
"email": ["data-oncall@daraz.com.bd"],
}
with DAG(
dag_id="daraz_daily_report",
description="Daily revenue report",
schedule_interval="0 2 * * *", # প্রতিদিন রাত ২টা
start_date=datetime(2026, 5, 1),
catchup=False,
default_args=default_args,
tags=["daraz", "report"],
) as dag:
extract = BashOperator(
task_id="extract",
bash_command="/opt/scripts/extract.sh",
)
transform = BashOperator(
task_id="transform",
bash_command="/opt/scripts/transform.sh",
)
load = BashOperator(
task_id="load",
bash_command="/opt/scripts/load.sh",
)
# Dependency: extract → transform → load
extract >> transform >> load
extract >> transform >> load — dependency directly লেখা। retry, alert, ownership সব declarative। UI-তে graph দেখা যাবে, প্রতিটি run-এর log click করে read করা যাবে।
৭ · কখন Airflow ব্যবহার করবেন (এবং কখন নয়)
Airflow ভাল:
- Batch ETL: hourly, daily, weekly schedule।
- Multi-step pipeline: dependency সহ ৫-৫০০ task।
- ML training pipeline: data prep → train → evaluate → deploy।
- Cross-system orchestration: S3 → Spark → Snowflake → email।
Airflow ভাল না:
- Sub-second latency: Airflow scheduler heartbeat ৫-১০ সেকেন্ড। Real-time fraud detection নয়।
- Streaming: continuous data flow — Kafka, Flink, Spark Streaming ভাল।
- Heavy computation: Airflow worker compute নয় — Spark/EMR-এ submit করুন।
- Single shell script: ১টি cron job — Airflow over-engineering।
৮ · Airflow alternative — দ্রুত পরিচিতি
২০২২-এর পর "modern data stack"-এ Airflow-এর প্রতিযোগী এসেছে:
- Prefect: Python-native, dynamic DAG (runtime-এ task তৈরি), modern API।
- Dagster: "asset"-centric — task-এর বদলে data asset-কে first-class।
- Temporal: long-running workflow (দিনের পর দিন wait), micro-service orchestration।
- Argo Workflows: Kubernetes-native, container-per-task।
কিন্তু এদের কারো community ও ecosystem Airflow-এর কাছাকাছি নয়। Airflow ২৫০+ provider package, ৩০,০০০ GitHub star, প্রতিটি cloud provider managed offering। নতুন project শুরু করলে Airflow safe choice।
৯ · প্রথম DAG run — local-এ
# দ্রুত local setup (Airflow 2.9+, Python 3.11)
export AIRFLOW_HOME=~/airflow
pip install "apache-airflow==2.9.0" \
--constraint "https://raw.githubusercontent.com/apache/airflow/constraints-2.9.0/constraints-3.11.txt"
# Initialize metadata DB (SQLite, only for dev)
airflow db init
# Create admin user
airflow users create \
--username admin --password admin \
--firstname Admin --lastname User \
--role Admin --email admin@example.com
# DAG file রাখুন
mkdir -p ~/airflow/dags
# (উপরের DAG কোড ~/airflow/dags/daraz_daily_report.py-এ save করুন)
# দু'টি terminal-এ:
airflow webserver --port 8080 # terminal 1
airflow scheduler # terminal 2
# ব্রাউজারে: http://localhost:8080
# Login: admin / admin
# DAG দেখুন → toggle on → trigger
ভাবনার প্রশ্ন
প্রতিটি প্রশ্ন নিজে কিছুক্ষণ ভাবুন — তারপর "→ উত্তর" চাপুন।
প্র ০১ Airflow-এর Executor চারটি — Sequential, Local, Celery, Kubernetes। প্রতিটির use case ও trade-off ব্যাখ্যা করুন। Bangladesh-এর কোন scale-এ কোনটি বাছবেন?
Executor — Airflow-এর সবচেয়ে গুরুত্বপূর্ণ deployment decision। Workload, scale, infrastructure সব এর উপর নির্ভর করে।
SequentialExecutor:
- একসাথে শুধু ১টি task চলে।
- SQLite metadata-এর সাথে default।
- Use: শুধু local development, "hello world" শেখা।
- Production-এ কখনো না।
LocalExecutor:
- Scheduler ও worker একই machine-এ। parallel task সম্ভব।
- PostgreSQL/MySQL metadata দরকার।
- Single VM-এ ১০-৫০টি concurrent task। সহজ setup।
- Use: ছোট startup (১-২০ DAG), single team, ৪-৮ vCPU machine।
- Bangladesh: ছোট fintech-এর internal reporting — ১টি ৪ vCPU instance যথেষ্ট।
CeleryExecutor:
- Worker pool — আলাদা machine-এ। Redis/RabbitMQ message broker।
- Horizontal scaling — worker সংখ্যা বাড়ান।
- Setup complex: Redis, Celery, multiple worker process।
- Use: medium scale (৫০-৫০০ DAG, ১০০-১,০০০ concurrent task)।
- Bangladesh: bKash, Daraz, Grameenphone — typical।
- Drawback: worker pool fixed size, idle resource waste।
KubernetesExecutor:
- প্রতিটি task একটি new pod-এ।
- Auto-scaling — task নেই, pod নেই, খরচ নেই।
- Per-task isolation — dependencies কাছাকাছি থাকে না।
- K8s cluster দরকার, complexity high।
- Use: large scale (১,০০০+ DAG), variable workload, cloud-native team।
- Bangladesh: বড় bank, telecom, AI startup-এর serious data team।
হাইব্রিড — CeleryKubernetesExecutor:
- Light task: Celery (fast)। Heavy/special: K8s (isolated)।
- ২.০+ থেকে available।
Bangladesh-এর scale-ভিত্তিক recommendation:
- Solo data engineer / startup: LocalExecutor
- Daraz/Pathao scale: CeleryExecutor + managed Postgres
- bKash/GP scale: KubernetesExecutor on EKS/AKS/GKE
- Don't want ops: MWAA (AWS) বা Cloud Composer (GCP) — managed
মূল কথা: Executor decision = trade-off ops complexity বনাম scaling। ছোট থেকে শুরু, প্রয়োজনে migrate করুন।
প্র ০২ "Idempotency" শব্দটি Airflow-এ বারবার আসে। কী মানে? কেন critical? একটি non-idempotent DAG কীভাবে fix করবেন?
Idempotency — Latin "idem" (একই) + "potent" (ক্ষমতা)। মানে: একই কাজ ১ বার বা ১০০ বার চালালেও result একই। Airflow-এ এটি pipeline reliability-র মূল ভিত্তি।
কেন critical:
- Retry: task fail হলে Airflow auto-retry। প্রতিটি attempt-এ একই side-effect হলে — duplicate data!
- Backfill: পুরাতন তারিখ পুনরায় চালাতে চান। Idempotent না হলে — সংখ্যা ভুল।
- Manual re-run: bug fix করে re-run — পুরাতন output overwrite হবে নাকি double-count?
Non-idempotent উদাহরণ (bKash style):
def increment_counter():
db.execute("UPDATE stats SET total = total + 1")
# Run 3 বার = total ৩ বার বাড়বে। ভুল!
Idempotent ভার্সন:
def set_counter_for_date(date):
db.execute(
"INSERT INTO stats (date, total) VALUES (%s, %s) "
"ON CONFLICT (date) DO UPDATE SET total = EXCLUDED.total",
(date, calculate_total(date))
)
# Run ১০০ বার = একই row, একই value
Idempotency patterns:
- Date-partitioned writes:
partitionBy("date")— re-run শুধু ওই partition overwrite। - UPSERT (MERGE): SQL-এ INSERT OR UPDATE — duplicate handle।
- Delete-before-insert:
DELETE WHERE date='X'; INSERT ...— full replace। - Content-addressable storage: hash-ভিত্তিক file name — একই content একই file।
Airflow-specific helpers:
execution_date/logical_date: task-এর "logical" তারিখ। প্রকৃত run time নয়। Backfill-এ এটি pass হয়।- Macro:
{{ ds }}= execution date — query-তে inject। প্রতি run-এ আলাদা। - External ID: Spark job-এ Airflow run_id pass — duplicate detect।
সাধারণ ভুল:
now()বাCURRENT_DATEব্যবহার — backfill ভাঙবে।{{ ds }}macro use করুন।- API call without idempotency key — payment double-charge!
- Append-only logs without dedup — multi-run-এ duplicate row।
Bangladesh fintech context: bKash-এ একটি payment confirmation API non-idempotent হলে — retry-তে ২ বার customer-এর account debit। Idempotency key (transaction_id) প্রতিটি API call-এ critical।
মূল কথা: Idempotent না হলে — Airflow চালাবেন না। প্রতিটি task design-এর সময় ভাবুন: "এটি ৫ বার চালালে result-এ পার্থক্য হবে কি?" না হলে ✓।
প্র ০৩ "DAG parsing" জিনিসটা কী? Scheduler কেন slow হয়? ১,০০০ DAG থাকলে কী optimize করবেন?
Airflow-এর scheduler প্রতি কয়েক সেকেন্ডে সব DAG file (Python) parse করে — শুধু schema বুঝতে নয়, পুরো Python execute করে! এটি ভয়াবহ slow হতে পারে।
Parsing process:
- Scheduler
~/airflow/dags/ফোল্ডার scan। - প্রতিটি .py file Python interpreter-এ load।
- DAG object ও task definition extract।
- Metadata DB-তে current state-এর সাথে compare।
- Need-to-run task queue-তে push।
সমস্যা যেখানে:
- Top-level code execution: DAG file-এ
requests.get(...), DB query — প্রতিবার parsing-এ চলে! ১,০০০ DAG × ১ সেকেন্ড = ১৬+ মিনিট/cycle। - Heavy import:
import pandas as pdtop-level — slow। - Dynamic DAG generation: for-loop-এ ১০০ DAG তৈরি করলে — সবগুলো প্রতি parse-এ।
Optimization কৌশল:
- Top-level কোড সরান: Variable, connection লোড task function-এর ভেতরে। DAG file শুধু structure।
- parse interval বাড়ান:
min_file_process_interval = 30— ৩০ সেকেন্ডে একবার। - parsing parallelism:
parsing_processes = 4— multi-core ব্যবহার। - DAG file split: ছোট ছোট ফোল্ডার-এ ভাগ।
.airflowignoreদিয়ে কিছু skip।
Airflow ২.০+ DAG Serialization:
- Webserver আর DAG file parse করে না — DB থেকে JSON serialized version পড়ে।
- UI অনেক দ্রুত।
- default-এ on।
Airflow ২.৩+ Dynamic Task Mapping:
- For-loop দিয়ে DAG তৈরি না করে —
.expand()ব্যবহার করুন। - Runtime-এ task সংখ্যা decide।
- Parsing time-এ scale handle।
Diagnostic command:
# প্রতিটি DAG file-এর parsing time
airflow dags list-import-errors
airflow dags list -o yaml | grep duration
# Webserver-এ Browse → DAG Parsing
# > ১ সেকেন্ড হলে problem
Production scale (Daraz/বড় telecom):
- ৫০০+ DAG: parse time monitor। ৩০s অতিক্রম হলে alert।
- Heavy DAG আলাদা scheduler instance-এ (Airflow ২.০+ multi-scheduler)।
- DAG বানানো-এর CI/CD-তে static check (linter)।
মূল কথা: "DAG file = pure structure declaration" — এই principle মাথায় রাখুন। Side-effect function-এর ভেতরে। তাহলে scheduler smooth।
প্র ০৪ আপনি Grameenphone-এর data team-এ যোগদান করেছেন। বর্তমানে ৩০টি cron job, কোন monitoring নেই, retry নেই। Airflow migration-এর roadmap কী হবে?
এটি প্রকৃত production scenario। ভুল migration = ৬ মাস কাজ ও ক্ষুব্ধ stakeholder। সঠিক migration = ৩ মাসে stable Airflow।
Phase 1: Inventory ও audit (১-২ সপ্তাহ)
- সব ৩০ cron job তালিকা: schedule, dependency, owner, criticality।
- প্রতিটির runtime, fail rate, last-failure (যদি jana থাকে)।
- "Hidden dependency" খুঁজুন — script-এ অন্য script-কে call?
- Stakeholder map: কোন report কে দেখে?
Phase 2: Airflow setup (২-৩ সপ্তাহ)
- Managed Airflow recommend (MWAA on AWS, Composer on GCP)। In-house ৬ মাস delay।
- PostgreSQL metadata DB।
- CeleryExecutor (GP scale-এ), ৫-১০ worker দিয়ে শুরু।
- SSO integration (LDAP/Azure AD) — security টিম-এর জন্য।
- Slack/email alert wiring।
Phase 3: Pilot migration (৪ সপ্তাহ)
- ৩-৫টি simple, low-risk DAG বাছুন।
- Cron-এর সাথে parallel চালান (dual write)।
- ২ সপ্তাহ result compare — সব মিলছে?
- Confidence build করুন।
- Decommission cron version।
Phase 4: Bulk migration (৬-৮ সপ্তাহ)
- প্রতি সপ্তাহে ৫-৭ DAG migrate।
- Pattern library তৈরি — common operator wrap।
- Idempotency add — retry safe করুন।
- Code review process — junior পেরিয়ে যান না।
- Critical DAG (CDR processing, billing) — শেষে, parallel run বেশি সময়।
Phase 5: Stabilize ও improve (চলমান)
- SLA monitor — কোন DAG সময়মতো শেষ না?
- Cost analysis — worker right-size।
- Documentation — runbook প্রতিটি DAG-এর।
- On-call rotation।
মূল কথা: "Big bang" migration-এ ভাঙবে। Iterative, parallel-run, low-risk first — এই pattern অনুসরণ করুন। Airflow আপনার team-এর জন্য — না এর বিপরীত।
অনুশীলন
-
চিনুন: এই pipeline-এ Airflow ভাল নাকি অন্য tool? কেন?
- (ক) প্রতি ১ সেকেন্ডে IoT sensor data → fraud detection।
- (খ) প্রতি রাত ৩টায় warehouse → 5 BI dashboard refresh।
- (গ) Customer signup-এ welcome email + SMS।
- (ক) Airflow নয় — sub-second latency। Kafka + Flink/Spark Streaming ব্যবহার করুন।
- (খ) Airflow আদর্শ — daily batch, multi-step dependency, retry/alert সব দরকার।
- (গ) Airflow নয় — event-driven, low latency। Application-level event queue (RabbitMQ/SQS + lambda) ভাল।
-
লিখুন: এই pipeline DAG-এ — extract, validate, load, notify। কোন
>>structure?extract >> validate >> load >> notify # অথবা যদি validation fail হলে অন্য branch: extract >> validate validate >> [load_success_branch, alert_failure_branch] load_success_branch >> notifySequential chain সবচেয়ে সাধারণ। Branch-এর জন্য
BranchPythonOperatorবাShortCircuitOperator। -
চিন্তা করুন: bKash-এর reconciliation DAG-এ
retry=3, retry_delay=5min। কোন ক্ষেত্রে retry safe, কোন ক্ষেত্রে dangerous?Safe (idempotent):
- S3 থেকে file download — একই file বারবার download safe।
- UPSERT into table by date — partition overwrite।
- Sending alert email (acceptable double-email)।
Dangerous (non-idempotent):
UPDATE balance SET amount = amount + X— retry মানে double credit!- API call to send money without idempotency key — double payment।
INSERT INTO logwithout dedup — duplicate row।
সমাধান: SQL-এ
ON CONFLICT, API-এ idempotency key, file-এ atomic rename।
আরও পড়ুন · ABCL TECH-এ আপনার পরবর্তী পদক্ষেপ
- পাঠ ১৪ · DAG লেখা ও schedule পরবর্তী পাঠ এবার সত্যিকারের DAG লিখুন — decorator, XCom, idempotency সহ।
- পাঠ ১২ · Spark optimization আগের পাঠ Airflow-এর অনেক DAG Spark কাজ orchestrate করে।
- পাঠ ১৫ · dbt — analytics engineering এই পাঠের সাথে সম্পর্কিত Airflow + dbt একসাথে — modern data stack-এর হৃদয়।
- সব AI Courses দেখুন ABCL TECH Python, ML, DL, NLP, CV, GenAI, RL, MLOps — সব AI কোর্স একসাথে।
pip install apache-airflow অথবা
Google Colab
ব্যবহার করুন — সরাসরি web থেকে Python চালান।