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

Apache Airflow পরিচিতি

Apache Airflow — workflow orchestration intro
৮ মিনিট পড়া মাঝারি · Intermediate DAG-ভিত্তিক

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

  • Workflow orchestration কী — কেন cron যথেষ্ট নয়
  • DAG, task, operator, sensor — Airflow-এর ভিত্তি শব্দ
  • Airflow architecture — চারটি component কীভাবে একসাথে কাজ করে
  • কখন Airflow ব্যবহার করবেন, কখন অন্য tool

১ · Workflow orchestration কী — সমস্যাটা কী

কল্পনা করুন আপনি Daraz-এর data team-এ। প্রতি রাতে এই কাজগুলো হতে হবে:

  1. MySQL থেকে গতকালের order data Snowflake-এ copy
  2. সেই data clean — duplicate বাদ, type fix
  3. Customer dimension table-এর সাথে JOIN
  4. Daily revenue, top products report তৈরি
  5. Email-এ executive-দের পাঠাও, Slack-এ team-কে

প্রতিটি কাজের dependency আছে। ৩ চলবে শুধু ১ ও ২ শেষ হলে। ৪ শেষ হওয়ার আগে ৫ চললে — empty email।

cron-এ কেন হয় না

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-কে ভাবুন bKash-এর daily settlement-এর foreman হিসেবে। তার একটি checklist: "প্রথমে A bank, তারপর B bank থেকে statement আনো; দু'টোই এলে reconcile করো; reconcile শেষে BTRC-কে report পাঠাও।" Foreman সময়মতো শুরু করে, কোনো step আটকে গেলে retry করে, না হলে manager-কে SMS দেয়।

৩ · 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।
Airflow architecture — চার component webserver · scheduler · executor · metadata DB 📁 DAG files Python (.py) ⏰ Scheduler parse DAGs · queue tasks ⚙️ Executor Celery / K8s / Local workers 🖥️ Webserver Flask UI · port 8080 🗄️ Metadata DB PostgreSQL / MySQL 👷 Workers actual task execution read/write state read heartbeat solid arrow = control flow · dashed = state read/write All four components must be running for Airflow to work
Airflow-এর চার component — Scheduler পরিকল্পনা, Executor চালায়, Metadata DB সব মনে রাখে, Webserver আপনাকে দেখায়।

৪ · 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।
Sensor default-এ slot dখল করে রাখে — অনেক sensor একসাথে চললে worker pool exhausted। Airflow ২.০+ থেকে mode="reschedule" ব্যবহার করুন — sensor task slot ছেড়ে wait time-এ ঘুমিয়ে যায়।

৫ · cron বনাম Airflow — তুলনা

সাধারণ ৩-task pipeline cron-এ:

crontab — ভঙ্গুর সমাধান
# /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 অপাঠ্য

    
cron সরল কাজে ভালো। কিন্তু ৩+ dependent task হলেই production-এ ভাঙে। প্রতিটি data team এক সময় cron থেকে Airflow-এ আসে — এটি rite of passage।

৬ · একই pipeline Airflow-এ

Python · Airflow DAG
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।
Bangladesh-এর context-এ: bKash-এর daily reconciliation, Daraz-এর order analytics, Grameenphone-এর CDR processing — সবই Airflow-এর জন্য আদর্শ। Pathao-এর live ETA prediction — Airflow নয়, Kafka + Flink।

৮ · 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-এ

bash · Airflow setup
# দ্রুত 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

    
Production-এ SQLite + SequentialExecutor কাজ করবে না — PostgreSQL + CeleryExecutor (বা KubernetesExecutor) দরকার। কিন্তু শেখার জন্য local SQLite যথেষ্ট।
Production-এ কখনো SQLite metadata DB ব্যবহার করবেন না — concurrency সমর্থন নেই। শুরু থেকেই PostgreSQL দিন। Astronomer, MWAA, Composer — managed Airflow এই সব auto-handle করে।

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

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

প্র ০১ 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:

  1. Scheduler ~/airflow/dags/ ফোল্ডার scan।
  2. প্রতিটি .py file Python interpreter-এ load।
  3. DAG object ও task definition extract।
  4. Metadata DB-তে current state-এর সাথে compare।
  5. Need-to-run task queue-তে push।

সমস্যা যেখানে:

  • Top-level code execution: DAG file-এ requests.get(...), DB query — প্রতিবার parsing-এ চলে! ১,০০০ DAG × ১ সেকেন্ড = ১৬+ মিনিট/cycle।
  • Heavy import: import pandas as pd top-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-এর জন্য — না এর বিপরীত।

অনুশীলন

  1. চিনুন: এই 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) ভাল।
  2. লিখুন: এই 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 >> notify

    Sequential chain সবচেয়ে সাধারণ। Branch-এর জন্য BranchPythonOperator বা ShortCircuitOperator।

  3. চিন্তা করুন: 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 log without dedup — duplicate row।

    সমাধান: SQL-এ ON CONFLICT, API-এ idempotency key, file-এ atomic rename।

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

Airflow চেষ্টা করতে চান? Local-এ pip install apache-airflow অথবা Google Colab ব্যবহার করুন — সরাসরি web থেকে Python চালান।
পূর্ববর্তী পাঠ
পাঠ ১২ · Spark optimization