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

Spark optimization — দ্রুত ও সাশ্রয়ী

Spark optimization — partitioning, joins, AQE, skew
৭ মিনিট পড়া মাঝারি · Intermediate Spark UI সহ

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

  • Partition, shuffle, ও task — Spark execution-এর তিন মৌল
  • Broadcast join, caching, file format choice — তিন বড় optimization
  • AQE কী করে এবং কীভাবে চালু করবেন
  • Data skew চিনতে ও salting দিয়ে সমাধান

১ · কেন Spark slow হয়? — তিন bottleneck

Spark-এ একটি job slow হলে প্রায় সবসময়ই কারণ একটি — অপ্রয়োজনীয় shuffleShuffleCluster-এর executor-গুলোর মধ্যে network দিয়ে data redistribute করা। JOIN, groupBy, repartition-এ ঘটে। Spark-এর সবচেয়ে costly operation।। Shuffle মানে — executor-গুলোর মধ্যে network দিয়ে ডেটা পাঠানো। CPU দ্রুত (ন্যানোসেকেন্ড), কিন্তু network slow (মিলিসেকেন্ড)। ফলে ১ TB ডেটা shuffle করলে — ১০ মিনিটের কাজ ১ ঘণ্টা লাগতে পারে।

তিনটি প্রধান bottleneck

১) Shuffle: JOIN, groupBy, repartition — সব ডেটা executor-এর মধ্যে move হয়।
২) Skew: এক partition-এ বেশি ডেটা — এক task সবার চেয়ে দেরিতে শেষ।
৩) Bad file format: CSV/JSON full scan; Parquet column-pruning করে দ্রুত।

bKash-এর উদাহরণ ভাবুন। প্রতিদিন ৫ কোটি transaction। যদি একটি Spark job-এ groupBy(user_id) করা হয় — তখন একই user-এর সব row একই executor-এ যেতে হবে। যদি এক user (যেমন একটি merchant account)-এর ১ কোটি transaction থাকে — সেই executor অন্য সবার চেয়ে অনেক দেরিতে শেষ হবে। এটাই data skew।

২ · Partition — Spark-এর atomic unit

Spark ডেটাকে partitionPartitionSpark DataFrame-এর একটি logical চuck — সাধারণত ১২৮ MB। প্রতিটি partition একটি task, একটি core-এ চলে। partition কম মানে under-utilization, বেশি মানে scheduler overhead।-এ ভাগ করে রাখে। প্রতিটি partition = একটি task = একটি core-এ চলবে। ১,০০০ partition + ১,০০০ core থাকলে — সবগুলো parallel চলে।

আদর্শ partition size: ১২৮ MB থেকে ২৫৬ MB। এর চেয়ে ছোট হলে scheduling overhead, বড় হলে memory pressure।

Default partition সংখ্যা spark.sql.shuffle.partitions = 200। ছোট cluster-এ ২০০ অনেক বেশি; বড় cluster-এ অনেক কম। সঠিক value calculate করুন: মোট ডেটা / ১২৮ MB।

Pathao-এর ১০০ rider আছে। যদি ১০টি delivery থাকে — ৯০ জন বসে। যদি ১,০০০ delivery থাকে — প্রত্যেক rider-এর ১০ delivery, কিন্তু কারও ১০০, কারও ১। এটাই partition planning ও skew-এর সমস্যা।

৩ · Broadcast join — shuffle-এর সবচেয়ে বড় remedy

একটি বড় table (১০০ GB) ও একটি ছোট table (৫০ MB) JOIN করছেন। Default Spark — দু'টোই shuffle করবে। কিন্তু ছোট table-টা প্রতিটি executor-এ memory-তে copy পাঠালে কোনো shuffle লাগে না। এটাই broadcast joinBroadcast Joinছোট table-টিকে সব executor-এর memory-তে duplicate করে — তারপর local-ভাবে JOIN। কোনো shuffle নেই। Spark default ১০ MB-র নিচে auto-broadcast করে; বড় হলে hint দিতে হয়।।

Spark default-এ spark.sql.autoBroadcastJoinThreshold = 10MB। এর চেয়ে বড় table হলে আপনাকে hint দিতে হবে: broadcast(small_df)।

