PySpark DataFrame ও SQL
এই পাঠে যা শিখবেন
- 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 — সব এর ভিতরে।
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-তে।
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/")
inferSchema=True পুরো file scan করে type detect করে (slow)। Production-এ schema fixed রাখুন; data type drift detect করতে আলাদা validation।
৩ · DataFrame transformations
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
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
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 করে কখনো কখনো।
৬ · Window function — running total, rank
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)
৭ · Spark SQL — DataFrame ও SQL একসাথে
# 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/")
৮ · Write modes ও partitionBy
- overwrite: আগের data মুছে নতুন লিখুন।
- append: existing-এ যোগ।
- ignore: already exists হলে কিছু না করুন।
- error (default): exists হলে throw।
# 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।
ভাবনার প্রশ্ন
প্রতিটি প্রশ্ন নিজে কিছুক্ষণ ভাবুন — তারপর "→ উত্তর" চাপুন।
প্র ০১ 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,writefile-এ। collect কোনদিন না — exception sample/test।
Pitfall ২: Loop-এ DataFrame
- Pandas-এ
for row in df.iterrows()চলে (slow, কিন্তু চলে)। - PySpark-এ
fordistributed না — driver-এ চলবে, পুরো data fetch। - Solution: vectorized operation, UDF, বা
foreachaction।
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:
- 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")) - Adaptive Query Execution (AQE): Spark ৩+ — automatic skew handling।
spark.sql.adaptive.skewJoin.enabled=true। - Broadcast join: ছোট side থাকলে — shuffle এড়ান।
- Filter skewed key separately: hot key আলাদা handle, বাকি normally; union।
- Increase shuffle partitions: default ২০০ — বড় data-তে ২০০০-৫০০০।
- Pre-aggregate before join: data কমিয়ে join — row count নেমে আসে।
- Bucketing: table bucketed by join key — shuffle eliminated।
- 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।
অনুশীলন
-
লিখুন: 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।
-
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।
-
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-এ আপনার পরবর্তী পদক্ষেপ
- পাঠ ১২ · Spark optimization পরবর্তী পাঠ Catalyst, AQE, partition tuning — production-grade Spark।
- পাঠ ১০ · Apache Spark পরিচিতি আগের পাঠ Theory ভিত্তি — RDD, DataFrame, lazy evaluation।
- পাঠ ১৩ · Apache Airflow পরিচিতি এই পাঠের সাথে সম্পর্কিত Spark job-কে schedule ও orchestrate করা।
- সব AI Courses দেখুন ABCL TECH Python, ML, DL, NLP, CV, GenAI, RL, MLOps — সব AI কোর্স একসাথে।
!pip install pyspark দিয়ে শুরু — single-node Spark। অথবা Databricks Community Edition (free tier) production-style cluster experience।