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

PySpark DataFrame ও SQL

PySpark — DataFrame transformations & Spark SQL
৮ মিনিট পড়া মাঝারি · Intermediate Python · PySpark

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

  • SparkSession তৈরি ও CSV/Parquet read
  • DataFrame transformation — select, filter, groupBy, join
  • Window function ও aggregate function
  • Spark SQL — DataFrame ও SQL একসাথে
  • Write modes ও partitionBy — production-ready output

১ · SparkSession — শুরু এখানে

Spark ২.০ থেকে SparkSessionSparkSessionSpark application-এর entry point। ২.০-র আগে SQLContext, HiveContext, SparkContext আলাদা ছিল — SparkSession সব unify করেছে। হলো একমাত্র entry point। পুরোনো SparkContext, SQLContext, HiveContext — সব এর ভিতরে।

Python · PySpark
from pyspark.sql import SparkSession

spark = (SparkSession.builder
         .appName("DarazSalesETL")
         .config("spark.sql.shuffle.partitions", "200")
         .config("spark.sql.adaptive.enabled", "true")    # AQE (Spark 3+)
         .getOrCreate())

# Spark version check
print(spark.version)
# 3.5.0

# logging কম করতে
spark.sparkContext.setLogLevel("WARN")

    
getOrCreate() — existing session থাকলে reuse, না থাকলে নতুন। spark.sql.adaptive.enabled — Adaptive Query Execution, runtime-এ plan adjust। Spark ৩+ থেকে চমৎকার feature।

২ · Read CSV ও Parquet

প্রতিদিন DE-রা CSV/Parquet/JSON/Avro — যেকোনো source থেকে read করেন। Spark সবগুলোই handle করে একই API-তে।

Python · PySpark
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, IntegerType, TimestampType

# Schema আগে define করা best practice (inferSchema slow + ভুল করতে পারে)
order_schema = StructType([
    StructField("order_id", StringType(), False),
    StructField("customer_id", StringType(), False),
    StructField("district", StringType(), True),
    StructField("amount", DoubleType(), False),
    StructField("items", IntegerType(), False),
    StructField("created_at", TimestampType(), False),
])

# CSV — Daraz-এর daily order export
orders = (spark.read
          .option("header", True)
          .option("dateFormat", "yyyy-MM-dd")
          .schema(order_schema)
          .csv("s3://daraz-raw/orders/2026/05/09/*.csv"))

print(f"Total rows: {orders.count():,}")
orders.printSchema()
orders.show(5, truncate=False)

# Parquet — অনেক দ্রুত (columnar, compressed)
products = spark.read.parquet("s3://daraz-curated/products/")

    
Schema explicitly দেওয়া উত্তম — inferSchema=True পুরো file scan করে type detect করে (slow)। Production-এ schema fixed রাখুন; data type drift detect করতে আলাদা validation।

৩ · DataFrame transformations

Python · PySpark
from pyspark.sql.functions import col, when, lit, upper, length, to_date

# select — শুধু কয়েকটি column
slim = orders.select("order_id", "district", "amount", "created_at")

# filter / where (synonym)
high_value = orders.filter(col("amount") > 5000)
dhaka_orders = orders.where(col("district") == "Dhaka")

# একাধিক condition
big_dhaka = orders.filter((col("amount") > 5000) & (col("district") == "Dhaka"))

# withColumn — নতুন column যোগ
enriched = (orders
            .withColumn("amount_usd", col("amount") / 110)
            .withColumn("order_date", to_date("created_at"))
            .withColumn("is_high_value",
                when(col("amount") > 10000, lit("HIGH"))
                .when(col("amount") > 1000, lit("MEDIUM"))
                .otherwise(lit("LOW")))
            .withColumnRenamed("district", "delivery_district"))

enriched.show(5)

    
col() — column reference; magic literal-এর জন্য lit()। when().otherwise() — SQL CASE WHEN-এর equivalent। Pandas user-দের চেনা — কিন্তু প্রতিটি transformation distributed cluster-এ চলবে।

৪ · groupBy ও aggregation

