পাঠ ০৭ · ২৯-এর মধ্যে · মডিউল ১
Home / AI Courses / Data Engineering / Advanced SQL

Window function ও advanced SQL

Window functions, recursive CTE, JSON, MERGE
৮ মিনিট পড়া মাঝারি · Intermediate SQL উদাহরণসহ

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

  • Window function — OVER, PARTITION BY, ORDER BY, frame clause
  • ROW_NUMBER, RANK, DENSE_RANK — top-N per group এবং deduplication
  • LAG/LEAD ও running aggregate — time-series analytics
  • Recursive CTE, JSON function, MERGE/UPSERT — production essentials

১ · Window function কী, কেন এত শক্তিশালী

GROUP BY rows-কে collapse করে — ১০০ row → ১০ group → ১০ row। কিন্তু কখনো আপনি চান প্রতিটি row রাখতে, পাশাপাশি group-এর তথ্য — যেমন "এই order-এর revenue, পাশাপাশি এই customer-এর গড় order"। এটাই window functionWindow Functionrow collapse না করে relative computation। OVER clause দিয়ে "window" define — partition + order + frame। DE-র power tool।।

Syntax: function_name(...) OVER (PARTITION BY ... ORDER BY ... ROWS BETWEEN ... )।

তিনটি অংশ

PARTITION BY: কোন column-এ ভাগ — যেমন প্রতি customer আলাদাভাবে।
ORDER BY: partition-এর মধ্যে sort — সময় বা amount অনুসারে।
Frame: "current row-এর আশেপাশে কোন rows" — running, sliding, expanding।

২ · ROW_NUMBER, RANK, DENSE_RANK — তিন ranking

সব ranking function দেখতে কাছাকাছি, কিন্তু tie-handling-এ পার্থক্য:

  • ROW_NUMBER(): tie ভাঙে arbitrarily — ১, ২, ৩, ৪ (একই value হলেও)।
  • RANK(): tie একই rank, পরের rank skip — ১, ২, ২, ৪।
  • DENSE_RANK(): tie একই rank, কোনো skip নেই — ১, ২, ২, ৩।
SQL · Top-N per group
-- Daraz: প্রতি category-তে top 3 best-selling product
WITH ranked AS (
  SELECT  p.product_id,
          p.name,
          p.category,
          SUM(o.units_sold) AS units,
          ROW_NUMBER() OVER (
            PARTITION BY p.category
            ORDER BY SUM(o.units_sold) DESC
          ) AS rn
  FROM    fact_orders o
  JOIN    dim_product p USING (product_key)
  WHERE   o.order_ts >= CURRENT_DATE - INTERVAL '90 days'
  GROUP BY p.product_id, p.name, p.category
)
SELECT product_id, name, category, units
FROM   ranked
WHERE  rn <= 3
ORDER BY category, rn;

    
Top-N per group — DE-তে সবচেয়ে frequent pattern। Window function-এর আগে এটা SQL-এ awkward subquery দিয়ে করতে হতো। এখন এটাই canonical solution।
Window function WHERE-এ ব্যবহার করা যায় না — কারণ window calculate হয় SELECT-এ, WHERE-এর পরে। তাই উপরে CTE বা subquery দিয়ে wrap করতে হয় — তারপর outer-এ WHERE rn <= 3।

৩ · LAG/LEAD — সময় ধরে comparison

LAG(col, n) — current row থেকে n সারি আগের value। LEAD — পরের। DE-তে এটি time-series analytics-এর কেন্দ্র।

SQL · LAG / week-over-week
-- bKash: প্রতিদিনের total transaction এবং পূর্ববর্তী দিনের তুলনায় change
SELECT  tx_date,
        daily_amount,
        LAG(daily_amount) OVER (ORDER BY tx_date) AS prev_day,
        daily_amount - LAG(daily_amount) OVER (ORDER BY tx_date) AS day_diff,
        ROUND(
          100.0 * (daily_amount - LAG(daily_amount) OVER (ORDER BY tx_date))
                / NULLIF(LAG(daily_amount) OVER (ORDER BY tx_date), 0),
          2
        ) AS pct_change
