পাঠ ১০ · ২৯-এর মধ্যে · মডিউল ২
Home / AI Courses / Data Engineering / Apache Spark পরিচিতি

Apache Spark পরিচিতি

Apache Spark — distributed processing made easy
৮ মিনিট পড়া মাঝারি · Intermediate Distributed

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

  • 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-এ। ।

Single machine কখন কম পড়ে

১) 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-এ রাখে।

Spark পাঁচটি library নিয়ে complete platform: Spark SQL (structured data), Spark Streaming (real-time), MLlib (machine learning), GraphX (graph processing), Structured Streaming (modern streaming)। সবই একই engine-এ।

৩ · 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।
ভাবুন Daraz-এর warehouse। Driver = warehouse manager — knows কী কী order, কারা কোন packing করছে। Executors = packing staff — ১০০ জন একসাথে box তৈরি করছে। Cluster manager = HR — কোন staff কোথায় থাকবে, কতজন assign হবে।
Apache Spark — Architecture Driver coordinates, Executors compute Driver SparkSession, DAG scheduler your main() lives here Cluster Manager YARN / K8s / Standalone Executor 1 RAM: cache partitions CPU: run tasks P1, P5, P9 Executor 2 RAM: cache partitions CPU: run tasks P2, P6, P10 Executor 3 RAM: cache partitions CPU: run tasks P3, P7, P11
Driver task assign করে; Cluster Manager resource ঠিক করে; Executor data process করে। ১২টি partition ৩ executor-এ ভাগ — parallel।

৪ · 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-এ।

Python · PySpark
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()

    
Spark internally filter দু'টি একত্রে করবে, select read time-ই apply (column pruning), groupBy partial aggregation প্রতিটি partition-এ — তারপর shuffle। এই সব optimization আপনাকে লিখতে হয়নি।

৭ · Spark vs Hadoop MapReduce

Bash
# 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()

    
Spark intermediate result memory-তে রাখে। MapReduce disk-এ লেখে। iterative algorithm-এ (যেমন PageRank) — Spark ১০০× দ্রুত। MapReduce আজও কোথাও কোথাও চলে — খুব বড়, একবারের batch job-এ যেখানে memory লিমিট সমস্যা।

৮ · 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।
Spark cluster setup tedious — তাই বেশিরভাগ company managed service ব্যবহার করে। বাংলাদেশে — Daraz, bKash analytics team Databricks বা AWS EMR-এ Spark চালায়। তাদের DE-রা cluster ops-এ সময় না দিয়ে business logic-এ মনোযোগ দিতে পারেন।

৯ · 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:

  1. Spot instance: ৭০-৯০% সস্তা। Driver on-demand, executor spot। Spark fault-tolerant — spot interruption survive করে।
  2. Auto-scaling: EMR/Databricks dynamic allocation — busy hour-এ বেশি, idle-এ কম।
  3. Job-cluster (transient): persistent cluster না — daily job-এর জন্য spawn-process-terminate। Idle হলেও cost $0।
  4. Graviton (ARM): AWS m6g/r6g — ১৫-২০% সস্তা, comparable performance।
  5. Right partition count: default ২০০ shuffle partition — ৫০০GB-এ ২৫০০-৩০০০ ভাল। কম হলে slow, বেশি হলে overhead।
  6. Storage: S3 standard নয় — Intelligent-Tiering, Glacier old data।
  7. Format: CSV → Parquet (10-50× সস্তা read)।
  8. Partition pruning: S3-তে date partition — শুধু আজকের data read।
  9. Caching strategically: repeated DataFrame cache; না হলে memory waste।
  10. 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 করলে — মাসে $২০০০-$৫০০০ সাশ্রয় সম্ভব।

অনুশীলন

  1. 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 করে — প্রায় সম্ভব filter parquet read time-এই apply করবে (predicate pushdown), শুধু দু'টি column read করবে (column pruning)। তারপর ১০ row নিয়ে driver-এ ফেরত পাঠাবে।

    Performance lesson: show(10)-এ সব data process হতে পারে না — Spark "limit pushdown"-এ smart, প্রথম ১০ row পেয়েই থামতে পারে।

  2. 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 তখন ১০-১৫ মিনিট।

  3. 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।

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

Spark-এ হাত মেলাতে চান? ব্রাউজারে Google Colab -এ !pip install pyspark দিয়ে শুরু করুন — single-node Spark, কোনো cluster setup ছাড়া। Production scale-এর জন্য Databricks Community Edition (free)।
পূর্ববর্তী পাঠ
পাঠ ০৯ · MongoDB ও NoSQL