Spark optimization — দ্রুত ও সাশ্রয়ী
এই পাঠে যা শিখবেন
- 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 করলে — ১০ মিনিটের কাজ ১ ঘণ্টা লাগতে পারে।
১) 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।
৩ · 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)।
৪ · 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 নয়।
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।
৮ · Salting — হাতে-কলমে কোড
নিচের কোডে — bKash-এর transaction data-তে skewed merchant_id-এর জন্য salting:
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")
৯ · Spark UI — কোথায় কী দেখবেন
Spark Web UI (default port 4040)-এ পাঁচটি tab গুরুত্বপূর্ণ:
# 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 বড় করুন
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:
-
Threshold বাড়ান:
spark.sql.autoBroadcastJoinThreshold = 100MB। ৫০ MB driver table এখন auto-broadcast। -
বা hint দিন:
rides.join(broadcast(drivers), "driver_id")। -
Read side: rides ডেটা যদি Parquet ও
partitionBy("ride_date")-এ থাকে — শুধু আজকের partition পড়ুন। - 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 = truespark.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):
- Sample (১% data) দিয়ে POC — speed/size measure।
- Schema lock — সব column type, nullable সিদ্ধান্ত।
- Backfill: পুরাতন CSV → Parquet, partitionBy date — সপ্তাহান্তে batch।
- New writes — সরাসরি Parquet।
- Dual-write period (১ মাস) — নতুন pipeline-এ দু'জায়গায় write, query Parquet, fallback CSV।
- Sunset CSV — verify ৩০ দিন কেউ পড়েনি।
মূল কথা: Parquet migration এক-দিনের কাজ নয়। Schema, partition, compaction strategy আগে set করে নিন। ভুল partition column বছরের পর বছর pain।
অনুশীলন
-
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 হবে।
-
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।
-
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-এ আপনার পরবর্তী পদক্ষেপ
- পাঠ ১৩ · Apache Airflow পরিচিতি পরবর্তী পাঠ Spark job কে schedule ও monitor করতে — workflow orchestration।
- পাঠ ১১ · PySpark DataFrame ও SQL আগের পাঠ Optimization-এর আগে DataFrame API ভাল করে বুঝে নিন।
- পাঠ ১৯ · Spark Structured Streaming এই পাঠের সাথে সম্পর্কিত একই optimization rule streaming workload-এ আরও critical।
- সব AI Courses দেখুন ABCL TECH Python, ML, DL, NLP, CV, GenAI, RL, MLOps — সব AI কোর্স একসাথে।
pip install pyspark অথবা
Google Colab
ব্যবহার করুন — ফ্রি Java + Python পরিবেশ, শুধু Gmail অ্যাকাউন্ট লাগে।