FROM (
  SELECT  DATE(transaction_ts) AS tx_date,
          SUM(amount_bdt)      AS daily_amount
  FROM    fact_transaction
  WHERE   transaction_ts >= CURRENT_DATE - INTERVAL '30 days'
    AND   status = 'success'
  GROUP BY DATE(transaction_ts)
) t
ORDER BY tx_date;

    
Self-JOIN ছাড়া আগের দিনের তুলনা — LAG-এর কারণে এক pass-এ। NULLIF(..., 0) divide-by-zero থেকে রক্ষা; এটি production SQL-এর জরুরি habit।
Window function-কে Excel-এর "প্রতিটি row-এর পাশের cell-এ একটি sum/avg/rank-এর mini-table" হিসেবে ভাবুন। Excel-এ একটি column-এর প্রতি row-এর জন্য সেই row-পর্যন্ত running total — exactly যা SUM(x) OVER (ORDER BY date) করে।

৪ · Frame clause — সবচেয়ে মিসঅনুধাবিত অংশ

ORDER BY দিলে default frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW — মানে partition শুরু থেকে এই row পর্যন্ত। এটি running total-এর জন্য ঠিক।

কিন্তু "৭-day moving average" চাইলে frame explicitly বলতে হবে:

SQL · Moving average
-- Pathao: প্রতি driver-এর শেষ ৭ দিনের earning moving average
SELECT  driver_id,
        ride_date,
        daily_earning_bdt,
        AVG(daily_earning_bdt) OVER (
          PARTITION BY driver_id
          ORDER BY ride_date
          ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
        ) AS ma_7day,
        SUM(daily_earning_bdt) OVER (
          PARTITION BY driver_id
          ORDER BY ride_date
        ) AS lifetime_total
FROM    agg_driver_daily
ORDER BY driver_id, ride_date;

    
ROWS BETWEEN 6 PRECEDING AND CURRENT ROW — current সহ ৭ row। ROWS physical row, RANGE value-based (gap-aware)। দু'টোর পার্থক্য subtle কিন্তু production-এ matter — সাবধানে বাছুন।

৫ · Deduplication — window-এর প্রিয় job

Source system থেকে duplicate row ingestion — DE-র দৈনন্দিন বাস্তবতা। CDC retry, Kafka at-least-once delivery, cron rerun — সব duplicate তৈরি করে। সমাধান:

SQL · Deduplication
-- Kafka stream-এ একই order একাধিকবার এসেছে — latest version রাখুন
WITH dedup AS (
  SELECT  *,
          ROW_NUMBER() OVER (
            PARTITION BY order_id
            ORDER BY updated_at DESC, ingestion_ts DESC
          ) AS rn
  FROM    raw_orders
)
SELECT  *
FROM    dedup
WHERE   rn = 1;

    
PARTITION BY natural_key, ORDER BY ... DESC, WHERE rn = 1 — staging-এ deduplication-এর canonical pattern। dbt model-এ অজস্রবার দেখবেন।

৬ · Recursive CTE — hierarchy traversal

Recursive CTERecursive CTEWITH RECURSIVE — base case ও recursive step-এর মাধ্যমে hierarchy বা graph traversal। SQL-এ tree query-র একমাত্র elegant উপায়। দিয়ে hierarchy ঘোরা যায় — manager-employee chain, category-subcategory tree, supply chain।

SQL · Recursive CTE
-- Daraz: category tree — একটি category-এর সব sub-category recursively
WITH RECURSIVE category_tree AS (
  -- base case: top-level category
  SELECT category_id, name, parent_id, 1 AS depth,
         CAST(name AS VARCHAR(500)) AS path
  FROM   dim_category
  WHERE  parent_id IS NULL

  UNION ALL

  -- recursive step: যেসব category-র parent এই tree-তে আছে
  SELECT c.category_id, c.name, c.parent_id, ct.depth + 1,
         CAST(ct.path || ' > ' || c.name AS VARCHAR(500))
  FROM   dim_category c
  JOIN   category_tree ct ON c.parent_id = ct.category_id
)
SELECT category_id, depth, path
FROM   category_tree
ORDER BY path;

    
Output-এ "Electronics > Mobile > Smartphone > Samsung" — পূর্ণ path। সাবধান: infinite loop সম্ভব — cycle থাকলে। অনেক engine CYCLE ... SET clause-এ rescue দেয়।

৭ · MERGE / UPSERT — idempotent pipeline