সাবধান: ৫০০ MB-র বেশি table broadcast করবেন না। প্রতিটি executor-এ duplicate হবে — out-of-memory ঘটতে পারে। ছোট dim table (lookup, country, category) broadcast করুন, fact table নয়।

৪ · Caching — যখন একই data বারবার লাগে

একটি DataFrame যদি ৫ বার ব্যবহার হয় — Spark default-এ ৫ বার পুরো computation করে। df.cache() দিলে — প্রথমবার compute, পরের ৪ বার memory থেকে।

Storage levels: MEMORY_ONLY (default, ছোট data), MEMORY_AND_DISK (বড় data, fallback), MEMORY_ONLY_SER (serialized, কম memory)।

যা cache করবেন না — একবার ব্যবহৃত data, বা যা cluster memory-র চেয়ে বড়। Cache হিতাহিত — সবসময় Spark UI-তে storage tab দেখে নিশ্চিত করুন।

৫ · File format — Parquet সবচেয়ে গুরুত্বপূর্ণ সিদ্ধান্ত

ParquetParquetApache Parquet — column-oriented binary file format। Snappy/Gzip compression built-in, predicate pushdown, schema evolution। Spark, Hive, Athena, BigQuery — সবাই native support করে। হলো Spark-এর native format। তিন কারণ:

  • Columnar: ১০০ column থেকে ৩টি লাগলে — শুধু সেই ৩টি পড়ে। CSV সব পড়তে হয়।
  • Compression: Snappy default — CSV-র ৫× ছোট।
  • Predicate pushdown: WHERE date='2026-05-09' — Parquet metadata দেখে অপ্রয়োজনীয় file skip করে।

Partition column: df.write.partitionBy("date").parquet(path) — তারিখ অনুসারে আলাদা ফোল্ডার। Daraz-এর order data partitionBy("order_date") করলে — গত ৭ দিনের query শুধু ৭ ফোল্ডার scan করে, পুরো ২ বছরের data নয়।

Partition column এমন কিছু রাখুন যার unique value কম (১০-১,০০০)। user_id দিয়ে partition করলে — কোটি ফোল্ডার তৈরি হবে; metadata-র জন্য Spark crash করবে। Date, country, category — ভাল choice।

৬ · AQE (Adaptive Query Execution) — Spark ৩.০-এর গেম-চেঞ্জার

Spark ৩.০-র আগে — query plan static ছিল। সব decision compile-time-এ। AQE এই plan runtime-এ adapt করতে পারে।

AQE তিনটি কাজ স্বয়ংক্রিয়ভাবে করে:

  • Coalesce shuffle partitions: ২০০ partition-এর ১৮০টি ছোট হলে merge করে ৫০-এ আনে।
  • Switch join strategy: runtime-এ যদি দেখে এক side ছোট — sort-merge থেকে broadcast-এ switch।
  • Skew handling: বড় partition split করে সমান-সমান করে।

চালু করুন (Spark ৩.০+):

spark.conf.set("spark.sql.adaptive.enabled", "true")

Spark ৩.২+ থেকে AQE default-এ চালু। কিন্তু পুরনো cluster-এ check করে নিন।

৭ · Data skew — এবং salting দিয়ে সমাধান

Skew মানে — কিছু partition-এ অনেক বেশি ডেটা, বাকিদের কম। উদাহরণ: bKash-এর data-তে user_id="merchant_001" (একটি বড় merchant)-এর ১ কোটি row, অন্য সবার গড়ে ১,০০০।

Salting — skewed key-তে random suffix যোগ করুন। তারপর দু'side-এই একই suffix যোগ করে JOIN। ফলে এক বড় partition → ১০টি ছোট partition।

Spark slow? — কোথায় সমস্যা Optimization decision tree ⚠️ Slow Spark job Open Spark UI → Stages এক task অনেক slow? → data skew suspected salt key + AQE skew spark.sql.adaptive.skewJoin Shuffle GB-পরিমাণ? → broadcast small side broadcast(small_df) + tune partitions Read stage slow? → format / partition issue parquet + partitionBy + predicate pushdown 🔧 সবচেয়ে আগে: AQE চালু করুন spark.sql.adaptive.enabled = true Spark UI → Stages tab → click slowest stage
Spark slow হলে — Spark UI দেখুন, কোন stage slow তা চিহ্নিত করুন, তারপর সঠিক remedy প্রয়োগ করুন।