Python · PySpark
from pyspark.sql.functions import sum as _sum, avg, count, countDistinct, max as _max

# District-wise sales summary
district_summary = (orders
    .groupBy("district")
    .agg(
        _sum("amount").alias("total_sales"),
        avg("amount").alias("avg_order"),
        count("*").alias("order_count"),
        countDistinct("customer_id").alias("unique_customers"),
        _max("amount").alias("biggest_order")
    )
    .orderBy(col("total_sales").desc()))

district_summary.show(10, truncate=False)

# একাধিক column-এ group
daily_district = (orders
    .withColumn("date", to_date("created_at"))
    .groupBy("date", "district")
    .agg(_sum("amount").alias("revenue"))
    .orderBy("date", col("revenue").desc()))

    
sum, max built-in Python function — তাই sum as _sum rename কর্মক্ষেত্রে standard practice। groupBy-এর পর shuffle হয় — partition cross করে — তাই এটা expensive। যত দেরিতে করা যায় ততই ভাল।

৫ · Join — সবচেয়ে ব্যবহৃত operation

Python · PySpark
from pyspark.sql.functions import broadcast

# customers — বড় table
customers = spark.read.parquet("s3://daraz-curated/customers/")

# orders + customers
joined = orders.join(customers, on="customer_id", how="inner")

# left join — সব order, customer info থাকলে যোগ
left = orders.join(customers, "customer_id", "left")

# multiple key
multi = orders.join(payments,
    (orders.order_id == payments.order_id) &
    (orders.created_at == payments.txn_at),
    "inner")

# broadcast join — ছোট table-এর জন্য (সব executor-এ copy পাঠায়)
# districts — মাত্র ৬৪ row, broadcast obvious choice
districts = spark.read.parquet("s3://reference/bd_districts/")
fast_join = orders.join(broadcast(districts), "district", "left")

    
broadcast() — ছোট table (<১০০MB) সব executor-এ copy। বড় table shuffle করতে হয় না — dramatic speedup। Spark ৩+ AQE auto-broadcast detect করে কখনো কখনো।
PySpark — Typical ETL Pipeline Read → Transform → Aggregate → Write Read spark.read.csv() S3, HDFS, JDBC Clean filter, dropna, withColumn Enrich · Join join, broadcast, window Aggregate groupBy, agg, SQL view Write df.write.partitionBy(date).parquet() overwrite | append Catalyst Optimizer — পুরো DAG analyze করে, action-এর সময় predicate pushdown · column pruning · join reorder · constant folding
প্রতিটি ETL job এই pattern-এ। Catalyst পেছন থেকে সব optimize করে — আপনি pure logic-এ মনোযোগ দিন।

৬ · Window function — running total, rank

Python · PySpark
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, rank, lag, sum as _sum

# প্রতিটি customer-এর order ranked by date
w = Window.partitionBy("customer_id").orderBy(col("created_at"))

ranked = orders.withColumn("order_seq", row_number().over(w))

# প্রতিটি customer-এর running total
running = orders.withColumn(
    "lifetime_spend",
    _sum("amount").over(w.rowsBetween(Window.unboundedPreceding, Window.currentRow))
)

# Top 3 order per customer
w2 = Window.partitionBy("customer_id").orderBy(col("amount").desc())
top3 = (orders.withColumn("rnk", rank().over(w2))
              .filter(col("rnk") <= 3))

top3.show(20)

    
Window function — SQL-এর মতো। customer-wise sequence, running aggregate, top-N — সব এক pattern-এ। বড় partition-এ slow হতে পারে — partition strategy চিন্তা করুন।

৭ · Spark SQL — DataFrame ও SQL একসাথে

Python · PySpark
# DataFrame-কে temp view হিসেবে register
orders.createOrReplaceTempView("orders")
customers.createOrReplaceTempView("customers")