Daily pipeline-এ same source data-কে warehouse-এ sync — duplicate এড়াতে UPSERTUPSERT"insert or update" — primary key match হলে update, না হলে insert। SQL standard MERGE; Postgres-এ INSERT ... ON CONFLICT। দরকার।

SQL · MERGE / UPSERT
-- Postgres style: INSERT ... ON CONFLICT
INSERT INTO dim_customer (customer_id, name, district, customer_tier, updated_at)
SELECT  customer_id, name, district, customer_tier, NOW()
FROM    stg_customer_daily
ON CONFLICT (customer_id) DO UPDATE SET
  name          = EXCLUDED.name,
  district      = EXCLUDED.district,
  customer_tier = EXCLUDED.customer_tier,
  updated_at    = EXCLUDED.updated_at
WHERE   dim_customer.name          IS DISTINCT FROM EXCLUDED.name
   OR   dim_customer.district      IS DISTINCT FROM EXCLUDED.district
   OR   dim_customer.customer_tier IS DISTINCT FROM EXCLUDED.customer_tier;

-- Standard MERGE (BigQuery, Snowflake, SQL Server)
MERGE INTO dim_customer t
USING stg_customer_daily s
  ON  t.customer_id = s.customer_id
WHEN MATCHED AND (
  t.district != s.district OR t.customer_tier != s.customer_tier
) THEN UPDATE SET
  district = s.district, customer_tier = s.customer_tier, updated_at = CURRENT_TIMESTAMP()
WHEN NOT MATCHED THEN INSERT (customer_id, name, district, customer_tier, updated_at)
VALUES (s.customer_id, s.name, s.district, s.customer_tier, CURRENT_TIMESTAMP());

    
WHERE ... IS DISTINCT FROM — শুধু changed row-এ update, write amplification কম। dbt-এর incremental model পেছনে এই pattern।
Idempotency = একই pipeline ১০ বার চালালেও result identical। DE-র golden rule। MERGE/UPSERT + deterministic transform = production-grade pipeline।

৮ · JSON function — semi-structured data

Modern API থেকে JSON ingestion সাধারণ — Pathao webhook, bKash API response। Postgres ও cloud warehouse-এ JSON path query সরাসরি SQL-এ:

SQL · JSON
-- Pathao webhook payload থেকে field extract
SELECT  payload->>'ride_id'                     AS ride_id,
        (payload->>'fare_bdt')::NUMERIC          AS fare_bdt,
        payload->'pickup'->>'district'           AS pickup_district,
        payload->'rider'->>'rating'              AS rider_rating,
        jsonb_array_length(payload->'stops')      AS num_stops
FROM    raw_pathao_events
WHERE   payload->>'event_type' = 'ride_completed'
  AND   event_ts >= CURRENT_DATE - INTERVAL '1 day';

-- BigQuery: JSON_EXTRACT_SCALAR(payload, '$.ride_id')
-- Snowflake: payload:ride_id::STRING

    
Postgres: -> returns JSON, ->> returns text। Cast করে numeric/date। JSON ingestion-এর পর staging-এ schema enforce — production reliability।

৯ · Window function visualization

Window function — PARTITION + ORDER + FRAME 📦 Customer A PARTITION (customer = A) Jan-01 · 500 BDT Jan-02 · 300 BDT Jan-03 · 700 BDT ⬅ current Jan-04 · 200 BDT Jan-05 · 800 BDT 📦 Customer B PARTITION (customer = B) Jan-01 · 1200 BDT Jan-02 · 900 BDT Jan-03 · 600 BDT independent partition ⬅ frame: 2 PRECEDING to CURRENT SUM OVER frame A: 500 A: 800 A: 1500 ⬅ A: 1200 A: 1700 SUM(amount) OVER (PARTITION BY customer ORDER BY date ROWS BETWEEN 2 PRECEDING AND CURRENT ROW) ⚙️ প্রতি customer-এর জন্য আলাদা partition · প্রতিটির মধ্যে date-এ sort 📊 frame = current row ও আগের ২টি row (3-day rolling sum)
Window function — row collapse না করে partition-এর মধ্যে sliding aggregation। প্রতিটি row নিজস্ব answer পায়।

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

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