৮ · Salting — হাতে-কলমে কোড

নিচের কোডে — bKash-এর transaction data-তে skewed merchant_id-এর জন্য salting:

Python · PySpark
from pyspark.sql import functions as F

# Skewed: এক merchant_id-তে কোটি row
N_SALTS = 10  # বড় partition-কে ১০ ভাগে ভাঙ্গবো

# left side: salt যোগ
tx_salted = transactions.withColumn(
    "salt", (F.rand() * N_SALTS).cast("int")
)

# right side (small dim): প্রতিটি merchant-এর জন্য ১০ copy
salts = spark.range(N_SALTS).withColumnRenamed("id", "salt")
merchants_salted = merchants.crossJoin(salts)

# এখন salt + merchant_id দু'টোতেই join
result = tx_salted.join(
    merchants_salted,
    on=["merchant_id", "salt"],
    how="inner"
)

# AQE skew handling-ও চালু রাখুন
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")

result.write.mode("overwrite").parquet("/data/tx_enriched")

    
Salting-এর পর — এক skewed merchant-এর ১ কোটি row → ১০টি partition-এ ১০ লাখ করে। প্রতিটি executor এখন প্রায় সমান কাজ পায়। AQE additionally runtime-এ আরও skew detect করতে পারে।

৯ · Spark UI — কোথায় কী দেখবেন

Spark Web UI (default port 4040)-এ পাঁচটি tab গুরুত্বপূর্ণ:

Spark UI · navigation
# http://localhost:4040 (local) বা EMR/Databricks UI

Jobs        → কোন job কত সময় নিল
Stages      → প্রতি stage-এর shuffle read/write, task duration
            → "Tasks" পেজে min/median/max time দেখুন
            → max ÷ median > 5 মানে skew
SQL         → query plan, AQE rewrite
Executors   → memory usage, GC time, failed tasks
Storage     → cached DataFrames — কত memory ব্যবহার

# Key metrics:
- Shuffle Read/Write: GB-পরিমাণ হলে broadcast সম্ভব?
- GC Time: total time-এর >১০% মানে memory pressure
- Spilled (Memory): disk-এ spill = memory কম, executor বড় করুন

    
প্রথম দেখার জায়গা: Stages tab → slowest stage → Tasks → max time vs median time। ৫× বেশি হলে skew, কাছাকাছি হলে partition বাড়ান বা broadcast দেখুন।
Production-এ Spark UI বন্ধ হয়ে যায় job শেষে। তাই spark.eventLog.enabled=true দিয়ে event log save করুন এবং Spark History Server-এ পরে দেখুন।

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

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

প্র ০১ Spark default-এ spark.sql.shuffle.partitions = 200। কেন এই নির্দিষ্ট সংখ্যা? ছোট cluster-এ ও বড় cluster-এ এটি কীভাবে বদলাবেন?

200 সংখ্যাটি Spark-এর designers ১০ বছর আগে বেছেছিলেন — যখন একটি "typical" Spark cluster-এ ২০-৪০ executor-এ মোট ২০০ core থাকত। এক core = এক task, তাই ২০০ partition = সব core busy। এটি একটি "একটি size সবার জন্য" default — এবং প্রায়ই ভুল।

সঠিক value-র সূত্র:

  • Total shuffle data ÷ আদর্শ partition size (১২৮-২৫৬ MB)।
  • উদাহরণ: ১ TB shuffle হলে — ১,০০০,০০০ MB ÷ ১২৮ MB = ৭,৮১২ partitions।
  • খুব ছোট data (১০ GB) — ১০,০০০ ÷ ১২৮ = ৭৮ partitions। ২০০ default এক্ষেত্রে ৩× বেশি।

ছোট cluster (৪ executor × ৪ core = ১৬ core):

  • ২০০ partitions = প্রতি core ১২.৫ task — scheduling overhead বেশি।
  • আদর্শ: ৩২-৬৪ partitions (২-৪× core সংখ্যা)।
  • spark.sql.shuffle.partitions = 32 দিন।

