Apache Spark পরিচিতি
এই পাঠে যা শিখবেন
- Distributed processing কেন দরকার, single machine কখন কম পড়ে
- Spark architecture — driver, executor, cluster manager
- RDD, DataFrame, Dataset — পার্থক্য ও কখন কোনটি
- Lazy evaluation, transformation vs action — Spark-এর জাদু
- Spark vs Hadoop MapReduce — কেন Spark জিতল
১ · কেন distributed processing দরকার
Grameenphone-এর একদিনের call detail record (CDR) — প্রায় ১.৫ TB। আপনার laptop-এর RAM ১৬ GB — পুরো ফাইল মেমোরিতে fit হবে না। Disk-এ stream করলেও — single CPU দিয়ে scan-aggregate করতে ১০-১৫ ঘণ্টা।
যদি ১০০টি machine থাকে — প্রতিটি ১৫ GB করে নিল, parallel-এ process করল — ৬ মিনিটেই কাজ শেষ। এটাই distributed processingDistributed Processingএকটি বড় কাজকে অনেক ছোট অংশে ভাগ করে — একাধিক machine-এ parallel চালানো। বিশাল data বা compute-heavy কাজে অপরিহার্য — Hadoop, Spark, Flink, Dask এই paradigm-এ। ।
১) Data > RAM: ১TB data, ৬৪GB RAM — disk swap, painfully slow।
২) Compute long: ১ CPU-তে ১০ ঘণ্টার job — production-এ unacceptable।
৩) Single point of failure: machine crash মানে কাজ পুরো হারানো।
৪) Scaling up costly: ১TB RAM-এর server $৫০K+; ১০×৬৪GB cluster অনেক সস্তা।
২ · Apache Spark কী, ইতিহাস
Spark UC Berkeley AMPLab-এ ২০০৯-এ শুরু — Matei Zaharia-র PhD প্রজেক্ট। ২০১৩-তে Apache top-level project, ২০১৪-তে Databricks কোম্পানি গঠন। আজ — distributed data processing-এর de-facto standard।
Spark-এর মূল ধারণা: RDD (Resilient Distributed Dataset) — একটি immutable, partitioned, fault-tolerant collection যা cluster-এর memory-তে থাকে। MapReduce-এর তুলনায় ১০-১০০× দ্রুত — কারণ MapReduce প্রতিটি step-এর পর disk-এ লিখত; Spark RAM-এ রাখে।
৩ · Architecture — driver, executor, cluster manager
Spark application-এ তিনটি প্রধান component:
- Driver: "main" program — আপনার code এখানে চলে। DAG তৈরি, scheduling, result collection।
- Executor: worker process — actual data processing। প্রতিটি executor-এর নিজস্ব RAM ও CPU core।
- Cluster manager: resource allocate করে — YARN (Hadoop), Mesos, Kubernetes, বা Spark-এর নিজস্ব Standalone।
৪ · Partition — parallelism-এর একক
PartitionPartitionএকটি বিশাল dataset-এর ছোট ছোট অংশ — প্রতিটি partition একটি executor-এর memory-তে fit হয়। Spark default partition size সাধারণত ১২৮MB (HDFS block size)। হলো dataset-এর একটি অংশ। ১TB data হয়তো ৮,০০০ partition-এ ভাগ (১২৮MB each)। প্রতিটি partition একটি task — একটি core-এ চলে।
Parallelism = total core (cluster-এ)। ১০০ core থাকলে ১০০ task একসাথে চলে। বাকি ৭,৯০০ task wait করে। তাই partition সংখ্যা > core সংখ্যা — সাধারণত ২-৪×।
৫ · RDD, DataFrame, Dataset — তিন abstraction
- RDD (Resilient Distributed Dataset): Spark-এর মূল primitive। distributed Java/Scala object collection। Type-safe, কিন্তু optimization manual।
- DataFrame: distributed table — column ও row, schema সহ। SQL-এর মতো API। Catalyst optimizer auto-optimize করে। Python/R-এ এটাই ব্যবহার্য।
- Dataset: DataFrame-এর type-safe version (Scala/Java only)। compile-time type check।
৯০% কাজে DataFrame। Catalyst optimizer-এর সুবিধা পান, কোড পরিষ্কার, multi-language (Python/Scala/R/SQL)। RDD শুধু low-level control দরকার হলে — যেমন custom partitioning। Dataset Scala-তে best-of-both, কিন্তু PySpark-এ অনুপস্থিত।
৬ · Lazy evaluation — Spark-এর জাদু
Spark-এ দু'ধরনের operation:
- Transformation: নতুন DataFrame তৈরি —
filter(),select(),groupBy(),join()। চলে না — শুধু plan জমে। - Action: result ফেরায় বা write করে —
show(),count(),collect(),write()। এখানেই execution শুরু।
Action-এ Spark পুরো DAG দেখে — কোন step আগে, কোন filter কোন join-এর আগে push করা যায়, কোন column-গুলোই শুধু দরকার — সব auto-optimize। এটা SQL planner-এর মতো, কিন্তু distributed scale-এ।
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum as _sum
spark = SparkSession.builder.appName("BkashAnalysis").getOrCreate()
# transformations — chained, lazy (চলে না এখনো)
txns = (spark.read
.parquet("s3://bkash-data/2026/05/")
.filter(col("status") == "success")
.filter(col("amount") > 100)
.select("sender", "amount", "district"))
# এই পর্যন্ত কিছুই execute হয়নি — শুধু DAG তৈরি
# আরও transformation
agg = (txns.groupBy("district")
.agg(_sum("amount").alias("total")))
# action — এখানে পুরো pipeline চলবে
agg.show()
filter দু'টি একত্রে করবে, select read time-ই apply (column pruning), groupBy partial aggregation প্রতিটি partition-এ — তারপর shuffle। এই সব optimization আপনাকে লিখতে হয়নি।
৭ · Spark vs Hadoop MapReduce
# MapReduce-এ word count — যেভাবে ২০০৬-এ লিখতে হত (পুরোটা Java boilerplate)
# Mapper: line → (word, 1)
# Reducer: (word, [1,1,1,...]) → (word, total)
# প্রতিটি stage-এর পর HDFS-এ disk-write। Slow।
# Spark-এ — এক লাইনে
spark.read.text("hdfs://logs/").rdd \
.flatMap(lambda r: r.value.split()) \
.map(lambda w: (w, 1)) \
.reduceByKey(lambda a, b: a + b) \
.collect()
৮ · Cluster modes — কোথায় চালাবেন
- Local: এক machine-এ — development ও test।
- Standalone: Spark-এর নিজস্ব cluster manager — সরল setup।
- YARN: Hadoop ecosystem-এ — অনেক enterprise-এ এখনো default।
- Kubernetes: ২০২০+ — cloud-native, dynamically scale, isolated containers।
- Managed (Databricks, EMR, Dataproc): production reality। autoscaling, monitoring, optimized runtime।
৯ · Fault tolerance — কোনো executor মারা গেলে?
RDD/DataFrame-এর "Resilient" — এই শব্দটাই key। Spark প্রতিটি partition-এর "lineage" রাখে — কীভাবে এটি তৈরি হয়েছে। কোনো executor crash হলে — Spark সেই partition অন্য executor-এ recompute করে। User intervention ছাড়াই।
Caching ব্যবহার করলে (df.cache()) — partition memory-তে থাকে। এটা lost হলে recompute lineage থেকে। তাই lineage = checkpoint, কিন্তু lightweight।
ভাবনার প্রশ্ন
প্রতিটি প্রশ্ন নিজে কিছুক্ষণ ভাবুন — তারপর "→ উত্তর" চাপুন।
প্র ০১ "Spark সবসময় MapReduce-এর চেয়ে দ্রুত" — এই দাবি কতটা সত্য? কী কী scenario-এ MapReduce আজও বেছে নেওয়া হয়, এবং Spark-এর কোন hidden cost আছে?
Marketing brochure-এ "১০০× faster" — বাস্তবে nuanced। DE হিসেবে এই myth-গুলো বুঝে নেওয়া জরুরি।
Spark জেতে যখন:
- Iterative workload: ML training (gradient descent), graph algorithm (PageRank) — একই data বহুবার scan। Memory-তে cache → MR-এর চেয়ে ১০-১০০× দ্রুত।
- Interactive query: notebook-এ ad-hoc analysis — sub-second response possible।
- Multi-stage pipeline: ১০ stage-এর pipeline, MR-এ ১০ disk-write; Spark-এ ১টা।
MapReduce আজও জেতে:
- Single massive scan: ১PB log-এ একবার count — সরল, predictable; memory pressure কম।
- Memory-constrained cluster: Spark in-memory — RAM কম থাকলে spill, slow।
- Truly batch, infrequent: দিনে একবার ৫ ঘণ্টার job — Spark-এর extra complexity-র মূল্য নেই।
- Existing Hadoop ecosystem: বহু legacy enterprise-এ MR pipeline ২০০৯ থেকে। rewrite cost প্রচুর।
Spark-এর hidden cost:
- Memory tuning: executor RAM, driver RAM, shuffle partition — wrongly configured মানে OOM error, retry storm।
- Skew problem: এক key-তে ৯০% data — সেই partition slow, পুরো job stuck।
- Garbage collection: long-running JVM-এ GC pause — performance variability।
- Cluster cost: idle executor RAM খরচ করে — managed cloud-এ $/hour।
- Debugging hard: distributed stack trace, lazy evaluation-এ stack mismatch — beginner-এর কাছে frustrating।
- Version compatibility: Spark major upgrade-এ API change, library incompatibility।
আজকের landscape: Hadoop MapReduce আর recommended না নতুন কাজে। Spark default; alternatives: Flink (streaming-first), Trino/Presto (interactive SQL), DuckDB/Polars (single-machine-এ Spark-replacement বেশিরভাগ workload-এ)।
মূল উপলব্ধি: "Spark fast" — যদি workload fit হয়, configuration ঠিক হয়, team expert হয়। Tool-এর সাথে discipline চাই।
প্র ০২ Lazy evaluation কেন এত শক্তিশালী, কিন্তু debug করতে এত কঠিন? "Action-এ গিয়ে error" — এই pattern কেন এবং কীভাবে handle করবেন?
Lazy evaluation Spark-এর ক্ষমতা ও যন্ত্রণা — দু'টোই।
কেন এত শক্তিশালী:
- Whole-plan optimization: Catalyst পুরো DAG দেখে — কোন filter কোন join-এর আগে push, কোন column drop, কোন stage merge — সব auto।
- Predicate pushdown: filter parquet read-time-ই apply — সব data না read হয়।
- Column pruning: select-এ যে column নাই — read-ই হবে না।
- Join reordering: small table-কে আগে। Bigger table broadcast।
- Constant folding, projection elimination — compiler-এর মতো।
Debug কেন কঠিন:
- Stack trace mismatch: error যেখানে appear করে — line ৫০, কিন্তু আসল bug line ১০-এ যেখানে transformation define হয়েছিল।
- Type errors at runtime: DataFrame schema runtime-এ resolve — typo IDE catch করে না।
- Repeated computation: এক DataFrame-এর উপর দু'টি action — ২ বার compute। অবাক করা slowness।
- Driver vs executor bug: error driver-এ থাকতে পারে, executor-এ থাকতে পারে — log আলাদা জায়গায়।
- Closure serialization: Python lambda-তে nonsense object capture — pickle error।
"Action-এ গিয়ে error" — common pattern:
- Schema mismatch — file-এ যা ছিল, code-এ যা ভাবলেন — ভিন্ন।
- Path doesn't exist — read define হলো, action-এ open চেষ্টা।
- Memory overflow — collect()-এ TB driver-এ আনতে চাইলে।
- Skew-induced timeout — এক task ১০ ঘণ্টায়, বাকি ১ মিনিট।
Debugging strategy:
- Schema first:
df.printSchema(),df.show(5)— early, often। - explain():
df.explain(extended=True)— Catalyst plan দেখান, optimization correct? - Action incrementally: chain-এর প্রতি step-এ
.count()— কোথায় ভাঙছে। - Spark UI: stage, task, time, shuffle — visual debugging।
- Cache strategically: repeated computation এড়াতে।
- Sample first: production data-র ১% নিয়ে develop, scale পরে।
- collect() নিষিদ্ধ production-এ — driver OOM।
মূল উপলব্ধি: Lazy evaluation একটি contract — Spark বেশি optimize পায়, আপনি deferred error পান। ভাল workflow ও observability ছাড়া distributed computing dangerous।
প্র ০৩ "DuckDB / Polars Spark-কে obsolete করে দিচ্ছে" — এই claim ২০২৩-২৪-এ অনেক viral। সত্যি কি? কখন single-machine tool যথেষ্ট, কখন Spark অপরিহার্য?
এটি modern data engineering-এর অন্যতম গুরুত্বপূর্ণ debate। "Big data is dead" — Jordan Tigani-র ২০২৩-এর viral post এই বিতর্ক উসকে দেয়।
DuckDB ও Polars কী:
- DuckDB: SQLite-এর OLAP version। Single-machine, embedded, columnar। parquet/CSV read super-fast।
- Polars: Rust-based DataFrame library। Pandas-এর successor — multi-threaded, lazy, memory-efficient।
কেন এরা Spark-কে threaten করছে:
- Hardware বদলেছে: ২০২৩-এ এক laptop-এ ৬৪GB RAM সাধারণ। AWS r6i.32xlarge-এ ১TB RAM, ১২৮ core। আগের "big data" এখন laptop-এ fit।
- Most data smaller than people think: survey-এ — ৮০%+ analytical workload < ১TB।
- Speed: ১০০GB data-তে DuckDB Spark-এর চেয়ে ৫-১০× দ্রুত (no shuffle overhead)।
- Simplicity: cluster setup নেই, JVM tuning নেই, network shuffle নেই। বিজনেস কোড লেখা যায়।
- Cost: Spark cluster idle-ও $/hour। DuckDB process-এই চলে।
কখন single-machine যথেষ্ট:
- Data ≤ ১TB।
- Workload ad-hoc analytics বা scheduled batch।
- Latency requirement মাঝারি (মিনিট-ঘণ্টা)।
- Team Spark expert নয়।
কখন Spark অপরিহার্য:
- Data কোটি TB scale (Grameenphone CDR, Daraz click stream multi-year)।
- Workload truly distributed চাই — fault tolerance critical।
- Streaming + batch একসাথে (Structured Streaming)।
- Existing Spark ecosystem (Delta Lake, MLlib, Iceberg)।
- Data warehouse external — শুধু processing layer চাই।
২০২৬-এর pragmatic approach:
- Default-এ DuckDB/Polars try করুন (যদি data laptop-এ fit হয়)।
- Spark যখন স্পষ্ট justify হয়।
- "Just in case" Spark choosing — outdated mindset।
মূল উপলব্ধি: Tools commoditize হয় hardware advance-এর সাথে। আজকের "big data" আগামীকালের "small data"। DE হিসেবে — fundamentals (SQL, data modeling, distributed concept) শিখুন; tool change করবে। Spark শিখুন কারণ industry এখন এতে চলে — কিন্তু সবকিছু Spark-এ লেখা ভুল।
প্র ০৪ Daraz analytics team-এ আপনি Spark cluster-এর জন্য AWS budget plan করছেন। প্রতিদিন ৫০০GB data process। কী cluster size, কোন instance, কীভাবে cost optimize করবেন?
Spark cluster sizing — DE-র অন্যতম practical দক্ষতা। ভুলে cost ১০× বেশি বা job ১০× slow।
Workload sizing:
- ৫০০GB raw → parquet-এ ~১০০-১৫০GB (compression)।
- Daily job: ETL clean → join → aggregate → write। Estimated ৩০-৬০ মিনিট acceptable।
- Spark guideline: total cluster RAM 2-3× data size। ৪০০-৫০০GB RAM target।
Cluster shape options:
- Option A — fewer big nodes: ৪× r6i.4xlarge (১২৮GB RAM, ১৬ vCPU each) = ৫১২GB, ৬৪ core। Less network shuffle।
- Option B — many small nodes: ১৬× m6i.2xlarge (৩২GB RAM, ৮ vCPU) = ৫১২GB, ১২৮ core। More parallelism, more shuffle।
- Option C — memory-optimized: ৪× r6i.4xlarge (RAM-heavy)। ML workload-এ ভাল।
Daraz ETL-এর জন্য — Option A বা C। compute-bound হলে B।
Cost optimization tactics:
- Spot instance: ৭০-৯০% সস্তা। Driver on-demand, executor spot। Spark fault-tolerant — spot interruption survive করে।
- Auto-scaling: EMR/Databricks dynamic allocation — busy hour-এ বেশি, idle-এ কম।
- Job-cluster (transient): persistent cluster না — daily job-এর জন্য spawn-process-terminate। Idle হলেও cost $0।
- Graviton (ARM): AWS m6g/r6g — ১৫-২০% সস্তা, comparable performance।
- Right partition count: default ২০০ shuffle partition — ৫০০GB-এ ২৫০০-৩০০০ ভাল। কম হলে slow, বেশি হলে overhead।
- Storage: S3 standard নয় — Intelligent-Tiering, Glacier old data।
- Format: CSV → Parquet (10-50× সস্তা read)।
- Partition pruning: S3-তে date partition — শুধু আজকের data read।
- Caching strategically: repeated DataFrame cache; না হলে memory waste।
- Photon/Tungsten engines: Databricks Photon ২× দ্রুত — same cost-এ কম time।
Estimated cost:
- Option A on-demand: $5-7/hour × ১ hr/day × ৩০ = $150-210/মাস।
- Option A spot + transient: $1-2/hour × ১ hr/day × ৩০ = $30-60/মাস।
- ৫-৭× savings — same job।
Monitoring:
- Spark UI — stage time, shuffle size, task skew।
- Ganglia/CloudWatch — CPU, RAM utilization (idle হলে cluster বড়, OOM হলে ছোট)।
- Cost Explorer tag-based — per-job cost track।
মূল উপলব্ধি: Production Spark = Spark + DevOps + FinOps। Bangladesh-এ data team প্রায়ই overspend করে — defaults মেনে। Tuning ১০০ ঘণ্টা invest করলে — মাসে $২০০০-$৫০০০ সাশ্রয় সম্ভব।
অনুশীলন
-
Lazy বনাম Eager চিন্তা: এই কোডের কোন লাইনে আসলে কাজ শুরু হবে এবং কেন?
df = spark.read.parquet("s3://logs/") df2 = df.filter("status='ERROR'") df3 = df2.select("timestamp","message") print("ready") df3.show(10)read,filter,select— সবই transformation, lazy। শুধু DAG তৈরি।print("ready")immediately চলবে কিন্তু Spark কিছু process করেনি।df3.show(10)— এখানে action। তখন Spark পুরো plan optimize করে — প্রায় সম্ভবfilterparquet read time-এই apply করবে (predicate pushdown), শুধু দু'টি column read করবে (column pruning)। তারপর ১০ row নিয়ে driver-এ ফেরত পাঠাবে।Performance lesson:
show(10)-এ সব data process হতে পারে না — Spark "limit pushdown"-এ smart, প্রথম ১০ row পেয়েই থামতে পারে। -
Architecture বুঝুন: ৫ executor × ৪ core প্রতি executor। ১০০ partition-এর একটি DataFrame — প্রথম "wave"-এ কতটি task parallel চলবে, কতটি wait?
Total parallel task slot = ৫ × ৪ = ২০।
প্রথম wave-এ ২০ task একসাথে চলবে, ৮০ task wait। প্রতিটি task শেষে নতুন task assigned। মোট ৫ wave (১০০ ÷ ২০)।
যদি প্রতিটি task ১ মিনিট নেয় — কাজ শেষ ৫ মিনিটে। ১ executor (৪ slot) হলে ২৫ মিনিট। Linear scale।
Caveat: data skew থাকলে এই নিখুঁত গণনা ভাঙে — এক task হয়তো ১০ মিনিট নেবে, পুরো job তখন ১০-১৫ মিনিট।
-
RDD নাকি DataFrame: এই use-cases-এর জন্য কোনটি? (ক) NLP-তে raw text-এ regex pattern matching। (খ) Sales aggregate by region by month। (গ) Customer graph traversal।
- (ক) NLP regex: DataFrame + UDF বা
regexp_extract। RDD একসময় preferred ছিল low-level control-এর জন্য, কিন্তু এখন built-in functions enough। - (খ) Sales aggregate: DataFrame, sure.
groupBy().agg()— Catalyst optimize করবে। - (গ) Graph traversal: GraphX (RDD-based) বা GraphFrames (DataFrame-based)। আজকে GraphFrames preferred।
সাধারণ rule: ৯০% কাজে DataFrame। RDD কেবল legacy code বা truly custom partitioning।
- (ক) NLP regex: DataFrame + UDF বা
আরও পড়ুন · ABCL TECH-এ আপনার পরবর্তী পদক্ষেপ
- পাঠ ১১ · PySpark DataFrame ও SQL পরবর্তী পাঠ Theory থেকে practice — Python-এ Spark write।
- পাঠ ০৯ · MongoDB ও NoSQL আগের পাঠ Spark যে variety-র data process করে — তার একটা।
- পাঠ ১২ · Spark optimization এই পাঠের সাথে সম্পর্কিত Cluster শক্তি কাজে লাগাতে — partition, skew, shuffle।
- সব AI Courses দেখুন ABCL TECH Python, ML, DL, NLP, CV, GenAI, RL, MLOps — সব AI কোর্স একসাথে।
!pip install pyspark দিয়ে শুরু করুন — single-node Spark, কোনো cluster setup ছাড়া। Production scale-এর জন্য Databricks Community Edition (free)।