প্র ০১ "Sessionization" — Daraz user-এর click event থেকে session ID তৈরি (৩০ মিনিট inactive হলে নতুন session)। Window function দিয়ে কীভাবে সমাধান করবেন?

Sessionization — DE interview-এর classic প্রশ্ন। Window function-এর শক্তি দেখার আদর্শ exercise। Goal: একটানা event sequence-কে session-এ ভাগ করা।

Algorithm intuition:

  • প্রতিটি event-এর আগের event থেকে time-gap বের করুন।
  • Gap > ৩০ মিনিট হলে — নতুন session শুরু (boundary marker)।
  • Cumulative sum of boundary markers = session number।

Solution:

WITH events_with_gap AS (
  SELECT  user_id, event_ts, event_type,
          LAG(event_ts) OVER (PARTITION BY user_id ORDER BY event_ts)
            AS prev_event_ts
  FROM    raw_clickstream
  WHERE   event_ts >= CURRENT_DATE - INTERVAL '7 days'
),
boundaries AS (
  SELECT  user_id, event_ts, event_type,
          CASE
            WHEN prev_event_ts IS NULL
              OR EXTRACT(EPOCH FROM event_ts - prev_event_ts) > 1800
            THEN 1 ELSE 0
          END AS new_session_flag
  FROM    events_with_gap
),
sessions AS (
  SELECT  user_id, event_ts, event_type,
          SUM(new_session_flag) OVER (
            PARTITION BY user_id ORDER BY event_ts
          ) AS session_num
  FROM    boundaries
)
SELECT  user_id,
        user_id || '_' || session_num AS session_id,
        MIN(event_ts) AS session_start,
        MAX(event_ts) AS session_end,
        COUNT(*)      AS event_count
FROM    sessions
GROUP BY user_id, session_num;

মূল ৩টি window:

  • LAG — আগের event-এর time।
  • CASE — boundary detect (gap > 1800s)।
  • Cumulative SUM — boundary count → session number।

Variations:

  • Inactivity threshold business-rule দিয়ে: web ৩০ মিনিট, mobile app ৫ মিনিট, IoT device ১০ সেকেন্ড।
  • Bot detect — session-এ event > ১০০০, bot probable।
  • Conversion attribution — session-এ "purchase" event থাকলে conversion session।

Real-world considerations:

  • Late-arriving event (mobile-এ network delay) — out-of-order handling।
  • Cross-device — same user mobile+desktop, fingerprint-based merge।
  • Privacy — user_id hashed, retention policy।
  • Streaming version — Flink/Spark Structured Streaming-এ session window built-in।

BD context: Daraz, Shohoz, Pathao — প্রত্যেকেই sessionization করে funnel analysis-এ। Add-to-cart থেকে purchase পর্যন্ত conversion rate, drop-off page identify। এটা data engineering থেকে product analytics-এর গুরুত্বপূর্ণ সেতু।

মূল উপলব্ধি: এক query-তে ৩-৪ window stack — DE-র সিগনেচার pattern। Window function ছাড়া এটা imperative loop-এ লিখতে হতো (Python/Pandas), যা ১০-১০০x slow।

প্র ০২ ROWS বনাম RANGE frame-এ পার্থক্য কী? bKash daily transaction-এ "৭-day moving average"-এ কোনটি সঠিক, এবং কেন?

ROWS vs RANGE — সবচেয়ে subtle SQL gotcha। প্রায় সব intermediate SQL developer এক বার ভুল করে।

ROWS — physical row-based:

  • "Current row-এর পেছনে n রো"।
  • Order column-এ duplicate বা gap থাকলেও exactly n row নেয়।
  • Most predictable, deterministic।

RANGE — logical value-based:

  • "Current row-এর order_value-র চেয়ে n কম পর্যন্ত সব row"।
  • Tie রো একসাথে treat করে।
  • Date-aware: RANGE BETWEEN INTERVAL '6 days' PRECEDING — gap-aware।

Concrete example — bKash daily moving average:

ধরা যাক ডেটা missing — Jan-01, 02, 03, 05, 06, 08 (৪ আর ৭ missing)।

  • ROWS BETWEEN 6 PRECEDING AND CURRENT ROW: Jan-08-এ Jan-01 থেকে Jan-08 অবধি ৬টি row সব নেবে। মোট ৭ row। কিন্তু এটি actual ৭ ক্যালেন্ডার দিনের avg নয় — গাণিতিকভাবে ভুল।
  • RANGE BETWEEN '6 days' PRECEDING AND CURRENT ROW: Jan-08 থেকে Jan-02 পর্যন্ত (৭ ক্যালেন্ডার দিন)। Jan-01 বাদ। সঠিক — even with missing days।

