অধ্যায় 5 · পাইপলাইন ও AI
ETL ও ELT: নির্ভরযোগ্য ডেটা পাইপলাইন বানানো
- পৃষ্ঠা 17 / 22
- 20 মিনিট পড়া
প্রতিটা ড্যাশবোর্ড, প্রতিটা ট্রেন করা মডেল আর প্রতিটা RAG ইনডেক্সের পেছনে থাকে একটা ডেটা পাইপলাইন (data pipeline), যা এদের ডেটা জোগায়: এমন কোড, যা নিয়মিত ডেটা যেখানে তৈরি হয় সেখান থেকে তুলে আনে, দরকারমতো রূপ বদলায়, আর যেখানে ব্যবহার হবে সেখানে পৌঁছে দেয়। পাইপলাইন নির্ভরযোগ্য না হলে শুরুতে কেউ টের পায় না। তারপর একদিন আয়ের চার্টে এমন একটা পতন দেখা যায়, যা আসলে ঘটেইনি, বা আধা দিনের ডুপ্লিকেট সারি নিয়ে মডেল আবার ট্রেন হয়ে যায়। ভালো পাইপলাইন একঘেয়ে: যেকোনো ব্যর্থতার পর আবার চালানো যায়, কিছুই দুবার গোনে না, আর কোনো গোলমাল হলে জোরেশোরে জানিয়ে দেয়।
এই পাতায় এই অধ্যায়ে শেখা জিনিসগুলো জুড়ে এমন একটা পাইপলাইন বানানো হবে। এটা ডেটা তোলে তিন জায়গা থেকে: shop ডেটাবেস, একটা কুরিয়ার "API" (আসল API-র বদলে একটা লোকাল JSON ফাইল) আর টার্গেটের একটা CSV। তারপর সারিগুলো রূপান্তর আর যাচাই করে আলাদা একটা ওয়্যারহাউস ডেটাবেসে লোড করে। আর এটা আইডেমপোটেন্ট (idempotent): দুবার চালালেও দ্বিতীয়বারে কিছুই বদলায় না।
যা শিখবেন
- পাইপলাইন কী, ETL বনাম ELT, আর ওয়াটারমার্ক দিয়ে পূর্ণ বনাম ইনক্রিমেন্টাল লোড (আর কেন একটা লুকব্যাক উইন্ডো লাগে)।
- আপসার্ট (upsert) দিয়ে আইডেমপোটেন্ট লোড, আর দুবার চালিয়ে সেটা প্রমাণ করা।
- যাচাই আর রিকনসিলিয়েশন (সারির সংখ্যা, যোগফল), লগিং, আর অনির্ভরযোগ্য উৎসের জন্য ব্যাকঅফ দিয়ে রিট্রাই (retry)।
- শিডিউলিং (cron, একটা Airflow DAG), মনিটরিং আর অ্যালার্ট, আর পাইপলাইন টেস্ট করা।
- CDC: change data capture কী, আর কখন দরকার হয়।
ETL আর ELT
দুটোই ডেটা উৎস (source) থেকে একটা টার্গেটে (ওয়্যারহাউস, লেক বা ফিচার স্টোর) নিয়ে যায়। পার্থক্য হলো রূপান্তর (transform) কোথায় হয়:
ETL: উৎস ──extract──► [ Python / Spark-এ transform ] ──load──► ওয়্যারহাউস (শুধু পরিষ্কার টেবিল)
ELT: উৎস ──extract──► ওয়্যারহাউস (কাঁচা টেবিল) ──ওয়্যারহাউসের ভেতরে SQL দিয়ে transform──► পরিষ্কার টেবিল| ETL | ELT | |
|---|---|---|
| রূপান্তর চলে | পাইপলাইনের কোডে (Python, Spark) | ওয়্যারহাউসে, SQL হিসেবে (প্রায়ই dbt দিয়ে সাজানো) |
| কাঁচা ডেটা রাখা হয়? | সাধারণত না | হ্যাঁ, তাই লজিক বদলালে আবার রূপান্তর করা যায় |
| কোথায় ভালো | যে কাজ SQL-এ ভালো হয় না (পার্সিং, API ডাকা, ML স্কোরিং), আর এমন সংবেদনশীল ডেটা, যা ওয়্যারহাউসে পৌঁছানোর আগেই মাস্ক করতে হয় | ক্লাউড ওয়্যারহাউসে (BigQuery, Snowflake) বিপুল ডেটা, যেখানে SQL চালানোর খরচ কম আর সহজেই বাড়ানো যায় |
বাস্তবের পাইপলাইনে দুটোই মেশানো থাকে, আমাদেরটাতেও: অর্ডার লাইনের সাথে API-র ডেটা জোড়া হয় Python-এ (ETL), তারপর মাসিক সারাংশ বানায় ওয়্যারহাউসের ভেতরের SQL (ELT)।
উৎসগুলো
প্রথম উৎস shop ডেটাবেস। বাকি দুটো এখানেই লিখে নেওয়া হচ্ছে, যাতে পাতাটা বাইরের কিছুর ওপর নির্ভর না করে: JSON হিসেবে সেভ করা কুরিয়ার API-র উত্তর, আর সেলস টিমের দেওয়া মাসিক টার্গেটের একটা ফাইল:
import json
deliveries = [
{"order_id": 1, "courier": "Pathao", "delivered_on": "2026-01-14"},
{"order_id": 2, "courier": "RedX", "delivered_on": "2026-01-23"},
{"order_id": 3, "courier": "Pathao", "delivered_on": "2026-02-05"},
{"order_id": 4, "courier": "Steadfast", "delivered_on": "2026-02-16"},
{"order_id": 6, "courier": "RedX", "delivered_on": "2026-03-11"},
{"order_id": 7, "courier": "Pathao", "delivered_on": "2026-03-17"},
{"order_id": 8, "courier": "Steadfast", "delivered_on": "2026-03-30"},
{"order_id": 9, "courier": "RedX", "delivered_on": "2026-04-06"},
{"order_id": 10, "courier": "Pathao", "delivered_on": "2026-04-21"},
{"order_id": 11, "courier": "Steadfast", "delivered_on": "2026-05-08"},
{"order_id": 12, "courier": "RedX", "delivered_on": None},
{"order_id": 14, "courier": "Pathao", "delivered_on": "2026-06-12"},
]
with open("courier_api.json", "w", encoding="utf-8") as f:
json.dump({"deliveries": deliveries}, f, indent=1)
with open("targets.csv", "w", encoding="utf-8") as f:
f.write("month,target\n2026-01,8000\n2026-02,8000\n2026-03,20000\n"
"2026-04,20000\n2026-05,15000\n2026-06,20000\n")
print("sources ready")sources readyপাইপলাইন বানানো
টার্গেট: স্টেটসহ একটা ওয়্যারহাউস
টার্গেট হলো আলাদা একটা SQLite ফাইল, warehouse.db। ডেটার টেবিল ছাড়াও এতে হিসাব রাখার দুটো টেবিল আছে: etl_state মনে রাখে শেষ রান কতদূর পৌঁছেছিল (ওয়াটারমার্ক, watermark), আর etl_runs মনিটরিংয়ের জন্য প্রতিটা রানের রেকর্ড রাখে। fact টেবিলে প্রতিটা অর্ডার লাইনের জন্য একটা সারি, আর এর প্রাইমারি কি-ই আইডেমপোটেন্ট লোড সম্ভব করে:
import csv
import logging
import sqlite3
import sys
import time
from datetime import date, timedelta
log = logging.getLogger("etl")
WAREHOUSE_DDL = """
CREATE TABLE IF NOT EXISTS fact_order_lines (
order_id INTEGER NOT NULL,
product_id INTEGER NOT NULL,
order_date TEXT NOT NULL,
month TEXT NOT NULL,
category TEXT NOT NULL,
status TEXT NOT NULL,
quantity INTEGER NOT NULL,
revenue INTEGER NOT NULL,
courier TEXT,
delivery_days INTEGER,
PRIMARY KEY (order_id, product_id)
);
CREATE TABLE IF NOT EXISTS monthly_sales (
month TEXT PRIMARY KEY, orders INTEGER NOT NULL, revenue INTEGER NOT NULL,
target INTEGER, pct_of_target REAL
);
CREATE TABLE IF NOT EXISTS etl_state (key TEXT PRIMARY KEY, value TEXT NOT NULL);
CREATE TABLE IF NOT EXISTS etl_runs (
run_id INTEGER PRIMARY KEY, mode TEXT, since TEXT, extracted INTEGER,
inserted INTEGER, updated INTEGER, status TEXT
);
"""এক্সট্র্যাক্ট (Extract)
তিনটা এক্সট্র্যাক্টর, প্রতিটা উৎসের জন্য একটা। SQL এক্সট্র্যাক্টর জয়েন আর ফিল্টারের কাজটা উৎস ডেটাবেসকে দিয়েই করিয়ে নেয়, আর সাধারণ ডিকশনারি ফেরত দেয়। API এক্সট্র্যাক্টরকে রিট্রাই দিয়ে মুড়ে রাখা হয়েছে, কারণ নেটওয়ার্ক মাঝে মাঝে ব্যর্থ হয়। বদলি ক্লাসটা ইচ্ছে করেই প্রথম কয়েকটা কলে ব্যর্থ হয়, যেমন আসল API মাঝে মাঝে টাইমআউট হয়:
def extract_order_lines(src, since):
"""`since` (YYYY-MM-DD) বা তার পরে দেওয়া অর্ডারের লাইনগুলো।"""
cur = src.execute("""
SELECT oi.order_id, oi.product_id, o.order_date, o.status,
COALESCE(c.name, 'Uncategorised') AS category,
oi.quantity, oi.quantity * oi.unit_price AS revenue
FROM order_items oi
JOIN orders o ON o.order_id = oi.order_id
JOIN products p ON p.product_id = oi.product_id
LEFT JOIN categories c ON c.category_id = p.category_id
WHERE o.order_date >= ?
ORDER BY oi.order_id, oi.product_id""", (since,))
cols = [d[0] for d in cur.description]
return [dict(zip(cols, row)) for row in cur]
class CourierAPI:
"""একটা HTTP API-র বদলি। আসল কোডে: requests.get(url, timeout=10).json()।"""
def __init__(self, path, fail_first=0):
self.path, self.fail_first = path, fail_first
def get_deliveries(self):
if self.fail_first > 0:
self.fail_first -= 1
raise ConnectionError("courier API timed out")
with open(self.path, encoding="utf-8") as f:
return json.load(f)["deliveries"]
def with_retries(fn, attempts=4, base_delay=0.01):
"""fn() কল করে; সাময়িক এররে base_delay-এর 1, 2, 4 ... গুণ অপেক্ষা করে আবার চেষ্টা করে।"""
for attempt in range(1, attempts + 1):
try:
return fn()
except ConnectionError as exc: # সাময়িক এরর: রিট্রাই করার মতো
if attempt == attempts:
raise
delay = base_delay * 2 ** (attempt - 1)
log.warning("attempt %d failed (%s), retrying in %.2fs", attempt, exc, delay)
time.sleep(delay)
def extract_targets(path):
with open(path, newline="", encoding="utf-8") as f:
return [{"month": r["month"], "target": int(r["target"])} for r in csv.DictReader(f)]রিট্রাই হয় শুধু ConnectionError-এ। ভুল API key বা ভাঙা ডেটা প্রতিবার একইভাবে ব্যর্থ হয়; সেটা রিট্রাই করলে শুধু অ্যালার্ট দেরিতে আসে। আসল কোডের জন্য দুটো কথা। requests ব্যবহার করলে requests.exceptions.ConnectionError আর requests.exceptions.Timeout ধরুন: এগুলো Python-এর বিল্ট-ইন ConnectionError-এর সাবক্লাস নয়। আর অপেক্ষার সময়ে জিটার (jitter), মানে একটা এলোমেলো অংশ যোগ করুন (যেমন delay * random.uniform(0.5, 1.5)), যাতে একসাথে ব্যর্থ হওয়া একশোটা ওয়ার্কার ঠিক একই মুহূর্তে আবার চেষ্টা করে API-কে আবার বসিয়ে না দেয়। tenacity-র মতো লাইব্রেরিতে এসব তৈরি করাই থাকে। (এখানে অপেক্ষার সময় খুব ছোট আর নির্দিষ্ট, যাতে পাতার আউটপুট প্রতিবার একই থাকে।)
ট্রান্সফর্ম আর যাচাই
ট্রান্সফর্ম ধাপ প্রতিটা লাইনে মাস যোগ করে, আর লাইনটাকে তার ডেলিভারির সাথে জোড়ে (কোন কুরিয়ার, কত দিন লেগেছে)। তারপর কিছু লেখার আগেই যাচাই ধাপ সারিগুলোকে নিয়মের সাথে মিলিয়ে দেখে। কোনো নিয়ম ভাঙলে রান থেমে যায়, ওয়্যারহাউসে হাতই পড়ে না:
def transform(lines, deliveries):
by_order = {d["order_id"]: d for d in deliveries}
rows = []
for line in lines:
d = by_order.get(line["order_id"], {})
days = None
if d.get("delivered_on"):
days = (date.fromisoformat(d["delivered_on"]) - date.fromisoformat(line["order_date"])).days
rows.append({**line, "month": line["order_date"][:7],
"courier": d.get("courier"), "delivery_days": days})
return rows
def validate(rows):
problems = []
for r in rows:
if r["quantity"] <= 0 or r["revenue"] <= 0:
problems.append(f"order {r['order_id']}: non-positive quantity or revenue")
if r["delivery_days"] is not None and not 0 <= r["delivery_days"] <= 30:
problems.append(f"order {r['order_id']}: delivery took {r['delivery_days']} days")
if problems:
raise ValueError("validation failed: " + "; ".join(problems))লোড: আইডেমপোটেন্ট আপসার্ট
আইডেমপোটেন্সি টিকবে কি না, তা ঠিক হয় লোড ধাপে। সাধারণ INSERT দুবার চালালে হয় প্রাইমারি কি-তে আটকে ব্যর্থ হয়, নয়তো (কি না থাকলে) সবকিছু ডুপ্লিকেট হয়ে যায়। আপসার্ট (upsert) নতুন লাইন ঢোকায় আর আগের লাইন হালনাগাদ করে (ON CONFLICT … DO UPDATE লেখার নিয়ম আছে INSERT, UPDATE, DELETE আর আপসার্ট, নিরাপদে পাতায়)। আগে থেকে থাকা সারি বাদ দিয়ে শুধু নতুন সারি ঢোকানো, যেমন ইমপোর্ট ও এক্সপোর্ট: CSV, JSON, Excel আর Parquet পাতার NOT EXISTS লোড, সেটাও আইডেমপোটেন্ট; কিন্তু আগে লোড হওয়া সারির কোনো পরিবর্তন সে কখনো ধরে না। এখানে অর্ডারের status বদলায়, তাই হালনাগাদ লাগবেই। WHERE … IS NOT … অংশটা কোনো সারি হালনাগাদ করে শুধু তখন, যখন বদলাতে পারে এমন কোনো কলাম (status, পরিমাণ, আয়, কুরিয়ার, ডেলিভারির দিন) সত্যিই আলাদা। তাই অপরিবর্তিত সারি আবার লেখাই হয় না, আর পরিবর্তনের সংখ্যা গুনলেই ঠিক বোঝা যায় রানটা কী করল। (<>-এর বদলে IS NOT, যাতে NULL কুরিয়ারের তুলনাও ঠিক হয়।)
UPSERT_FACT = """
INSERT INTO fact_order_lines (order_id, product_id, order_date, month, category, status,
quantity, revenue, courier, delivery_days)
VALUES (:order_id, :product_id, :order_date, :month, :category, :status,
:quantity, :revenue, :courier, :delivery_days)
ON CONFLICT (order_id, product_id) DO UPDATE SET
order_date = excluded.order_date, month = excluded.month, category = excluded.category,
status = excluded.status, quantity = excluded.quantity, revenue = excluded.revenue,
courier = excluded.courier, delivery_days = excluded.delivery_days
WHERE (fact_order_lines.status, fact_order_lines.quantity, fact_order_lines.revenue,
fact_order_lines.courier, fact_order_lines.delivery_days)
IS NOT (excluded.status, excluded.quantity, excluded.revenue,
excluded.courier, excluded.delivery_days)
"""
BUILD_MONTHLY = """
INSERT INTO monthly_sales (month, orders, revenue, target, pct_of_target)
SELECT f.month, COUNT(DISTINCT f.order_id), SUM(f.revenue), t.target,
ROUND(100.0 * SUM(f.revenue) / t.target, 1)
FROM fact_order_lines f LEFT JOIN stg_targets t ON t.month = f.month
WHERE f.status <> 'cancelled'
GROUP BY f.month
ON CONFLICT (month) DO UPDATE SET
orders = excluded.orders, revenue = excluded.revenue,
target = excluded.target, pct_of_target = excluded.pct_of_target
WHERE (monthly_sales.orders, monthly_sales.revenue, monthly_sales.target)
IS NOT (excluded.orders, excluded.revenue, excluded.target)
"""
def count(con, table):
return con.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0]
def load(wh, rows, targets):
"""fact সারিগুলো upsert করে, তারপর SQL দিয়ে মাসিক সারাংশ নতুন করে বানায় (ELT ধাপ)।"""
before, changes = count(wh, "fact_order_lines"), wh.total_changes
wh.executemany(UPSERT_FACT, rows)
inserted = count(wh, "fact_order_lines") - before
updated = wh.total_changes - changes - inserted
wh.execute("CREATE TEMP TABLE IF NOT EXISTS stg_targets (month TEXT PRIMARY KEY, target INTEGER)")
wh.execute("DELETE FROM stg_targets")
wh.executemany("INSERT INTO stg_targets VALUES (:month, :target)", targets)
changes = wh.total_changes
wh.execute(BUILD_MONTHLY)
return inserted, updated, wh.total_changes - changesরিকনসাইল: মোট হিসাব মেলানো
যাচাই দেখে প্রতিটা সারি; রিকনসিলিয়েশন (reconciliation) দেখে মোট হিসাব: টার্গেটের সাথে উৎস মিলছে কি না। এখানে shop-এর প্রতিটা অর্ডার লাইন ওয়্যারহাউসে থাকতে হবে, আর মোট আয়ও একই হতে হবে। এটা সবচেয়ে সস্তা টেস্ট, অথচ এক ধাক্কায় পুরো এক ধরনের বাগ ধরে ফেলে: যে জয়েন সারি হারায়, যে ফিল্টার বেশি কড়া, যে লোড আধাআধি ব্যর্থ হয়েছে।
def reconcile(src, wh):
source = src.execute("SELECT COUNT(*), SUM(quantity * unit_price) FROM order_items").fetchone()
target = wh.execute("SELECT COUNT(*), SUM(revenue) FROM fact_order_lines").fetchone()
if source != target:
raise RuntimeError(f"reconciliation failed: source {source} vs warehouse {target}")
return targetরান: পূর্ণ বা ইনক্রিমেন্টাল, ওয়াটারমার্কসহ
পূর্ণ লোড (full load) প্রতিবার সবকিছু তোলে: সহজ আর সবসময় সঠিক, কিন্তু উৎসে কয়েক বছরের ডেটা থাকলে ধীর। ইনক্রিমেন্টাল লোড (incremental load) তোলে শুধু শেষ রানের পরে যা নতুন এসেছে; কতদূর লোড হয়েছে তা মনে রাখা হয় একটা ওয়াটারমার্ক দিয়ে (এখানে লোড হওয়া সবচেয়ে নতুন অর্ডারের তারিখ)। সমস্যা হলো: যে সারি বদলেছে কিন্তু ওয়াটারমার্কের চেয়ে পুরোনো, সেটা বাদ পড়ে যায়। অর্ডার ১৩ এখন pending; পরের সপ্তাহে শিপ হলেও এর তারিখ বদলাবে না। তাই এক্সট্র্যাক্টর ওয়াটারমার্ক থেকে একটা লুকব্যাক উইন্ডো (lookback window) পরিমাণ পিছিয়ে পড়া শুরু করে। কয়েক দিনের ডেটা দুবার পড়লে কোনো ক্ষতি নেই, ঠিক এই কারণেই যে লোডটা আইডেমপোটেন্ট।
তুলনাটাও খেয়াল করুন: order_date >= ?, > নয়। ওয়াটারমার্ক হলো এ পর্যন্ত দেখা সবচেয়ে বড় মান, আর ঠিক ওই মানেরই আরও সারি পরে আসতে পারে (একই দিনে আরেকটা অর্ডার, বা টাইমস্ট্যাম্প হলে একই সেকেন্ডে)। > দিলে সেগুলো চিরকালের জন্য বাদ পড়ত; >= দিলে সীমানার মানটা আবার পড়া হয়, আর আপসার্টের কারণে তাতে কোনো ক্ষতি হয় না।
LOOKBACK_DAYS = 14
def run_pipeline(api, src_path="shop.db", wh_path="warehouse.db", full=False):
logging.basicConfig(level=logging.INFO, format="%(levelname)s %(name)s: %(message)s",
stream=sys.stdout, force=True)
src, wh = sqlite3.connect(src_path), sqlite3.connect(wh_path)
mode = None
try:
wh.executescript(WAREHOUSE_DDL)
state = wh.execute("SELECT value FROM etl_state WHERE key = 'orders_watermark'").fetchone()
mode = "full" if full or state is None else "incremental"
since = "0001-01-01" if mode == "full" else \
(date.fromisoformat(state[0]) - timedelta(days=LOOKBACK_DAYS)).isoformat()
lines = extract_order_lines(src, since)
deliveries = with_retries(api.get_deliveries)
targets = extract_targets("targets.csv")
log.info("%s extract since %s: %d order lines, %d deliveries, %d targets",
mode, since, len(lines), len(deliveries), len(targets))
rows = transform(lines, deliveries)
validate(rows)
with wh: # একটা ট্রানজ্যাকশন: হয় সবটা, নয় কিছুই না
inserted, updated, monthly = load(wh, rows, targets)
if rows:
wh.execute("INSERT INTO etl_state VALUES ('orders_watermark', ?) "
"ON CONFLICT (key) DO UPDATE SET value = excluded.value",
(max(r["order_date"] for r in rows),))
wh.execute("INSERT INTO etl_runs (mode, since, extracted, inserted, updated, status) "
"VALUES (?, ?, ?, ?, ?, 'ok')", (mode, since, len(rows), inserted, updated))
n, revenue = reconcile(src, wh)
log.info("load: %d inserted, %d updated, %d monthly rows changed; warehouse %d lines, %d taka",
inserted, updated, monthly, n, revenue)
return inserted, updated, monthly
except Exception:
log.exception("run failed")
with wh:
wh.execute("INSERT INTO etl_runs (mode, status) VALUES (?, 'failed')", (mode,))
raise
finally:
src.close()
wh.close()লোডের চারপাশের ট্রানজ্যাকশনটা খেয়াল করুন। প্রসেস মাঝপথে বন্ধ হয়ে গেলে ওয়্যারহাউসে আগের সম্পূর্ণ অবস্থাই থেকে যায়, ওয়াটারমার্কও সরে না, তাই পরের রান কাজটা আবার করে নেয়।
দুবার চালানো
প্রথম রান কোনো ওয়াটারমার্ক পায় না, তাই পূর্ণ লোড করে। উত্তর দেওয়ার আগে কুরিয়ার API দুবার টাইমআউট হয়:
run_pipeline(CourierAPI("courier_api.json", fail_first=2))WARNING etl: attempt 1 failed (courier API timed out), retrying in 0.01s
WARNING etl: attempt 2 failed (courier API timed out), retrying in 0.02s
INFO etl: full extract since 0001-01-01: 20 order lines, 12 deliveries, 6 targets
INFO etl: load: 20 inserted, 0 updated, 6 monthly rows changed; warehouse 20 lines, 103550 takaএবার আসল পরীক্ষা। ওয়্যারহাউসের টেবিলগুলোর একটা ফিঙ্গারপ্রিন্ট নিন, পাইপলাইন আবার চালান, তারপর দুটো মিলিয়ে দেখুন:
import hashlib
def fingerprint(path="warehouse.db"):
con = sqlite3.connect(path)
rows = con.execute("SELECT * FROM fact_order_lines ORDER BY order_id, product_id").fetchall()
rows += con.execute("SELECT * FROM monthly_sales ORDER BY month").fetchall()
con.close()
return hashlib.sha256(repr(rows).encode()).hexdigest()[:12]
before = fingerprint()
run_pipeline(CourierAPI("courier_api.json"))
print("unchanged:", fingerprint() == before)INFO etl: incremental extract since 2026-05-26: 3 order lines, 12 deliveries, 6 targets
INFO etl: load: 0 inserted, 0 updated, 0 monthly rows changed; warehouse 20 lines, 103550 taka
unchanged: Trueদ্বিতীয় রানটা ছিল ইনক্রিমেন্টাল: ওয়াটারমার্ক থেকে ১৪ দিন পিছিয়ে অর্ডার ১৩ আর ১৪-এর ৩টা লাইন আবার পড়েছে, কোনো পার্থক্য পায়নি, তাই কিছুই লেখেনি। ক্র্যাশের পর আবার চালানো, শিডিউলারের রিট্রাই, বা কেউ ভুল করে দুবার "run" চাপা, কোনোটাতেই এখন আর ভয় নেই। ওয়্যারহাউসে ফলাফল দাঁড়াল এমন:
wh = sqlite3.connect("warehouse.db")
for row in wh.execute("SELECT * FROM monthly_sales ORDER BY month"):
print(row)
print(wh.execute("""SELECT courier, COUNT(DISTINCT order_id), ROUND(AVG(delivery_days), 1)
FROM fact_order_lines WHERE courier IS NOT NULL
GROUP BY courier ORDER BY courier""").fetchall())
wh.close()('2026-01', 2, 6600, 8000, 82.5)
('2026-02', 2, 7000, 8000, 87.5)
('2026-03', 3, 21400, 20000, 107.0)
('2026-04', 2, 23450, 20000, 117.3)
('2026-05', 2, 8500, 15000, 56.7)
('2026-06', 2, 18100, 20000, 90.5)
[('Pathao', 5, 2.3), ('RedX', 4, 3.5), ('Steadfast', 3, 2.0)]উৎস যখন বদলায়
এক সপ্তাহ পরে একটা নতুন অর্ডার আসে, আর pending অর্ডার ১৩ শিপ হয়। shop ডেটাবেসে পরিবর্তনগুলো এমন:
INSERT INTO orders (order_id, customer_id, order_date, status) VALUES (15, 8, '2026-06-20', 'pending');
INSERT INTO order_items (order_id, product_id, quantity, unit_price) VALUES (15, 1, 1, 1200), (15, 4, 2, 900);
UPDATE orders SET status = 'shipped' WHERE order_id = 13;
SELECT order_id, order_date, status FROM orders WHERE order_id >= 13 ORDER BY order_id;+----------+------------+-----------+
| order_id | order_date | status |
+----------+------------+-----------+
| 13 | 2026-06-02 | shipped |
| 14 | 2026-06-09 | delivered |
| 15 | 2026-06-20 | pending |
+----------+------------+-----------+অর্ডার ১৩-এর তারিখ ২ জুন ২০২৬, ওয়াটারমার্কের (৯ জুন) এক সপ্তাহ আগে। যে এক্সট্র্যাক্টর শুধু ওয়াটারমার্কের পরের সারি পড়ে, সে এই পরিবর্তন কখনোই দেখতে পেত না:
src = sqlite3.connect("shop.db")
newer_only = extract_order_lines(src, "2026-06-10")
print("without lookback:", sorted({r["order_id"] for r in newer_only}))
src.close()
run_pipeline(CourierAPI("courier_api.json"))without lookback: [15]
INFO etl: incremental extract since 2026-05-26: 5 order lines, 12 deliveries, 6 targets
INFO etl: load: 2 inserted, 1 updated, 1 monthly rows changed; warehouse 22 lines, 106550 takaলুকব্যাক উইন্ডো থাকায় রানটা দুটো পরিবর্তনই ধরেছে: অর্ডার ১৫-এর দুটো নতুন লাইন আর অর্ডার ১৩-এর status বদল, সাথে জুনের সারাংশও হালনাগাদ হয়েছে। রিকনসিলিয়েশন এখনো মিলে যাচ্ছে।
লুকব্যাক উইন্ডো একটা আন্দাজনির্ভর কৌশল: অর্ডার দেওয়ার ২০ দিন পরে কিছু বদলালে সেটা তবুও বাদ পড়বে। নির্ভরযোগ্য ইনক্রিমেন্টাল লোডের জন্য উৎসকেই জানাতে হবে কী বদলেছে। হয় প্রতিটা লেখার সময় হালনাগাদ হওয়া একটা
updated_atকলাম, যা ওয়াটারমার্ক হিসেবে ব্যবহার হবে (একটা ট্রিগার দিয়ে এটা রাখা যায়, দেখুন ভিউ, স্টোরড প্রসিডিওর, ফাংশন আর ট্রিগার), নয়তো চেঞ্জ ডেটা ক্যাপচার (change data capture, CDC): Debezium-এর মতো টুল ডেটাবেসের নিজস্ব পরিবর্তনের লগ (MySQL-এর binlog, PostgreSQL-এর WAL) পড়ে প্রতিটা insert, update আর delete স্ট্রিম করে পাঠায়। CDC delete-ও ধরে: মুছে ফেলা সারি টেবিলেই নেই, তাই ওয়াটারমার্ক কোয়েরি তাকে খুঁজে পায় না। ব্যতিক্রম শুধু তখন, যখন উৎস সারি না মুছে কেবল "মুছে ফেলা হয়েছে" বলে একটা ফ্ল্যাগ (soft delete) বসায়।
প্রোডাকশনে চালানো
শিডিউলিং
সবচেয়ে সহজ শিডিউলার হলো cron। এই লাইনটা প্রতি রাত ২টা ১৫ মিনিটে পাইপলাইন চালায় আর লগটা একটা ফাইলের শেষে যোগ করে। পাঁচটা ঘর হলো মিনিট, ঘণ্টা, মাসের দিন, মাস আর সপ্তাহের দিন; * মানে "প্রতিটা"। cron চলে সার্ভারের টাইম জোনে, আর ক্লাউড সার্ভারে সেটা প্রায়ই UTC, বাংলাদেশের সময় নয়:
# m h dom mon dow command
15 2 * * * cd /srv/shop-etl && .venv/bin/python run_etl.py >> logs/etl.log 2>&1পাইপলাইনগুলো যখন একটা আরেকটার ওপর নির্ভর করে (অর্ডার লোড, তারপর ফিচার নতুন করে বানানো, তারপর মডেল আবার ট্রেন), তখন একটা অর্কেস্ট্রেটর (orchestrator) ব্যবহার করুন: Apache Airflow, Dagster বা Prefect। এরা টাস্কগুলো ঠিক ক্রমে চালায়, ব্যর্থ টাস্ক রিট্রাই করে, প্রতিটা রানের ইতিহাস রাখে আর ব্যর্থ হলে অ্যালার্ট পাঠায়। এই পাইপলাইনের জন্য একটা Airflow DAG (টাস্কের একটা গ্রাফ) দেখতে এমন। এটা Airflow 3-এর জন্য লেখা, Airflow 2-এর import পাথ কমেন্টে দেওয়া আছে; Airflow ইনস্টল করা লাগে বলে এখানে চালানো হয়নি:
from datetime import datetime, timedelta
from airflow.sdk import DAG # Airflow 2-তে: from airflow import DAG
from airflow.providers.standard.operators.python import PythonOperator
# Airflow 2-তে: from airflow.operators.python import PythonOperator
with DAG(
dag_id="shop_etl",
schedule="15 2 * * *", # cron লাইনের মতোই
start_date=datetime(2026, 6, 1),
catchup=False,
default_args={"retries": 3, "retry_delay": timedelta(minutes=5)},
) as dag:
load_sales = PythonOperator(task_id="load_sales", python_callable=run_etl_main)
build_features = PythonOperator(task_id="build_features", python_callable=build_features_main)
load_sales >> build_features # লোড সফল হলে তবেই features চলেschedule-এ একই cron এক্সপ্রেশন দেওয়া হয়। Airflow 2.4-এ এটাschedule_interval-এর জায়গা নিয়েছে, আর Airflow 3 শুধুschedule-ই মানে।start_date-এ টাইম জোন না দিলে Airflow সময়টা UTC ধরে, তাই এই DAG চলবে UTC ২:১৫-তে, মানে ঢাকায় সকাল ৮:১৫-তে।catchup=Falseথাকলে DAG চালু করার সময় Airflowstart_dateথেকে এ পর্যন্ত পেরিয়ে যাওয়া প্রতিটা ইন্টারভ্যালের জন্য আলাদা করে রান চালায় না (Airflow 3-এ এটাই ডিফল্ট)।retriesব্যর্থ টাস্ক আবার চালায়। এটা নিরাপদ শুধু টাস্কটা আইডেমপোটেন্ট বলেই: আধাআধি শেষ হওয়া রানের পর রিট্রাই হলে কিছু যেন দুবার গোনা না হয়।- Airflow-এ TaskFlow ধাঁচেও লেখা যায়: সাধারণ Python ফাংশনে
@dagআর@taskডেকোরেটর বসালেই সেগুলো DAG আর তার টাস্ক হয়ে যায়।
মনিটরিং আর অ্যালার্ট
যে পাইপলাইন চুপচাপ ব্যর্থ হয়, সেটা জোরেশোরে ব্যর্থ হওয়া পাইপলাইনের চেয়েও খারাপ। etl_runs টেবিলটা নিজেই একটা তৈরি মনিটর। যেসব ব্যর্থতায় কোনো এরর ওঠে না, সেগুলোর জন্য আলাদা চেক যোগ করুন: কয়েক দিন ধরে কিছুই লোড না হওয়া, বা খুব পুরোনো ডেটা (ফ্রেশনেস, freshness):
def health_check(wh_path, today, max_age_days=2):
"""সতর্কবার্তার লিস্ট ফেরত দেয়; খালি লিস্ট মানে সব ঠিক।"""
wh = sqlite3.connect(wh_path)
alerts = []
last = wh.execute("SELECT status FROM etl_runs ORDER BY run_id DESC LIMIT 1").fetchone()
if last is None or last[0] != "ok":
alerts.append("last run did not succeed")
newest = wh.execute("SELECT MAX(order_date) FROM fact_order_lines").fetchone()[0]
age = (today - date.fromisoformat(newest)).days
if age > max_age_days:
alerts.append(f"newest order is {age} days old")
wh.close()
return alerts
print(sqlite3.connect("warehouse.db").execute(
"SELECT run_id, mode, extracted, inserted, updated, status FROM etl_runs ORDER BY run_id").fetchall())
print(health_check("warehouse.db", today=date(2026, 6, 21)))
print(health_check("warehouse.db", today=date(2026, 6, 30)))[(1, 'full', 20, 20, 0, 'ok'), (2, 'incremental', 3, 0, 0, 'ok'), (3, 'incremental', 5, 2, 1, 'ok')]
[]
['newest order is 10 days old']প্রোডাকশনে এই লিস্ট print না করে পাঠানো হয় Slack, ইমেইল বা পেজারে। today একটা প্যারামিটার, date.today() নয়, যাতে চেকটা টেস্ট করা যায়; আসল কোডে যে ফাংশন এটাকে ডাকে, সে আজকের তারিখ পাঠায়।
পাইপলাইন টেস্ট করা
অন্য যেকোনো কোডের মতো পাইপলাইনেরও টেস্ট দরকার। ট্রান্সফর্ম টেস্ট করুন হাতে বানানো ছোট ইনপুট দিয়ে, আর আইডেমপোটেন্সি টেস্ট করুন শুরু থেকে শেষ পর্যন্ত, ডেটার এমন একটা কপিতে যা পরে ফেলে দেওয়া যায়:
import shutil
def test_transform_computes_delivery_days():
lines = [{"order_id": 1, "product_id": 1, "order_date": "2026-01-12", "status": "delivered",
"category": "Books", "quantity": 1, "revenue": 1200}]
rows = transform(lines, [{"order_id": 1, "courier": "Pathao", "delivered_on": "2026-01-14"}])
assert rows[0]["delivery_days"] == 2 and rows[0]["month"] == "2026-01"
def test_pipeline_is_idempotent():
shutil.copy("shop.db", "test_shop.db")
logging.disable(logging.CRITICAL) # টেস্টের সময় লগ বন্ধ
try:
first = run_pipeline(CourierAPI("courier_api.json"), "test_shop.db", "test_wh.db", full=True)
snapshot = fingerprint("test_wh.db")
second = run_pipeline(CourierAPI("courier_api.json"), "test_shop.db", "test_wh.db", full=True)
finally:
logging.disable(logging.NOTSET)
assert first[0] > 0 and second == (0, 0, 0)
assert fingerprint("test_wh.db") == snapshot
test_transform_computes_delivery_days()
test_pipeline_is_idempotent()
print("pipeline tests passed")pipeline tests passedপূর্ণ রিলোডও দ্বিতীয়বার কিছু বদলায় না। এই গুণটাই আগলে রাখতে হবে: ভবিষ্যতের কোনো পরিবর্তন এটা ভেঙে দিলে প্রোডাকশনে পৌঁছানোর আগেই এই টেস্ট ব্যর্থ হবে।
বাস্তবে কোথায় দেখবেন
- অ্যানালিটিক্স: অ্যাপের ডেটাবেস থেকে প্রতি রাতে BigQuery বা Snowflake-এ ELT, আর dbt মডেল দিয়ে রিপোর্টের টেবিল বানানো।
- মেশিন লার্নিং: ফিচার পাইপলাইন, যা প্রতিদিন গ্রাহকের ফিচার নতুন করে হিসাব করে; পাইপলাইন আইডেমপোটেন্ট হলে আবার চালালেও মডেল কখনো ডুপ্লিকেট সারিতে ট্রেন হয় না।
- RAG অ্যাপ: ডকুমেন্ট পাইপলাইন, যা নতুন আর বদলানো পাতা তোলে, চাঙ্কে ভাগ করে এমবেড করে, আর ডকুমেন্ট id ও চাঙ্ক id-কে কি ধরে একটা ভেক্টর টেবিলে আপসার্ট করে।
সাধারণ ভুল
- লোডে সাধারণ INSERT। যতবার আবার চালাবেন, ততবার ডেটা ডুপ্লিকেট হবে। টার্গেটে একটা স্বাভাবিক কি (natural key) রাখুন আর আপসার্ট করুন, অথবা একটা ট্রানজ্যাকশনের ভেতরে পুরো একটা পার্টিশন (একটা দিন, একটা মাস) মুছে আবার লোড করুন।
- লুকব্যাক ছাড়া ওয়াটারমার্ক। পুরোনো সারিতে দেরিতে আসা পরিবর্তন চিরকালের জন্য বাদ পড়ে যায়। লুকব্যাক উইন্ডো, একটা
updated_atকলাম, বা CDC ব্যবহার করুন। - লোড commit হওয়ার আগেই ওয়াটারমার্ক সরানো। এরপর লোড ব্যর্থ হলে পরের রান শুরু হবে এমন ডেটার পর থেকে, যা আসলে কখনো লোডই হয়নি। ডেটা আর স্টেট একই ট্রানজ্যাকশনে হালনাগাদ করুন, যেমন
run_pipelineকরে। - সব এররে রিট্রাই। যাচাইয়ের এরর বা ভুল পাসওয়ার্ড রিট্রাই করলে শুধু অ্যালার্ট দেরিতে আসে। রিট্রাই করুন শুধু সাময়িক এররে (টাইমআউট, কানেকশন রিসেট, HTTP 429/503), ব্যাকঅফ আর জিটার দিয়ে, আর চেষ্টার একটা সীমা রেখে।
- রিকনসিলিয়েশন না করা। যে জয়েন চুপচাপ ৩% সারি হারিয়ে ফেলে, সেটা সারি ধরে ধরে করা প্রতিটা চেক পার হয়ে যায়। প্রতিটা রানে উৎস আর টার্গেটের সারির সংখ্যা আর যোগফল মেলান।
নিজে চেষ্টা করুন
- সহজ:
warehouse.db-এরfact_order_linesথেকে ক্যাটাগরি অনুযায়ী আয় বের করুন, বাতিল অর্ডার বাদ দিয়ে, সবচেয়ে বেশি আয় আগে। - মাঝারি: রিকনসিলিয়েশন আরও কড়া করুন: উৎস আর ওয়্যারহাউসে প্রতিটা
status-এর আয়ও মিলিয়ে দেখুন, আর কোনো status-এ অমিল পেলে এরর raise করুন। এখনকার ডেটায় ফাংশনটা চালান। - কঠিন:
shop.db-তে অর্ডার ১৫-এর প্রোডাক্ট ৪-এর পরিমাণ ২ থেকে বাড়িয়ে ৩ করুন, পাইপলাইন চালান, আর দেখান যে ঠিক একটা fact সারি হালনাগাদ হয়েছে আরmonthly_sales-এর জুনের সারিটা বদলেছে। তারপর আরেকবার চালিয়ে দেখান যে এবার কিছুই বদলায়নি।
উত্তর
# ১. সহজ
wh = sqlite3.connect("warehouse.db")
print(wh.execute("""SELECT category, SUM(revenue) FROM fact_order_lines
WHERE status <> 'cancelled' GROUP BY category ORDER BY 2 DESC""").fetchall())
wh.close()
# ২. মাঝারি
def reconcile_by_status(src_path="shop.db", wh_path="warehouse.db"):
src, wh = sqlite3.connect(src_path), sqlite3.connect(wh_path)
source = dict(src.execute("""SELECT o.status, SUM(oi.quantity * oi.unit_price)
FROM orders o JOIN order_items oi ON oi.order_id = o.order_id
GROUP BY o.status""").fetchall())
target = dict(wh.execute("SELECT status, SUM(revenue) FROM fact_order_lines GROUP BY status").fetchall())
src.close()
wh.close()
if source != target:
raise RuntimeError(f"status totals differ: {source} vs {target}")
return target
print(reconcile_by_status())
# ৩. কঠিন
src = sqlite3.connect("shop.db")
with src:
src.execute("UPDATE order_items SET quantity = 3 WHERE order_id = 15 AND product_id = 4")
src.close()
print(run_pipeline(CourierAPI("courier_api.json")))
print(sqlite3.connect("warehouse.db").execute(
"SELECT * FROM monthly_sales WHERE month = '2026-06'").fetchone())
print(run_pipeline(CourierAPI("courier_api.json")))সারসংক্ষেপ
- পাইপলাইন ডেটা তোলে (extract), রূপ বদলায় (transform) আর লোড করে (load)। ETL লোডের আগে রূপান্তর করে; ELT কাঁচা ডেটা লোড করে, তারপর ওয়্যারহাউসের ভেতরে SQL দিয়ে রূপান্তর করে। বেশিরভাগ বাস্তব পাইপলাইনে দুটোই মেশানো থাকে।
- ইনক্রিমেন্টাল লোডে লাগে একটা ওয়াটারমার্ক আর একটা লুকব্যাক উইন্ডো (বা
updated_at, বা CDC), যাতে দেরিতে আসা পরিবর্তন বাদ না পড়ে। - আইডেমপোটেন্সি আসে কি আর আপসার্ট থেকে, সাথে ডেটা আর স্টেট একই ট্রানজ্যাকশনে হালনাগাদ করা থেকে। প্রমাণ করুন: দুবার চালিয়ে মিলিয়ে দেখুন।
- লোডের আগে সারি যাচাই করুন, পরে মোট হিসাব মেলান, প্রতিটা রান লগ করুন, আর রিট্রাই করুন শুধু সাময়িক এররে, ব্যাকঅফ দিয়ে।
- cron বা একটা অর্কেস্ট্রেটর (Airflow) দিয়ে শিডিউল করুন, রান আর ফ্রেশনেস মনিটর করুন, সমস্যায় অ্যালার্ট দিন, আর রূপান্তর ও আইডেমপোটেন্সি টেস্ট করুন।
এরপর: SQL ও pandas দিয়ে ML-এর উপযোগী ডেটা তৈরি পাতায় এমন একটা পাইপলাইনের পৌঁছে দেওয়া ডেটা থেকে ট্রেনিং সেট বানানো হবে: ফিচার আর টার্গেট, লিকেজ এড়াতে পয়েন্ট-ইন-টাইম শুদ্ধতা, আর এমন একটা পুনরুৎপাদনযোগ্য ডেটাসেট, যা দিয়ে একটা scikit-learn মডেল ট্রেন হয়।