# এখন pure SQL
result = spark.sql("""
    SELECT
        c.district,
        COUNT(DISTINCT o.customer_id) AS active_customers,
        SUM(o.amount) AS revenue,
        AVG(o.amount) AS avg_order
    FROM orders o
    JOIN customers c ON o.customer_id = c.customer_id
    WHERE o.created_at >= '2026-05-01'
      AND o.amount > 0
    GROUP BY c.district
    HAVING revenue > 100000
    ORDER BY revenue DESC
    LIMIT 10
""")

result.show()
result.write.mode("overwrite").parquet("s3://daraz-curated/district_kpi/")

    
DataFrame API ও Spark SQL — একই Catalyst optimizer-এ যায়। Performance equivalent। SQL analyst-friendly, DataFrame programmatic-friendly। দু'টোই project-এ মিশিয়ে ব্যবহার করতে পারেন।

৮ · Write modes ও partitionBy

  • overwrite: আগের data মুছে নতুন লিখুন।
  • append: existing-এ যোগ।
  • ignore: already exists হলে কিছু না করুন।
  • error (default): exists হলে throw।
Python · PySpark
# Production-style write
(district_summary.write
    .mode("overwrite")
    .partitionBy("district")          # প্রতি district আলাদা folder
    .option("compression", "snappy")
    .parquet("s3://daraz-curated/district_summary/"))

# CSV (analyst-এর জন্য)
(district_summary
    .coalesce(1)                       # ১টি ফাইল
    .write.mode("overwrite")
    .option("header", True)
    .csv("s3://daraz-curated/district_summary_csv/"))

# JDBC — Postgres-এ লেখা
(district_summary.write
    .format("jdbc")
    .option("url", "jdbc:postgresql://db.example.com:5432/dwh")
    .option("dbtable", "district_summary")
    .option("user", "etl")
    .option("password", "...")
    .mode("overwrite")
    .save())

    
partitionBy — Hive-style folder structure। পরে যখন WHERE district='Dhaka' দিয়ে read করবেন — Spark শুধু সেই folder পড়বে (partition pruning)। ১০০× speedup possible। coalesce(1) dangerous — driver-এ memory pressure।
partitionBy বনাম bucketBy: partitionBy directory তৈরি — high-cardinality column-এ ভয়ংকর (১ লাখ folder)। bucketBy fixed bucket — uniform। District ৬৪ — partition OK; customer_id লাখ — ভুল।

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

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

প্র ০১ Pandas-এ যেভাবে DataFrame use করেন PySpark-এ ঠিক সেভাবে করলে production-এ অনেক mistake হয়। সবচেয়ে common pitfall কী, কেন এগুলো subtle?

PySpark-এর syntactic similarity Pandas-এর সাথে — উপকারী, কিন্তু distinct semantics। পার্থক্য না বুঝলে production fire।

Pitfall ১: collect() ব্যবহার

  • Pandas-এ df মানে data RAM-এ — সব সময় access।
  • PySpark-এ df = distributed। df.collect() = সব data driver-এ আনুন। ১TB হলে driver crash।
  • Solution: show() debug-এর জন্য, take(n) small sample, write file-এ। collect কোনদিন না — exception sample/test।

Pitfall ২: Loop-এ DataFrame

  • Pandas-এ for row in df.iterrows() চলে (slow, কিন্তু চলে)।
  • PySpark-এ for distributed না — driver-এ চলবে, পুরো data fetch।
  • Solution: vectorized operation, UDF, বা foreach action।

Pitfall ৩: Repeated computation

  • df1.show(); df1.count(); df1.write() — tিনবার compute।
  • Solution: df1.cache() first action-এর আগে। তবে memory খরচ — চিন্তা করুন।

Pitfall ৪: Schema inference

  • inferSchema=True — full scan। 1TB CSV-এ ১৫ মিনিট waste।
  • Solution: Schema explicitly define। File metadata cache।

Pitfall ৫: == None filter

  • Pandas df[df.col == None] ভুলভাবে সব return করে।
  • PySpark filter(col("x") == None) — কিছুই match হয় না (NULL handling)।
  • Solution: filter(col("x").isNull())।

Pitfall ৬: UDF abuse

  • Python UDF — JVM ↔ Python serialization, slow।
  • Solution: built-in functions আগে। না পেলে — Pandas UDF (vectorized), Scala UDF (fastest), বা SQL function।