সঠিক উত্তর: daily metric-এ RANGE গাণিতিকভাবে সঠিক। কিন্তু bKash-এ সাধারণত daily data-তে gap থাকে না (active business)। তাই ROWS practical, simpler।

Production pattern:

  • Date dimension JOIN: calendar table-এ all dates, LEFT JOIN fact, missing day-এ 0 fill। তারপর ROWS-এ ৭-day MA নিরাপদ।
  • এটা সবচেয়ে rigorous solution — gap explicit, MA mathematically সঠিক।

Engine-specific gotchas:

  • RANGE-এ default frame ORDER BY column tie-aware। ১,২,২,৩-এর জন্য SUM-এ tie একসাথে।
  • ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW vs RANGE same — running total-এ identical।
  • BigQuery-তে RANGE-এ INTERVAL syntax limited; Snowflake-এ full।
  • Postgres ১১+ — full RANGE support; পুরোনো version-এ ROWS।

Stock price বা financial time series example:

  • Trading day only — weekend/holiday gap ৯৯% guaranteed।
  • "৫ trading day MA" — ROWS exact (৫ market day)।
  • "৭ calendar day MA" — RANGE বা date dimension JOIN।

মূল উপলব্ধি: Frame choice business semantics-এর প্রশ্ন। "৭ দিন" মানে কী — ৭ row, না ৭ calendar day, না ৫ business day? এই precision-এ disagreement মানে dashboard ভুল answer। Senior DE এই subtleties stakeholder-এর সাথে clarify করেন।

প্র ০৩ Recursive CTE-এর performance characteristics কী? কখন graph problem-এ SQL-এর বদলে graph database (Neo4j) ব্যবহার করবেন?

Recursive CTE শক্তিশালী কিন্তু performance-এ subtle। ভুল ব্যবহার = pipeline crash।

Recursive CTE-এর working:

  • Base case execute — initial result set।
  • Recursive step — previous result-কে input হিসেবে নিয়ে নতুন rows generate।
  • UNION ALL — সব iteration concat।
  • Termination — recursive step empty result দিলে।

Performance characteristics:

  • Iteration overhead: প্রতি step আলাদা JOIN execution, optimizer plan rebuild।
  • Materialization: intermediate results memory/disk-এ — depth বাড়লে memory spike।
  • No index on intermediate: recursive step-এ JOIN করছেন un-indexed temp result-এ — slow।
  • Cycle detection: built-in protection নেই (Postgres-এ optional)। Cycle থাকলে infinite loop।

Postgres benchmark (নমুনা):

  • ১০০-row hierarchy, depth ৫ — ১০ ms।
  • ১০,০০০-row, depth ১০ — ২-৫ সেকেন্ড।
  • ১ মিলিয়ন-row, depth ১৫ — ৩০+ সেকেন্ড, প্রায়ই OOM।

SQL-এ ভালো use cases:

  • Org chart — manager-employee chain (৫-১০ level, ১০-১০০ rows per level)।
  • Category tree — ecommerce 4-5 level hierarchy।
  • Bill of materials — manufacturing parts hierarchy।
  • Sequential generation — date series, number series।
  • Path enumeration — small graph।

SQL-এ খারাপ use cases:

  • Social network: "এই user-এর friend-of-friend" — million users, billions edges। Recursive CTE মৃত্যু।
  • Shortest path: Dijkstra/A* algorithm — SQL-এ implementable কিন্তু painful, slow।
  • PageRank-style: iterative algorithm convergence — SQL-এ অসম্ভব।
  • Recommendation graph: "people also bought" — graph traversal-heavy।
  • Fraud detection: connected component detection in transaction graph।

Graph database (Neo4j, Neptune, ArangoDB) কখন:

  • Native graph storage — node ও edge both indexed।
  • Traversal optimized — O(degree) per hop, SQL-এ O(JOIN cost)।
  • Cypher/Gremlin syntax — graph-native।
  • Built-in algorithms — community detection, centrality, shortest path।
  • Whiteboard friendly — query model match real graph।