বড় cluster (Daraz scale: ১০০ executor × ৮ core = ৮০০ core, ৫ TB data):

  • ২০০ partitions মানে — প্রতি partition ২৫ GB। OOM নিশ্চিত।
  • আদর্শ: ৫,০০০-১০,০০০ partitions (৫ TB ÷ ৫১২ MB = ১০,০০০)।
  • বেশি partition হলে driver-এ task scheduling overhead — কিন্তু ১,০০,০০০-এর নিচে সাধারণত ঠিক।

আজকের Best practice — AQE: Spark ৩.২+ AQE default-এ চালু। তখন আপনি initial partition বেশি দিতে পারেন (যেমন ২,০০০) — AQE runtime-এ ছোটগুলোকে coalesce করে। তাই over-partition কম খরচের, under-partition বেশি খরচের।

মূল উপলব্ধি: Default value blind ভাবে accept না করে — data size, cluster size, ও workload analyze করুন। এটাই junior থেকে senior data engineer-এ পার্থক্য।

প্র ০২ আপনি Pathao-এর জন্য daily report বানাচ্ছেন — ১৫০ GB ride data + ৫০ MB driver dim table JOIN। Cluster-এ ১৬ executor × ৪ core। সবচেয়ে দ্রুত strategy কী?

এটি একটি classic dim-table JOIN — broadcast join-এর জন্য perfect। চলুন step-by-step ভাবি।

Default strategy কী হবে?

  • ৫০ MB < ১০ MB threshold — তাই Spark auto-broadcast করবে না।
  • Default sort-merge join — দু'side shuffle।
  • ১৫০ GB shuffle write + ১৫০ GB shuffle read = ৩০০ GB network transfer।
  • ~১০ Gbps network-এ ৪-৫ মিনিট শুধু shuffle-এ।

Optimal strategy:

  1. Threshold বাড়ান: spark.sql.autoBroadcastJoinThreshold = 100MB। ৫০ MB driver table এখন auto-broadcast।
  2. বা hint দিন: rides.join(broadcast(drivers), "driver_id")।
  3. Read side: rides ডেটা যদি Parquet ও partitionBy("ride_date")-এ থাকে — শুধু আজকের partition পড়ুন।
  4. Output: result-ও Parquet-এ partition by date লিখুন।

তুলনা:

  • Sort-merge: ~৬ মিনিট (১৫০ GB shuffle)।
  • Broadcast: ~৪০ সেকেন্ড (৫০ MB × ১৬ executor = ৮০০ MB একবার broadcast)।
  • ৯× speedup।

Memory check: ৫০ MB × ১৬ executor = ৮০০ MB মোট network usage broadcast-এ। প্রতিটি executor-এ ৫০ MB extra memory। ৪ GB executor-এ — ১.২৫% — কোনো সমস্যা নেই।

যদি driver table ১ GB হতো?

  • ১ GB × ১৬ executor = ১৬ GB network — broadcast হলে ভাল।
  • প্রতি executor-এ ১ GB extra memory — ৪ GB executor-এ tight।
  • Border case — পরিমাপ করে decide করুন।