Pitfall ৭: Wide transformation overuse

  • প্রতিটি groupBy/join/orderBy = shuffle = network expensive।
  • Solution: filter আগে, narrow transformation আগে; aggregate সবশেষে।

Pitfall ৮: show() excessive

  • প্রতিটা show() action — পুরো DAG re-execute।
  • Solution: development cache + show; production show কমান।

মূল উপলব্ধি: Pandas single-machine, eager। PySpark distributed, lazy। API similar — mental model আলাদা। Spark UI বুঝুন — DAG visualization আপনার শিক্ষক।

প্র ০২ একই কাজ DataFrame API ও Spark SQL দু'ভাবেই করা যায়। Production-এ কোনটি বাছবেন কেন? কোন factor-এ DataFrame win, কোথায় SQL?

এই debate প্রতিটি data team-এ হয়। সঠিক উত্তর nuanced — দু'টোই পাশাপাশি ব্যবহার করা সবচেয়ে ভাল।

Performance: Identical। দু'টোই Catalyst-এ যায়, একই plan, একই execution।

DataFrame API জেতে যেখানে:

  • Programmatic composition: dynamic column list, conditional logic, function abstraction। SQL string-build ভঙ্গুর ও SQL injection risk।
  • Type safety (Scala): compile-time error catch।
  • Refactoring: IDE rename column reference সহজে; SQL string-এ painful।
  • Reusable transformations: Python function-এ wrap। Module-এ share।
  • Testability: unit test লিখতে সহজ। SQL test framework কম mature।
  • Linting: mypy, pylint syntax catch করে।

Spark SQL জেতে যেখানে:

  • Analyst-friendly: SQL সবাই পড়তে পারে। Python না জানা stakeholder-ও বুঝবে।
  • Complex queries: 7-table join, multi-CTE — SQL-এ পরিষ্কার, DataFrame chain-এ confusing।
  • BI tool integration: Tableau, Metabase, Looker — সবই SQL।
  • Migration from warehouse: existing PostgreSQL/Snowflake SQL — minimal change।
  • Window function syntax: arguably SQL-এ পরিষ্কার।
  • EXPLAIN debugging: SQL plan পড়া সহজ।

Production-এ pragmatic approach:

  • Heavy lifting (read, complex transformation, schema management) — DataFrame API।
  • Business logic (join, aggregate, KPI) — SQL string বা spark.sql(...)।
  • SQL file-এ রাখুন (.sql) — version control, peer review সহজ।
  • dbt-style project-এ — SQL-first, transformation models।

Hybrid example:

raw = spark.read.parquet(path)              # DataFrame
clean = raw.filter(col("status").isNotNull())  # DataFrame
clean.createOrReplaceTempView("clean_orders")
result = spark.sql(open("kpi.sql").read())  # SQL
result.write.parquet(out_path)              # DataFrame

Anti-patterns:

  • SQL string Python f-string-এ build with user input — SQL injection।
  • DataFrame chain ৩০ লাইন long — readability dies।
  • Same logic দু'বার — DataFrame ও SQL — diverge করে।

মূল উপলব্ধি: Tool ভাল-খারাপ নয় — readability, maintainability, team skill সব মিলিয়ে। বড় team-এ SQL audit-friendly; engineering-heavy team-এ DataFrame composable।

প্র ০৩ "Data skew" — Spark-এর সবচেয়ে কুখ্যাত problem। PySpark-এ এটি কীভাবে চিহ্নিত করবেন এবং কী কী technique দিয়ে fix করবেন?

Skew — Spark-এর নীরব killer। ১০০ task-এর ৯৯টা ১ মিনিটে শেষ, ১টা ৩ ঘণ্টা — পুরো job stuck।

Skew কীভাবে আসে:

  • Power-law distribution: "Dhaka" district-এ ৬০% data; বাকি ৬৩ district-এ ৪০%।
  • Hot key: NULL গুলো এক partition-এ যায়। "anonymous" customer_id সব ঐ একই hash bucket-এ।
  • Bot traffic: এক IP থেকে লক্ষ event — সেই key skewed।
  • Date partitioning: latest day-তে data বেশি, archive কম।