BD context examples:

  • bKash AML investigation: "এই suspicious account থেকে ৩-hop-এ কোন accounts" — graph DB, Neo4j perfect। SQL recursive CTE-তে ৫ মিনিট, Neo4j-এ ৫ সেকেন্ড।
  • Daraz recommendation: co-purchase graph। SQL-এ pre-computed adjacency table workable; real-time traversal Neo4j।
  • Pathao route optimization: driver-area-time graph। Neo4j বা specialized routing engine।
  • Telco call network analysis: Grameenphone fraud detection — call graph। Neo4j বা Spark GraphX।

Hybrid approach:

  • Source of truth — relational DB (Postgres)।
  • Graph queries — Neo4j replica, CDC-synced।
  • Best of both — relational integrity + graph performance।

মূল উপলব্ধি: Recursive CTE — small hierarchy-র জন্য brilliant। কিন্তু "graph problem" দেখলে SQL-এ force করার চেষ্টা করবেন না। Right tool for right job — DE-র maturity-র লক্ষণ।

প্র ০৪ Daraz প্রতিদিনের new orders incremental-ভাবে dimensional warehouse-এ load করতে চায়। MERGE কখন কাজ করে, কখন full refresh ভালো? Late-arriving data কীভাবে handle করবেন?

Incremental load — DE pipeline-এর সবচেয়ে impactful optimization, কিন্তু সবচেয়ে error-prone area-র একটি।

MERGE/UPSERT যখন আদর্শ:

  • Source-এ stable primary key — order_id unique এবং কখনো বদলায় না।
  • Update pattern bounded — daily ১% rows update, ৯৯% append।
  • Update-able sink — Snowflake, BigQuery, Postgres সবেই MERGE support।
  • Modest data volume — ১-১০ million row/day।

MERGE-এর challenge:

  • Cost on cloud warehouse: MERGE = full table rewrite (small partitions affected, but still I/O heavy)।
  • Concurrency: MERGE locks table — concurrent reader block। Snowflake-এ multi-version concurrency কমায়।
  • Slow on huge dim: dim_customer ১০০ million rows, daily ১০০K change — MERGE ১০ মিনিট। Append-only pattern-এ ৩০ সেকেন্ড।

Full refresh যখন ভালো:

  • Source ছোট ও fast — full snapshot daily সম্ভব।
  • State-based dimension — entire current state replace।
  • Logic complex — incremental বুঝতে গেলে bug-prone।
  • Compute cheap, storage cheap — modern cloud DW।

Append-only পদ্ধতি (যখন possible best):

  • Immutable event stream — fact table-এ ideal।
  • Insert only, update নেই — partition pruning-এ super fast।
  • Late-arriving event — late partition update, কিন্তু rest unaffected।

Late-arriving data handling:

এটি বাস্তব challenge। Pathao webhook ৩ ঘণ্টা delay-এ আসতে পারে; mobile app offline mode থেকে event ২৪ ঘণ্টা পর sync হতে পারে।

  • Lookback window: incremental load-এ "last 7 days" reload — recent data refresh, পুরোনোতে hand off।
    -- daily incremental: last 3 days reprocess
    DELETE FROM fact_orders WHERE order_date >= CURRENT_DATE - 3;
    INSERT INTO fact_orders SELECT * FROM stg_orders
      WHERE order_date >= CURRENT_DATE - 3;
  • Watermark + late event partition: normal flow daily partition; late event আলাদা fact_orders_late table, periodic merge।
  • Soft delete with version: append all versions with updated_at, query-time pick latest via window function (ROW_NUMBER() OVER PARTITION BY pk ORDER BY updated_at DESC)।
  • Slowly changing — SCD Type 2: dim_customer-এ late-arriving change — valid_from/valid_to backfill।

Production pattern (Daraz scale):

  • Hot zone (last 7 days): daily full reload — late event handle হয়।
  • Warm zone (8-90 days): immutable, partition-locked।
  • Cold zone (90+ days): archived, read-only।
  • Late event detection: event_ts vs ingestion_ts — gap track, dashboard-এ "data freshness" metric।

dbt-এ implementation:

  • materialized='incremental', unique_key='order_id', incremental_strategy='merge'।
  • {{ var('lookback_days', 7) }} dynamic lookback।
  • full_refresh mode for backfills।