যদি driver table ১০ GB হতো?

  • Broadcast = OOM নিশ্চিত।
  • Sort-merge join + bucketing (যদি দু'side একই key-এ bucketed)।
  • অথবা bloom filter join (Spark ৩.৩+)।

মূল কথা: Dimension table-এর size জানা — JOIN strategy-র সবচেয়ে গুরুত্বপূর্ণ ইনপুট। সবসময় EXPLAIN বা Spark UI-তে query plan check করে দেখুন কোন strategy চলছে।

প্র ০৩ BTRC-এর telecom CDR (call detail records) data — দেশের ১৭ কোটি sim-এর মধ্যে ০.১% sim-এর ৪০% traffic। কেন simple groupBy(sim_id) দিনের পর দিন hang করে? সমাধানের তিন পদ্ধতি ব্যাখ্যা করুন।

এটি extreme data skew-এর classic case। ০.১% × ১৭ কোটি = ১.৭ লাখ "heavy hitter" sim — এরা business clients (টেলিমার্কেটিং, IVR, automated alerts)। এদের প্রত্যেকের কয়েক হাজার call/day, যেখানে গড় sim-এর ৫টা।

কেন hang করে?

  • groupBy(sim_id) = shuffle by hash(sim_id)।
  • একটি heavy sim-এর সব record একই partition-এ যায়।
  • সেই partition-এ ৩-৫ কোটি row, অন্যগুলোতে ১,০০০।
  • ৯৯৯ task ৩০ সেকেন্ডে শেষ; ১ task ৩ ঘণ্টা চলে।
  • Spark UI-তে — "1/1000 tasks in progress" stuck।

সমাধান ১: AQE skewJoin (সবচেয়ে সহজ)

  • spark.sql.adaptive.enabled = true
  • spark.sql.adaptive.skewJoin.enabled = true
  • Spark runtime-এ skewed partition detect করে split করে।
  • Threshold: median × ৫ এবং ২৫৬ MB-র বেশি।
  • Limitation: শুধু JOIN-এ কাজ করে, pure groupBy-তে limited।

সমাধান ২: Two-stage aggregation (groupBy-এর জন্য আদর্শ)

  • Stage 1: groupBy(sim_id, salt) — N salts দিয়ে।
  • Stage 2: groupBy(sim_id) — partial result merge।
  • Heavy sim-এর data ১০০ partition-এ ছড়ানো; ১০০টি ছোট sum।
  • তারপর ১০০ → ১ — ছোট operation।

সমাধান ৩: Separate heavy hitters

  • প্রথমে heavy sim list বের করুন: SELECT sim_id FROM cdr GROUP BY sim_id HAVING COUNT(*) > 100000।
  • Heavy ও non-heavy আলাদা DataFrame-এ ভাগ করুন।
  • Non-heavy: normal groupBy।
  • Heavy: প্রতিটি sim-কে আলাদা job বা salting দিয়ে।
  • শেষে UNION।
  • Best for extreme cases (১,০০০× skew)।

BTRC scale-এ practical recommendation: AQE চালু রাখুন। তারপর salting দিয়ে two-stage aggregation। শুরুতে complex separate-heavy strategy দরকার নেই — সাধারণত প্রথম দু'টি যথেষ্ট।

মূল উপলব্ধি: Skew = real-world data-র ভিত্তিগত বৈশিষ্ট্য। Power-law distribution সর্বত্র — telecom, ecommerce, social network। Spark engineer হিসেবে আপনার অর্ধেক সময় skew-এর সাথে যাবে।

প্র ০৪ আপনার team CSV ছেড়ে Parquet-এ migrate করতে চায়। ১০ TB historical data। কী কী trade-off আছে? Compression, partition column, schema evolution — কী strategy?

CSV → Parquet migration প্রতিটি data team-এর "growing up" milestone। ভুল করলে অপরিবর্তনীয় cost; সঠিক করলে বছরের পর বছর সাশ্রয়।

Benefit (কেন migrate):

  • Storage: Snappy compression — ৩-৫× ছোট। ১০ TB → ২-৩ TB।
  • Query speed: Column pruning — যদি query-তে ৫টি column লাগে, ১০০ column-এর CSV-র সব পড়তে হবে; Parquet শুধু সেই ৫টি।
  • Predicate pushdown: WHERE date='2026-05-01' — Parquet metadata দেখে অর্ধেক file skip।
  • Type safety: CSV সব string; Parquet-এ int, timestamp, struct সবই native।

Trade-off (কেন careful):

  • Human-readable নয়: CSV Excel-এ খুলতে পারেন; Parquet — special tool (parquet-tools, DuckDB) দরকার।
  • Append cost: CSV-তে নতুন row append সহজ; Parquet-এ পুরো file rewrite।
  • Small files: অনেক ছোট Parquet file performance ক্ষতি করে। Compaction job দরকার।

Compression choice:

  • Snappy (default): দ্রুত compress/decompress, মাঝারি ratio। সব analytical workload-এ standard।
  • Gzip: বেশি ছোট, কিন্তু slow decompress। Cold archive-এর জন্য।
  • Zstd: Spark ৩.২+ — Snappy-র চেয়ে ভাল ratio ও comparable speed। নতুন project-এ এটাই recommend।

Partition column strategy:

  • সবচেয়ে ভাল: Date (order_date, event_date)। Time-based query-এ massive speedup।
  • Cardinality rule: ১০-১০,০০০ unique value আদর্শ। বেশি = file explosion, কম = বড় partition।
  • মাল্টি-level: partitionBy("year", "month", "day") — হায়ারার্কিকাল।
  • সাবধান: একটি partition-এ কমপক্ষে ১২৮ MB হওয়া উচিত। কম হলে — over-partitioning।

Schema evolution:

  • Parquet built-in schema evolution support করে।
  • Safe changes: Column যোগ (পুরাতন file-এ null), nullable → not-null (যদি data clean)।
  • Risky: Type change (int → string) — full rewrite।
  • Best: Delta Lake বা Iceberg ব্যবহার করুন — ACID + time travel + schema evolution managed।

Migration plan (Bangladesh-bank scale):

  1. Sample (১% data) দিয়ে POC — speed/size measure।
  2. Schema lock — সব column type, nullable সিদ্ধান্ত।
  3. Backfill: পুরাতন CSV → Parquet, partitionBy date — সপ্তাহান্তে batch।
  4. New writes — সরাসরি Parquet।
  5. Dual-write period (১ মাস) — নতুন pipeline-এ দু'জায়গায় write, query Parquet, fallback CSV।
  6. Sunset CSV — verify ৩০ দিন কেউ পড়েনি।

মূল কথা: Parquet migration এক-দিনের কাজ নয়। Schema, partition, compaction strategy আগে set করে নিন। ভুল partition column বছরের পর বছর pain।

অনুশীলন

  1. Decide করুন: ২ TB fact table + ২৫০ MB dim table JOIN। Cluster: ৪ GB executor। Default broadcast threshold ১০ MB। কী strategy?

    ২৫০ MB dim table — broadcast-এর জন্য ভালো (৪ GB executor-এ ৬%)। তিন কাজ:

    • spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "300MB")
    • বা explicit hint: fact.join(F.broadcast(dim), "key")
    • AQE চালু রাখুন (spark.sql.adaptive.enabled=true)

    Sort-merge হলে ২ TB shuffle — ১০-১৫× slow হবে।

  2. Skew detect: Spark UI-তে stage-এর ১০০টি task-এর median time ১০ সেকেন্ড, max time ১০ মিনিট। কী সমস্যা ও সমাধান?

    Max ÷ median = ৬০× → severe skew।

    • সম্ভাব্য কারণ: groupBy/join key-তে hot value (এক customer/merchant/date)।
    • Solution 1: AQE skewJoin চালু (spark.sql.adaptive.skewJoin.enabled=true)।
    • Solution 2: Salting — skewed key-তে random suffix যোগ।
    • Solution 3: Two-stage aggregation।
    • Solution 4: Heavy hitters আলাদা করে process।
  3. Configure করুন: bKash-এর daily summary job — ২০০ GB transactions, ১০ executor × ৮ core। কী partition সংখ্যা ও কোন format রাখবেন?
    # Cluster: 80 cores
    # Data: 200 GB shuffle expected
    # Ideal partition: 200,000 MB / 128 MB ≈ 1,560
    
    spark.conf.set("spark.sql.shuffle.partitions", "1600")
    spark.conf.set("spark.sql.adaptive.enabled", "true")  # AQE will coalesce if needed
    spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
    
    # Read: Parquet partitioned by date
    df = spark.read.parquet("/data/transactions/date=2026-05-09")
    
    # Output: Parquet, partition by hour for finer granularity
    result.write.mode("overwrite") \
      .partitionBy("date", "hour") \
      .parquet("/data/tx_summary")

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

Spark চালাতে চান? Local-এ pip install pyspark অথবা Google Colab ব্যবহার করুন — ফ্রি Java + Python পরিবেশ, শুধু Gmail অ্যাকাউন্ট লাগে।
পূর্ববর্তী পাঠ
পাঠ ১১ · PySpark DataFrame ও SQL