চিহ্নিত:

  • Spark UI Stages tab: task duration distribution। বেশিরভাগ ৩০ sec, ১টি ১ ঘণ্টা — clear skew।
  • Stage time vs total CPU time: wall-clock ১০ মিনিট কিন্তু aggregate CPU time ১০০ মিনিট — parallelism worked; ১০ মিনিট both — straggler।
  • df.groupBy(key).count().orderBy(desc("count")) — top-key analysis।
  • Partition size: df.rdd.glom().map(len).collect() — কত row per partition।

Fix techniques:

  1. Salting: hot key-তে random suffix যোগ। groupBy salted_key, তারপর re-aggregate।
    from pyspark.sql.functions import rand, floor
    df = df.withColumn("salt", floor(rand() * 100))
    salted = df.groupBy("district", "salt").agg(_sum("amount").alias("partial"))
    final = salted.groupBy("district").agg(_sum("partial").alias("total"))
  2. Adaptive Query Execution (AQE): Spark ৩+ — automatic skew handling। spark.sql.adaptive.skewJoin.enabled=true।
  3. Broadcast join: ছোট side থাকলে — shuffle এড়ান।
  4. Filter skewed key separately: hot key আলাদা handle, বাকি normally; union।
  5. Increase shuffle partitions: default ২০০ — বড় data-তে ২০০০-৫০০০।
  6. Pre-aggregate before join: data কমিয়ে join — row count নেমে আসে।
  7. Bucketing: table bucketed by join key — shuffle eliminated।
  8. Repartition by hash: df.repartition(N, "key") — explicit।

NULL key special case: WHERE key IS NOT NULL filter; NULL row separate union (যদি দরকার থাকে)।

মূল উপলব্ধি: Skew unique distributed problem — single-machine-এ এটি নেই। Spark mastery মানে data distribution বুঝে কাজ। Spark UI ছাড়া skew debug প্রায় অসম্ভব — UI-এর সাথে বন্ধুত্ব করুন।

প্র ০৪ Daraz-এ ৩ বছরের order data PySpark-এ process করছেন। CSV নাকি Parquet — কোন format কেন? partitionBy কোন column-এ? compression algorithm কী?

File format ও storage strategy — query performance ও cost-এ ১০০× difference দিতে পারে।

CSV vs Parquet:

  • CSV: row-based, text, no schema, no compression default। Human-readable। বহু source এই format-এ data দেয়।
  • Parquet: columnar binary, embedded schema, compressed। Apache standard, Spark/Snowflake/BigQuery — সবাই native।

Parquet কেন জেতে:

  • Compression: ৫-১০× ছোট। Snappy-compressed parquet ১০০GB CSV → ১০-২০GB।
  • Column pruning: ১০০ column-এর মধ্যে শুধু ৫টা select করলে — ৫টা-ই read। CSV সব read।
  • Predicate pushdown: footer statistics দিয়ে whole row-group skip।
  • Type preserve: int, decimal, timestamp ঠিক type-এ। CSV সব string।
  • Encoding: dictionary encoding repeated value কমিয়ে দেয়।

CSV কেন এখনো:

  • Source system Parquet support করে না।
  • Excel-এ open করতে।
  • Inter-team handoff (analyst-এ)।

Best practice: Raw landing zone — যা পান (CSV often)। Curated/staged — Parquet always।

partitionBy strategy — Daraz orders:

  • By date: partitionBy("year", "month", "day")। সবচেয়ে সাধারণ পাঠ pattern (date range query)।
  • By district: ৬৪ district — manageable folder count।
  • Combined: partitionBy("year", "month") ভাল; day যোগ করলে ১১০০ folder/year — borderline।

Partition cardinality rule of thumb:

  • Total partition < ১০০০ ideally।
  • প্রতিটি partition ১০০MB-১GB data।
  • High-cardinality (customer_id, order_id) — কখনো partition করবেন না।