Idempotency check:

  • Pipeline 2x রান হলে — sink table identical থাকা চাই।
  • Test-এ deliberately rerun, row count + checksum compare।
  • প্রোডাকশনে retry-safe pipeline = data engineer-এর peace of mind।

Cost reality (Daraz scale):

  • Full daily refresh, ১০ million order, BigQuery — ~$৫/day, ~$১৫০/month।
  • Incremental MERGE — ~$০.৫/day, ~$১৫/month।
  • Append-only with lookback partition — ~$০.৩/day, ~$৯/month।
  • Engineering cost vs cloud cost — সিদ্ধান্ত team size ও ROI-অনুযায়ী।

মূল উপলব্ধি: Incremental load একটি spectrum — full refresh, MERGE, append-only — সব valid, context-dependent। Late-arriving data সবসময় assume করুন (কখনো নেই-এর ৯৯% নিশ্চিত না হলে)। Lookback window + watermark — DE-র production survival kit।

অনুশীলন

  1. Top-N per group: bKash-এ প্রতি district-এ গত ৩০ দিনে top 5 highest-spending customer বের করুন।
    WITH spend_30d AS (
      SELECT  c.district, c.customer_id, c.name,
              SUM(t.amount_bdt) AS spent
      FROM    fact_transaction t
      JOIN    dim_customer c USING (customer_key)
      WHERE   t.transaction_ts >= CURRENT_DATE - INTERVAL '30 days'
        AND   t.status = 'success'
      GROUP BY c.district, c.customer_id, c.name
    ),
    ranked AS (
      SELECT *,
             DENSE_RANK() OVER (
               PARTITION BY district ORDER BY spent DESC
             ) AS rk
      FROM   spend_30d
    )
    SELECT district, customer_id, name, spent
    FROM   ranked
    WHERE  rk <= 5
    ORDER BY district, rk;

    DENSE_RANK ব্যবহার — tie হলে ৫-এর বেশি customer পেতে পারেন, কিন্তু rank skip হয় না। ব্যবসায়িকভাবে সাধারণত সঠিক।

  2. Running total + day-over-day: Pathao-তে last ৩০ দিনের প্রতিদিনের ride count, cumulative count, এবং পূর্ববর্তী দিনের সাথে % change।
    WITH daily AS (
      SELECT DATE(ride_ts) AS d, COUNT(*) AS rides
      FROM   fact_ride
      WHERE  ride_ts >= CURRENT_DATE - INTERVAL '30 days'
        AND  status = 'completed'
      GROUP BY DATE(ride_ts)
    )
    SELECT  d, rides,
            SUM(rides) OVER (ORDER BY d) AS cum_rides,
            ROUND(
              100.0 * (rides - LAG(rides) OVER (ORDER BY d))
                    / NULLIF(LAG(rides) OVER (ORDER BY d), 0), 2
            ) AS pct_change
    FROM    daily
    ORDER BY d;
  3. UPSERT লিখুন: daily পচাচ্ছ stg_product থেকে dim_product-এ MERGE — শুধু changed row update।
    INSERT INTO dim_product (product_id, name, category, price_bdt, updated_at)
    SELECT product_id, name, category, price_bdt, NOW()
    FROM   stg_product
    ON CONFLICT (product_id) DO UPDATE SET
      name       = EXCLUDED.name,
      category   = EXCLUDED.category,
      price_bdt  = EXCLUDED.price_bdt,
      updated_at = EXCLUDED.updated_at
    WHERE  dim_product.name      IS DISTINCT FROM EXCLUDED.name
       OR  dim_product.category  IS DISTINCT FROM EXCLUDED.category
       OR  dim_product.price_bdt IS DISTINCT FROM EXCLUDED.price_bdt;

    IS DISTINCT FROM — NULL-safe comparison। শুধু changed row update — write amplification কম, audit trail পরিষ্কার।

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

Window function practice কোথায়? Google Colab ব্যবহার করুন — DuckDB ও SQLite (3.25+) দু'টোতেই window function full support। sample dataset দিয়ে এই query সব trivially চলবে।
পূর্ববর্তী পাঠ
পাঠ ০৬ · SQL refresher — DE-এর জন্য