Compression:

  • Snappy (default): fast compress/decompress, ভাল ratio। সবচেয়ে balanced।
  • ZSTD (Spark ৩.২+): better ratio, similar speed। Modern choice।
  • Gzip: better ratio কিন্তু slow। Cold archive।
  • Brotli, LZ4: niche use-case।

সাধারণ recommendation: Parquet + ZSTD + partitionBy(year, month) + sortBy(customer_id)।

Production architecture:

  • Bronze (raw): CSV/JSON থেকে আসে।
  • Silver (cleaned): Parquet, partitioned by date।
  • Gold (aggregate): small Parquet, BI-ready।
  • Delta Lake/Iceberg uses Parquet underneath, ACID transaction যোগ।

মূল উপলব্ধি: Storage choice = future query performance। ভাল decision আজ — পরে dramatic ROI।

অনুশীলন

  1. লিখুন: bKash transaction log (parquet)। প্রতিটি sender-এর গত ৭ দিনে total send_money এবং transaction count বের করুন। Top ১০ sender-এর data Postgres-এ লিখুন।
    from pyspark.sql.functions import col, sum as _sum, count, current_timestamp, expr
    
    txns = (spark.read.parquet("s3://bkash-curated/transactions/")
            .filter(col("status") == "success")
            .filter(col("txn_type") == "send_money")
            .filter(col("created_at") >= expr("current_timestamp() - INTERVAL 7 DAYS")))
    
    top10 = (txns.groupBy("sender")
                 .agg(_sum("amount").alias("total_amount"),
                      count("*").alias("txn_count"))
                 .orderBy(col("total_amount").desc())
                 .limit(10))
    
    (top10.write
          .mode("overwrite")
          .format("jdbc")
          .option("url", "jdbc:postgresql://...:5432/dwh")
          .option("dbtable", "top_senders_7d")
          .option("user", "etl")
          .option("password", "...")
          .save())

    Filter আগে → groupBy পরে। limit-এর পর ১০ row — driver-এ আনা OK।

  2. SQL দিয়ে: উপরের same কাজ Spark SQL-এ লিখুন।
    txns.createOrReplaceTempView("txns")
    
    result = spark.sql("""
        SELECT
            sender,
            SUM(amount) AS total_amount,
            COUNT(*) AS txn_count
        FROM txns
        WHERE status = 'success'
          AND txn_type = 'send_money'
          AND created_at >= CURRENT_TIMESTAMP() - INTERVAL 7 DAYS
        GROUP BY sender
        ORDER BY total_amount DESC
        LIMIT 10
    """)
    
    result.write.mode("overwrite").format("jdbc")...

    লক্ষ্য করুন — same Catalyst plan, identical performance।

  3. Window function: Daraz-এ প্রতিটি customer-এর first ও latest order-এর পার্থক্য (lifetime in days) বের করুন।
    from pyspark.sql.window import Window
    from pyspark.sql.functions import min as _min, max as _max, datediff
    
    # Approach 1: groupBy
    lifetime = (orders.groupBy("customer_id")
                      .agg(_min("created_at").alias("first_order"),
                           _max("created_at").alias("last_order"))
                      .withColumn("lifetime_days",
                                  datediff("last_order", "first_order")))
    
    # Approach 2: window function (যদি প্রতিটি row-এ context চান)
    w = Window.partitionBy("customer_id")
    enriched = (orders
        .withColumn("first_order", _min("created_at").over(w))
        .withColumn("last_order", _max("created_at").over(w))
        .withColumn("lifetime_days", datediff("last_order", "first_order")))

    Approach 1 দ্রুত (one row per customer)। Approach 2 প্রয়োজনীয় যদি raw row-এ extra column চান।

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

ব্রাউজারে PySpark চালাতে চান? Google Colab -এ একটি cell-এ !pip install pyspark দিয়ে শুরু — single-node Spark। অথবা Databricks Community Edition (free tier) production-style cluster experience।
পূর্ববর্তী পাঠ
পাঠ ১০ · Apache Spark পরিচিতি